ML Engineer MasterClass (October) | 4 seats left

Spark · Union and schema alignment
AmazonAmazon Analytics
Amazon · Combine DataFrames

Union and schema alignment

Stack rows safely, align fields by name, and handle missing columns.

Step 1 of 6 · Learn

Stack rows

unionByName appends rows after matching column names; it does not remove duplicates. The example combines the same two orders twice, producing four rows. Applying distinct afterwards removes exact duplicate rows. Neither operation guarantees output order, so we sort explicitly.

Lesson reference: PySpark and Scala

Stack rows

unionByName appends rows after matching column names; it does not remove duplicates. The example combines the same two orders twice, producing four rows. Applying distinct afterwards removes exact duplicate rows. Neither operation guarantees output order, so we sort explicitly.

PySpark example

from pyspark.sql.functions import col, lit, struct, array, size, element_at, filter, transform


result = orders.filter(col("order_id") <= 1002).unionByName(
    orders.filter(col("order_id") <= 1002)
).orderBy("order_id")
result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._


val result = orders.filter(col("order_id") <= 1002).unionByName(
    orders.filter(col("order_id") <= 1002)
).orderBy("order_id")
result.show(false)

Align by name

The second batch has customer_id before order_id. union aligns by position, silently mixing these numeric fields; the example deliberately shows that incorrect result. unionByName instead uses names. Matching names does not excuse incompatible data types: validate the schemas too.

PySpark example

from pyspark.sql.functions import col, lit, struct, array, size, element_at, filter, transform

first = orders.filter(col("order_id") <= 1002).select("order_id", "customer_id")
second = orders.filter(col("order_id") >= 1005).select("customer_id", "order_id")

result = first.union(second).orderBy("order_id", "customer_id")
result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._

val first = orders.filter(col("order_id") <= 1002).select("order_id", "customer_id")
val second = orders.filter(col("order_id") >= 1005).select("customer_id", "order_id")

val result = first.union(second).orderBy("order_id", "customer_id")
result.show(false)

Handle missing columns

The second batch lacks status. allowMissingColumns adds that column with null values. The missing status is not automatically inferred from the original orders table. A placeholder can make the absence explicit, but is not a recovered fact.

PySpark example

from pyspark.sql.functions import col, lit, struct, array, size, element_at, filter, transform

first = orders.filter(col("order_id") <= 1002).select("order_id", "status")
second = orders.filter(col("order_id") >= 1005).select("order_id")

result = first.unionByName(second, allowMissingColumns=True).orderBy("order_id")
result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._

val first = orders.filter(col("order_id") <= 1002).select("order_id", "status")
val second = orders.filter(col("order_id") >= 1005).select("order_id")

val result = first.unionByName(second, allowMissingColumns = true).orderBy("order_id")
result.show(false)
example.pyPySpark
Batch oneAppend batch twoCombined rows
Source first batch2 rows
order_idcustomer_idstatustotalitem_count
1001101Delivered89.52
1002102Shipped1493
Second input second batch2 rows
order_idcustomer_idstatustotalitem_count
1001101Delivered89.52
1002102Shipped1493
Stack rows
Result4 rows
order_idcustomer_idstatustotalitem_count
1001101Delivered89.52
1001101Delivered89.52
1002102Shipped1493
1002102Shipped1493
Compare the source columns with the transformed result.

The example is loaded in the editor. Run it as written, then try a small change.

Runs on the Spark backend. First startup may take a moment.

Run your code to see the result.