CertSafari
    Databricks Certified Associate Developer for Apache Spark· Lessons

    Domain 3 · Lesson 16/32

    Spark DataFrame Joins and Unions: Inner, Left, Broadcast, Cross, Union

    Combine DataFrames with operations such as Inner join, left join, broadcast join, multiple keys, cross join, union, and union all.

    16 min read
    3.12% of exam
    10 sources
    Published 3 Oct 2026
    Docs as of 30 Sep 2026

    What you will be able to do

    • Write a DataFrame join with the right on argument (a column name, a list of names, or a Column expression) and the right how string
    • Predict the rows returned by inner, left and cross joins
    • Join on multiple keys and avoid ending up with duplicate key columns
    • Mark a small DataFrame for a broadcast join with broadcast(), and explain what spark.sql.autoBroadcastJoinThreshold controls
    • Choose between union, distinct(), SQL UNION/UNION ALL and unionByName when stacking DataFrames

    Key concept

    join(other, on, how) — Every DataFrame join is one call with three parts: the right-side DataFrame, an on argument that says how rows match, and a how string that picks the join type. Inner is the default. If on is a column name or a list of names, the join is an equi-join on columns that exist on both sides.

    1.The anatomy of DataFrame.join

    Every join in this lesson goes through DataFrame.join(other, on=None, how=None), apart from crossJoin, which has its own method. other is the right side of the join. on can take four forms: a single column name as a string, a list of column names, a join expression built from Columns, or a list of Column expressions. how names the join type. If you leave it out, you get an inner join.

    The form of on matters more than it looks. If you pass a string or a list of strings, the named columns must exist in both DataFrames, and Spark runs an equi-join on them. If you pass a Column expression such as df.name == df2.name, you can write any boolean condition, but both sides' columns stay in the result.

    Joining on a column name: the default `how` is inner, so only Bob, who appears on both sides, survives, and `name` appears oncepython
    import pyspark.sql.functions as sf
    from pyspark.sql import Row
    df = spark.createDataFrame([Row(name="Alice", age=2), Row(name="Bob", age=5)])
    df2 = spark.createDataFrame([Row(name="Tom", height=80), Row(name="Bob", height=85)])
    
    df.join(df2, "name").show()
    # +----+---+------+
    # |name|age|height|
    # +----+---+------+
    # | Bob|  5|    85|
    # +----+---+------+

    The how parameter accepts a fixed set of strings, and several are synonyms. left, leftouter and left_outer all mean the same join, and so do outer, full, fullouter and full_outer. The list also includes cross, semi/leftsemi/left_semi and anti/leftanti/left_anti. On the exam, a spelling like left_outer is a valid choice, not a typo.

    Checkpoint 1 of 9· Check yourself

    A candidate writes df.join(df2, "name") and passes no third argument. Which join type does Spark run?

    Sources1

    2.Inner join vs left join: which rows survive

    Inner and left joins differ in what happens to a row that has no match. Take an employee table with six people in departments 1–6 and a department table that lists only departments 1, 2 and 3.

    Left join: all six employees come back, and the three with no matching department get NULLsql
    -- Use employee and department tables to demonstrate left join.
    > SELECT id, name, employee.deptno, deptname
       FROM employee
       LEFT JOIN department ON employee.deptno = department.deptno;
     105 Chloe      5        NULL
     103  Paul      3 Engineering
     101  John      1   Marketing
     102  Lisa      2       Sales
     104  Evan      4        NULL
     106   Amy      6        NULL

    Change LEFT JOIN to INNER JOIN and only Paul, John and Lisa come back, because an inner join returns only rows with matching values in both tables. The DataFrame API follows the same rules. how="inner" (or no how) is the inner join, and how="left" is the left outer join. The left side is the DataFrame you call .join on, and the right side is other.

    Results of the employee/department join for each join type
    Join type (how / SQL)Unmatched left rowsUnmatched right rowsRows in this example
    inner / INNER JOINDroppedDropped3
    left / LEFT [OUTER] JOINKept, with NULLs on the rightDropped6
    full / FULL [OUTER] JOINKept, with NULLsKept, with NULLs6
    cross / CROSS JOINNo condition: every pairingNo condition: every pairing18

    Checkpoint 2 of 9· Check yourself

    You need a list of all customers. For each one, show their order total if they have one, and NULL if they don't. customers is the DataFrame you call .join on. Which how value is correct?

    Checkpoint 3 of 9· Exam question

    ```python employees = spark.createDataFrame( [(1, "Ana", 10), (2, "Ben", 20), (3, "Cara", None)], ["emp_id", "name", "dept_id"] ) departments = spark.createDataFrame( [(10, "Sales"), (20, "Engineering"), (30, "Marketing")], ["dept_id", "dept_name"] ) result = employees.join(departments, on="dept_id", how="inner") ``` How many rows does `result` contain, and why?

    Sources2

    3.Joining on multiple keys and avoiding duplicate key columns

    To join on several keys, pass a list. With a list of column names such as ["a", "b"], every named column must exist on both sides, and Spark runs an equi-join on all of them. A row matches only when every key is equal. SQL does the same thing with USING (a, b), which the SQL reference defines as an ON clause that joins the equality tests with AND, as in ON first.a = second.a AND first.b = second.b. For conditions that aren't plain equality, pass Column expressions instead, either as one expression combined with & or as a list of Columns.

    An expression join keeps both key columns. The result has two `name` columns, which makes a later reference to `name` ambiguouspython
    joined = df.join(df2, df.name == df2.name, "outer").sort(sf.desc(df.name))
    joined.show()
    # +-----+----+----+------+
    # | name| age|name|height|
    # +-----+----+----+------+
    # |  Bob|   5| Bob|    85|
    # |Alice|   2|NULL|  NULL|
    # | NULL|NULL| Tom|    80|
    # +-----+----+----+------+

    Joining by name (a string or list of strings in the DataFrame API, or USING/NATURAL in SQL) collapses each key into a single output column, placed first. An expression join keeps the key from each side. That is usually why a later select("name") fails as ambiguous. When you can't avoid an expression join, for example in a self-join, give each side an alias with df.alias("a") and refer to columns as "a.name" and "b.name".

    Checkpoint 4 of 9· Check yourself

    orders and shipments both have region and order_id columns. Which call joins on both keys and leaves one region and one order_id column in the result?

    Sources12

    4.Broadcast joins: sending the small side to every worker

    Joins often pair a large fact table with a small lookup table. A broadcast join copies the whole small table to all worker nodes, so each worker can join its part of the large table locally. A Databricks engineering post puts it this way: broadcast hash join is usually the fastest strategy when one side fits comfortably in memory.

    You can trigger it in two ways. The explicit way is pyspark.sql.functions.broadcast(df), which marks a DataFrame as small enough to broadcast. You then join as usual. The automatic way is the setting spark.sql.autoBroadcastJoinThreshold, which defaults to 10485760 bytes (10 MB). Spark plans a broadcast hash join on its own when a relation's estimated size is below this threshold. Setting it to -1 turns automatic broadcasting off. A BROADCAST join hint tells Spark to use that strategy for the relation you mark.

    Checkpoint 5 of 9· Fill the gap

    Which function marks df_small for a broadcast join?

    from pyspark.sql import functions as dbf
    df = spark.createDataFrame([1, 2, 3, 3, 4], "int")
    df_small = spark.range(3)
    df_b = dbf. ? (df_small)
    df.join(df_b, df.value == df_small.id).show()

    Checkpoint 6 of 9· Exam question

    A developer runs `df_small.join(df_large, "key")` without any join hints, and Spark's optimizer automatically converts it into a broadcast hash join. Under the default configuration, what determines whether Spark performs this automatic conversion?

    Sources345

    5.Cross joins: the Cartesian product

    A cross join has no join condition. It pairs every row on the left with every row on the right, so the result has left × right rows. In the DataFrame API you write it as df.crossJoin(other). In SQL you write CROSS JOIN. The same employee/department tables (6 rows and 3 rows) produce 18 rows.

    `crossJoin`: 3 people × 2 heights = 6 rowspython
    df.crossJoin(df2.select("height")).select("age", "name", "height"
        ).orderBy("age", "name", "height").show()
    # +---+-----+------+
    # |age| name|height|
    # +---+-----+------+
    # | 14|  Tom|    80|
    # | 14|  Tom|    85|
    # | 16|  Bob|    80|
    # | 16|  Bob|    85|
    # | 23|Alice|    80|
    # | 23|Alice|    85|
    # +---+-----+------+

    You can also get a cross join by accident. The SQL reference warns that if you leave out the join criteria, any join type behaves as a CROSS JOIN. Because the output grows with the product of the two inputs, Databricks' join-performance guidance says cross joins are expensive and should be removed from workloads that need low latency or frequent recomputation.

    Checkpoint 7 of 9· Check yourself

    A SQL query reads SELECT * FROM a INNER JOIN b with no ON or USING clause. What does it return?

    Sources627

    6.Stacking rows: union, UNION ALL and unionByName

    Joins add columns side by side. Unions stack rows, and the main question is what happens to duplicate rows.

    `union` keeps duplicates, and adding `.distinct()` removes thempython
    df1 = spark.createDataFrame([(1, 'A'), (2, 'B'), (3, 'C')], ['id', 'value'])
    df2 = spark.createDataFrame([(3, 'C'), (4, 'D')], ['id', 'value'])
    df3 = df1.union(df2).distinct().sort("id")
    df3.show()
    # +---+-----+
    # | id|value|
    # +---+-----+
    # |  1|    A|
    # |  2|    B|
    # |  3|    C|
    # |  4|    D|
    # +---+-----+

    SQL works the other way round. In SQL, UNION without a qualifier means UNION DISTINCT and removes duplicates. You write UNION ALL to keep them. Both sides must have the same number of columns with compatible types, otherwise you get NUM_COLUMNS_MISMATCH or INCOMPATIBLE_COLUMN_TYPE. So the DataFrame method union() gives you SQL's UNION ALL result. A note on scope: the sources for this lesson don't document a separate DataFrame unionAll method. They list only union and unionByName. What this lesson can teach is the union-all behaviour, which keeps duplicates.

    union() matches columns by position, not by name. If one DataFrame has id, value and the other has value, id, union() silently puts each value under the wrong column. unionByName matches columns by name instead. With allowMissingColumns=True, it also accepts columns that exist on only one side and fills the gaps with NULL.

    `unionByName` matches columns by name, so 4, 5, 6 land under `col1`, `col2`, `col0`python
    df1 = spark.createDataFrame([[1, 2, 3]], ["col0", "col1", "col2"])
    df2 = spark.createDataFrame([[4, 5, 6]], ["col1", "col2", "col0"])
    df1.unionByName(df2).show()
    # +----+----+----+
    # |col0|col1|col2|
    # +----+----+----+
    # |   1|   2|   3|
    # |   6|   4|   5|
    # +----+----+----+

    Checkpoint 8 of 9· Match them up

    Match each operation to how it treats duplicates and columns

    Tap a term, then the definition that fits it.

    Checkpoint 9 of 9· Exam question

    A pipeline joins a 500 GB `transactions` DataFrame with a 2 MB `country_codes` lookup DataFrame. `spark.sql.autoBroadcastJoinThreshold` has been set to `-1` in the cluster's Spark config, disabling automatic broadcast selection, but the developer still wants this specific join to avoid a shuffle. Which line correctly forces `country_codes` to be broadcast? ```python from pyspark.sql.functions import broadcast result = transactions.join( ____________, on="country_id", how="inner" ) ```

    Sources8910

    Exam traps

    Each one states something that sounds right. Open it to see what is actually true.

    1. 1.DataFrame.union() removes duplicate rows the way SQL UNION does.Why is that wrong?

      union() keeps duplicates, so it behaves like SQL UNION ALL. Add .distinct() if you want them removed.

      Covered in Stacking rows: union, UNION ALL and unionByName

    2. 2.union() lines columns up by name, so the column order of the two DataFrames doesn't matter.Why is that wrong?

      union() matches columns by position. When the column orders differ, use unionByName.

      Covered in Stacking rows: union, UNION ALL and unionByName

    3. 3.A SQL join with no ON or USING clause fails to run.Why is that wrong?

      Without join criteria, any join type behaves as a cross join and returns the Cartesian product, which can be very large.

      Covered in Cross joins: the Cartesian product

    4. 4.Spark only runs a broadcast join when you call broadcast().Why is that wrong?

      Spark also broadcasts on its own when a table's estimated size is below spark.sql.autoBroadcastJoinThreshold, which defaults to 10 MB. Setting the threshold to -1 turns this off.

      Covered in Broadcast joins: sending the small side to every worker

    Sources

    Every claim above is drawn from one of these pages, quoted as it was written on the date shown.

    1. 1.
      “the column(s) must exist on both sides, and this performs an equi-join.”
      ↩︎ The anatomy of DataFrame.join
      “default inner. Must be one of: inner, cross, outer, full, fullouter, full_outer, left, leftouter, left_outer”
      ↩︎ The anatomy of DataFrame.join
      “the column(s) must exist on both sides, and this performs an equi-join.”
      ↩︎ Joining on multiple keys and avoiding duplicate key columns
      “a string for the join column name, a list of column names, a join expression (Column), or a list of Columns.”
      ↩︎ Key concept
    2. 2.
      “Returns the rows that have matching values in both table references.”
      ↩︎ Inner join vs left join: which rows survive
      “Returns all values from the left table reference and the matched values from the right table reference”
      ↩︎ Inner join vs left join: which rows survive
      “Matches the rows by comparing equality for list of columns column_name which must exist in both relations.”
      ↩︎ Joining on multiple keys and avoiding duplicate key columns
      “SELECT * will only show one occurrence for each of the columns used to match first”
      ↩︎ Joining on multiple keys and avoiding duplicate key columns
      “If you omit the join_criteria the semantic of any join_type becomes that of a CROSS JOIN.”
      ↩︎ Cross joins: the Cartesian product
      “If you omit the join_criteria the semantic of any join_type becomes that of a CROSS JOIN.”
      ↩︎ Exam trap 3
    3. 4.
      “Configures the maximum size in bytes for a table that will be broadcast to all worker nodes when performing a join.”
      ↩︎ Broadcast joins: sending the small side to every worker
      “By setting this value to -1, broadcasting can be disabled.”
      ↩︎ Broadcast joins: sending the small side to every worker
      “instruct Spark to use the hinted strategy on each specified relation when joining them with another relation”
      ↩︎ Broadcast joins: sending the small side to every worker
      “Configures the maximum size in bytes for a table that will be broadcast to all worker nodes when performing a join.”
      ↩︎ Exam trap 4
    4. 5.
      “broadcast hash join is usually the most performant if one side of the join can fit well in memory.”
      ↩︎ Broadcast joins: sending the small side to every worker
      “Spark plans a broadcast hash join if the estimated size of a join relation is lower than the broadcast-size threshold.”
      ↩︎ Broadcast joins: sending the small side to every worker
    5. 8.
      “with no automatic deduplication of elements.”
      ↩︎ Stacking rows: union, UNION ALL and unionByName
      “The method resolves columns by position (not by name), following the standard behavior in SQL.”
      ↩︎ Stacking rows: union, UNION ALL and unionByName
      “with no automatic deduplication of elements.”
      ↩︎ Exam trap 1
      “The method resolves columns by position (not by name), following the standard behavior in SQL.”
      ↩︎ Exam trap 2
    6. 9.
      “If ALL is specified duplicate rows are preserved. If DISTINCT is specified the result does not contain any duplicate rows. This is the default.”
      ↩︎ Stacking rows: union, UNION ALL and unionByName
      “Both subqueries must have the same number of columns and share a least common type for each respective column.”
      ↩︎ Stacking rows: union, UNION ALL and unionByName
    7. 10.
      “resolving columns by name (rather than position). When allowMissingColumns is True, missing columns will be filled with null.”
      ↩︎ Stacking rows: union, UNION ALL and unionByName

    Ready to test yourself?

    Practise Databricks Certified Associate Developer for Apache Spark in quiz mode.

    Spotted a mistake, or was something unclear? Tell us.