Lesson reference: PySpark and Scala
Expand arrays
The setup makes two quantities per order, except cancelled order 1003 has an empty array. explode creates one row per element, repeating order_id. It drops empty and null arrays; explode_outer instead preserves those source rows with a null element.
PySpark example
from pyspark.sql import Window
from pyspark.sql.functions import (
col, lit, array, when, explode, explode_outer, posexplode,
create_map, to_json, struct, from_json, get_json_object, count, row_number
)
array_orders = orders.select(
"order_id", when(col("status") == "Cancelled", array().cast("array<int>"))
.otherwise(array(col("item_count"), col("item_count") + 1)).alias("quantities")
)
result = array_orders.select(
col("order_id"), explode(col("quantities")).alias("quantity")
).orderBy("order_id", "quantity")
result.show(truncate=False)
Scala example
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
val array_orders = orders.select(
col("order_id"), when(col("status") === "Cancelled", array().cast("array<int>"))
.otherwise(array(col("item_count"), col("item_count") + 1)).alias("quantities")
)
val result = array_orders.select(
col("order_id"), explode(col("quantities")).alias("quantity")
).orderBy("order_id", "quantity")
result.show(false)
Keep element positions
posexplode returns both the zero-based position and the element value. Positions restart for each source array. Its indexing differs from element_at, whose first index is 1. Empty arrays produce no rows in this example.
PySpark example
from pyspark.sql import Window
from pyspark.sql.functions import (
col, lit, array, when, explode, explode_outer, posexplode,
create_map, to_json, struct, from_json, get_json_object, count, row_number
)
array_orders = orders.select(
"order_id", when(col("status") == "Cancelled", array().cast("array<int>"))
.otherwise(array(col("item_count"), col("item_count") + 1)).alias("quantities")
)
result = array_orders.select(
"order_id", posexplode("quantities").alias("position", "quantity")
).orderBy("order_id", "position")
result.show(truncate=False)
Scala example
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
val array_orders = orders.select(
col("order_id"), when(col("status") === "Cancelled", array().cast("array<int>"))
.otherwise(array(col("item_count"), col("item_count") + 1)).alias("quantities")
)
val result = array_orders.select(
col("order_id"), posexplode(col("quantities")).as(Seq("position", "quantity"))
).orderBy("order_id", "position")
result.show(false)
Filter expanded rows
Once an array is expanded, a DataFrame filter removes individual element rows. Here quantities at least 2 remain. An order disappears from the output if none of its expanded elements pass; this differs from keeping an empty array in the source row.
PySpark example
from pyspark.sql import Window
from pyspark.sql.functions import (
col, lit, array, when, explode, explode_outer, posexplode,
create_map, to_json, struct, from_json, get_json_object, count, row_number
)
array_orders = orders.select(
"order_id", when(col("status") == "Cancelled", array().cast("array<int>"))
.otherwise(array(col("item_count"), col("item_count") + 1)).alias("quantities")
)
result = array_orders.select(
col("order_id"), explode(col("quantities")).alias("quantity")
).filter(col("quantity") >= 2).orderBy("order_id", "quantity")
result.show(truncate=False)
Scala example
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
val array_orders = orders.select(
col("order_id"), when(col("status") === "Cancelled", array().cast("array<int>"))
.otherwise(array(col("item_count"), col("item_count") + 1)).alias("quantities")
)
val result = array_orders.select(
col("order_id"), explode(col("quantities")).alias("quantity")
).filter(col("quantity") >= 2).orderBy("order_id", "quantity")
result.show(false)