PySpark courseLesson 9 of 10
PySpark course · Lesson 9 of 10
Nested and Complex Types in PySpark: Structs, Arrays, Maps and explode
Work with nested JSON in PySpark: query structs, arrays and maps, use higher-order functions, flatten with explode and posexplode, and rebuild nested output.
On this page
- Sample data
- Structs, arrays and maps
- The three complex types
- Reading struct fields
- Arrays and maps, safely
- Higher-order functions: work inside arrays without exploding
- Changing structs
- Building nested output
- JSON strings inside columns
- Pitfalls
- In interviews
- explode and posexplode
- What they do
- Keeping the position: posexplode
- Maps and arrays of structs
- Two explodes multiply rows
- Re-nesting after exploding
- In SQL: LATERAL VIEW
- Pitfalls
- In interviews
- Practice questions
- Key takeaways
Event payloads, API responses and document databases deliver data as nested JSON: an order containing a customer object, an array of line items and a bag of attributes. Spark represents these as complex types (structs, arrays and maps) and can query them without flattening first. Knowing when to keep data nested and when to explode it into rows is a daily skill for Data Engineers, and exploding carelessly is a common cause of wrong counts and runaway row volumes.
Sample data
Three orders written as JSON Lines and read with an explicit nested schema. Order 2 has an empty tag list and a missing zip; order 3 has no address, no items and no tags at all.
import json, os, tempfile
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder.master("local[2]").appName("nested").getOrCreate()
spark.sparkContext.setLogLevel("ERROR")
path = os.path.join(tempfile.mkdtemp(), "orders.jsonl")
records = [
{"order_id": 1, "customer": {"name": "asha", "address": {"city": "Pune", "zip": "411001"}},
"items": [{"sku": "A1", "qty": 2, "price": 10.0}, {"sku": "B7", "qty": 1, "price": 25.0}],
"tags": ["gift", "express"], "attrs": {"channel": "web", "coupon": "NEW10"}},
{"order_id": 2, "customer": {"name": "ben", "address": {"city": "Leeds", "zip": None}},
"items": [{"sku": "A1", "qty": 1, "price": 10.0}],
"tags": [], "attrs": {"channel": "app"}},
{"order_id": 3, "customer": {"name": "chen", "address": None},
"items": None, "tags": None, "attrs": {}},
]
with open(path, "w") as f:
for r in records:
f.write(json.dumps(r) + "\n")
schema = """
order_id INT,
customer STRUCT<name: STRING, address: STRUCT<city: STRING, zip: STRING>>,
items ARRAY<STRUCT<sku: STRING, qty: INT, price: DOUBLE>>,
tags ARRAY<STRING>,
attrs MAP<STRING, STRING>
"""
orders = spark.read.schema(schema).json(path)
orders.printSchema()
root
|-- order_id: integer (nullable = true)
|-- customer: struct (nullable = true)
| |-- name: string (nullable = true)
| |-- address: struct (nullable = true)
| | |-- city: string (nullable = true)
| | |-- zip: string (nullable = true)
|-- items: array (nullable = true)
| |-- element: struct (containsNull = true)
| | |-- sku: string (nullable = true)
| | |-- qty: integer (nullable = true)
| | |-- price: double (nullable = true)
|-- tags: array (nullable = true)
| |-- element: string (containsNull = true)
|-- attrs: map (nullable = true)
| |-- key: string
| |-- value: string (valueContainsNull = true)
Structs, arrays and maps
The three complex types
| Type | Holds | DDL | Access |
|---|---|---|---|
| Struct | A fixed set of named fields, each with its own type (like a nested row) | STRUCT<name: STRING, age: INT> |
col("customer.name"), col("customer")["name"] |
| Array | An ordered list of values of one type | ARRAY<STRING> |
get(col, 0) (0-based), element_at(col, 1) (1-based) |
| Map | Key-value pairs; keys of one type, values of one type | MAP<STRING, STRING> |
col("attrs")["channel"] |
They nest freely: an array of structs (items) is the most common shape for line items, events and history records. Use a struct when the field names are known and stable, a map when keys vary per record (free-form attributes, labels), and an array for repeated values.
Parquet, ORC, Avro and Delta all store these types natively, so you can keep nested data nested in your lake. CSV cannot; serialise to JSON strings first.
Reading struct fields
Dot paths reach into structs; "customer.*" expands the top level of a struct into columns. A NULL anywhere along the path gives NULL, not an error:
orders.select("order_id", "customer.name", "customer.address.city",
F.col("customer")["address"]["zip"].alias("zip")).show()
orders.select("order_id", "customer.*").show()
+--------+----+-----+------+
|order_id|name| city| zip|
+--------+----+-----+------+
| 1|asha| Pune|411001|
| 2| ben|Leeds| NULL|
| 3|chen| NULL| NULL|
+--------+----+-----+------+
+--------+----+--------------+
|order_id|name| address|
+--------+----+--------------+
| 1|asha|{Pune, 411001}|
| 2| ben| {Leeds, NULL}|
| 3|chen| NULL|
+--------+----+--------------+
Selecting only customer.name from Parquet reads only that nested column (nested schema pruning), so deep structs do not have to cost more than flat tables.
Arrays and maps, safely
Under Spark 4’s ANSI mode, indexing past the end of an array raises an error instead of returning NULL:
try:
orders.select(F.col("tags")[0]).collect()
except Exception as e:
print(type(e).__name__, str(e).split(".")[0])
ArrayIndexOutOfBoundsException [INVALID_ARRAY_INDEX] The index 0 is out of bounds
Order 2’s empty tag list caused it. Use get (0-based) or try_element_at (1-based), which return NULL for missing positions:
orders.select("order_id",
F.size("tags").alias("n_tags"),
F.get("tags", 0).alias("first_tag"),
F.try_element_at("tags", F.lit(2)).alias("second_tag"),
F.array_contains("tags", "gift").alias("is_gift"),
F.col("items.sku").alias("skus"),
F.col("attrs")["channel"].alias("channel"),
F.col("attrs")["coupon"].alias("coupon"),
F.map_keys("attrs").alias("attr_keys"),
).show(truncate=False)
+--------+------+---------+----------+-------+--------+-------+------+-----------------+
|order_id|n_tags|first_tag|second_tag|is_gift|skus |channel|coupon|attr_keys |
+--------+------+---------+----------+-------+--------+-------+------+-----------------+
|1 |2 |gift |express |true |[A1, B7]|web |NEW10 |[channel, coupon]|
|2 |0 |NULL |NULL |false |[A1] |app |NULL |[channel] |
|3 |NULL |NULL |NULL |NULL |NULL |NULL |NULL |[] |
+--------+------+---------+----------+-------+--------+-------+------+-----------------+
Things to notice:
items.skuon an array of structs returns an array of that field:[A1, B7].- A missing map key returns NULL (no error).
size(NULL)returns NULL in Spark 4 (it returned-1in older versions with ANSI off), so an empty array (0) and a missing one (NULL) can be told apart.
Higher-order functions: work inside arrays without exploding
transform, filter, aggregate, exists and forall apply a lambda to each element and keep the row structure. They run in the JVM, so they are much faster than a Python UDF over the array, and they avoid the shuffle that an explode-then-group round trip needs.
orders.select("order_id",
F.transform("items", lambda i: i["qty"] * i["price"]).alias("line_totals"),
F.aggregate("items", F.lit(0.0), lambda acc, i: acc + i["qty"] * i["price"]).alias("order_total"),
F.filter("items", lambda i: i["price"] > 15).alias("pricey"),
F.exists("items", lambda i: i["sku"] == "B7").alias("has_b7"),
).show(truncate=False)
+--------+------------+-----------+---------------+------+
|order_id|line_totals |order_total|pricey |has_b7|
+--------+------------+-----------+---------------+------+
|1 |[20.0, 25.0]|45.0 |[{B7, 1, 25.0}]|true |
|2 |[10.0] |10.0 |[] |false |
|3 |NULL |NULL |NULL |NULL |
+--------+------------+-----------+---------------+------+
Other useful array functions: array_distinct, array_union, array_intersect, array_except, sort_array/array_sort, flatten, slice, arrays_zip; for maps, map_values, map_entries, map_from_entries, transform_values, map_filter.
Changing structs
Rebuilding a whole struct to change one nested field is error-prone. withField and dropFields change fields in place, even several levels down:
fixed = (orders
.withColumn("customer", F.col("customer").withField("address.country", F.lit("unknown")))
.withColumn("customer", F.col("customer").dropFields("address.zip")))
fixed.select("customer").printSchema()
fixed.select("order_id", "customer").show(truncate=False)
root
|-- customer: struct (nullable = true)
| |-- name: string (nullable = true)
| |-- address: struct (nullable = true)
| | |-- city: string (nullable = true)
| | |-- country: string (nullable = false)
+--------+-----------------------+
|order_id|customer |
+--------+-----------------------+
|1 |{asha, {Pune, unknown}}|
|2 |{ben, {Leeds, unknown}}|
|3 |{chen, NULL} |
+--------+-----------------------+
Order 3 had a NULL address, and withField does not create the parent: the address stays NULL.
Building nested output
struct, array and create_map build complex columns; to_json serialises them, for example to produce a Kafka message value or an API payload:
built = orders.select(
"order_id",
F.struct(F.col("customer.name").alias("name"), F.size("tags").alias("tag_count")).alias("summary"),
F.array(F.lit("x"), F.lit("y")).alias("arr"),
F.create_map(F.lit("source"), F.lit("jsonl")).alias("meta"),
)
built.printSchema()
print(built.select(F.to_json("summary").alias("j")).first().j)
root
|-- order_id: integer (nullable = true)
|-- summary: struct (nullable = false)
| |-- name: string (nullable = true)
| |-- tag_count: integer (nullable = true)
|-- arr: array (nullable = false)
| |-- element: string (containsNull = false)
|-- meta: map (nullable = false)
| |-- key: string
| |-- value: string (valueContainsNull = false)
{"name":"asha","tag_count":2}
JSON strings inside columns
Often the nested data arrives as a JSON string in one column (a Kafka value, a database TEXT field). from_json parses it with a schema; schema_of_json infers a schema from a sample; get_json_object extracts one path without a schema:
raw = spark.createDataFrame([('{"id": 7, "tags": ["a","b"], "extra": {"k": "v"}}',)], "payload STRING")
parsed = raw.select(F.from_json("payload", "id INT, tags ARRAY<STRING>, extra MAP<STRING,STRING>").alias("p"))
parsed.select("p.*").show()
print(raw.select(F.schema_of_json(F.lit('{"id": 7, "tags": ["a","b"]}'))).first()[0])
print(raw.select(F.get_json_object("payload", "$.tags[1]")).first()[0])
+---+------+--------+
| id| tags| extra|
+---+------+--------+
| 7|[a, b]|{k -> v}|
+---+------+--------+
STRUCT<id: BIGINT, tags: ARRAY<STRING>>
b
Parse once with from_json and select the fields you need; calling get_json_object many times re-parses the string for every call. Spark 4 also has a VARIANT type (parse_json, variant_get) for semi-structured data whose shape varies too much for a fixed schema.
Pitfalls
- JSON inference merges every record’s fields into one schema, so a rarely seen field or a type conflict (a number in one record, a string in another) changes the schema. Supply schemas in pipelines.
- Field name case and dots. A field literally named
a.bmust be quoted with backticks:F.col("`a.b`"). - Maps cannot be grouping or join keys in many contexts, and map key order is not guaranteed; convert with
map_entriesand sort if you need stable output. - Very wide structs and deeply nested arrays are hard to query in BI tools. Flatten into a curated table for consumers while keeping the raw nested data in the bronze layer.
In interviews
Typical prompts: “how do you process nested JSON in Spark?”, “struct vs map, when would you use each?”, and “compute the order total from an array of line items”. Strong answers: declare the schema, use dot paths and safe accessors, use higher-order functions such as aggregate instead of exploding when the result stays at order level, and know that explode changes the grain.
explode and posexplode
What they do
A generator function turns one row into zero or more rows. explode(array) produces one row per element; explode(map) produces one row per entry with key and value columns. Exploding changes the grain of the data: from one row per order to one row per order line.
flat = orders.select("order_id", F.explode("items").alias("item"))
flat.select("order_id", "item.*").show()
+--------+---+---+-----+
|order_id|sku|qty|price|
+--------+---+---+-----+
| 1| A1| 2| 10.0|
| 1| B7| 1| 25.0|
| 2| A1| 1| 10.0|
+--------+---+---+-----+
Order 3 disappeared: explode emits no rows for a NULL or empty array, like an inner join. explode_outer keeps the row with NULLs, like a left join:
orders.select("order_id", F.explode_outer("items").alias("item")).select("order_id", "item.sku").show()
+--------+----+
|order_id| sku|
+--------+----+
| 1| A1|
| 1| B7|
| 2| A1|
| 3|NULL|
+--------+----+
Keeping the position: posexplode
posexplode adds the 0-based position of each element, which you need when order carries meaning (the first touchpoint in an attribution path, the step number in a funnel, the line number on an invoice):
orders.select("order_id", F.posexplode("tags").alias("pos", "tag")).show()
orders.select("order_id", F.posexplode_outer("tags").alias("pos", "tag")).show()
+--------+---+-------+
|order_id|pos| tag|
+--------+---+-------+
| 1| 0| gift|
| 1| 1|express|
+--------+---+-------+
+--------+----+-------+
|order_id| pos| tag|
+--------+----+-------+
| 1| 0| gift|
| 1| 1|express|
| 2|NULL| NULL|
| 3|NULL| NULL|
+--------+----+-------+
posexplode_outer kept order 2 (empty array) and order 3 (NULL array) with NULL position and value.
Maps and arrays of structs
orders.select("order_id", F.explode("attrs").alias("key", "value")).show()
orders.select("order_id", F.inline("items")).show()
+--------+-------+-----+
|order_id| key|value|
+--------+-------+-----+
| 1|channel| web|
| 1| coupon|NEW10|
| 2|channel| app|
+--------+-------+-----+
+--------+---+---+-----+
|order_id|sku|qty|price|
+--------+---+---+-----+
| 1| A1| 2| 10.0|
| 1| B7| 1| 25.0|
| 2| A1| 1| 10.0|
+--------+---+---+-----+
inline explodes an array of structs straight into columns, equivalent to explode followed by item.*. (inline_outer keeps empty rows.)
Two explodes multiply rows
Exploding two independent arrays of the same row gives their Cartesian product, not a side-by-side pairing:
orders.select("order_id", F.explode("items").alias("item"), F.explode("tags").alias("tag")) \
.select("order_id", "item.sku", "tag").show()
+--------+---+-------+
|order_id|sku| tag|
+--------+---+-------+
| 1| A1| gift|
| 1| A1|express|
| 1| B7| gift|
| 1| B7|express|
+--------+---+-------+
Order 1 has 2 items and 2 tags, so 4 rows; any sum over qty is now doubled. The same happens if you explode the two arrays into separate DataFrames and join them on order_id. To pair elements by position use arrays_zip(a, b) and explode once; to keep them independent, explode each into its own table at its own grain.
Re-nesting after exploding
The inverse of explode is groupBy with collect_list / collect_set. Note that this round trip costs a shuffle, which is why higher-order functions are preferable when the result stays at the original grain:
items = orders.select("order_id", F.inline("items"))
regrouped = (items.groupBy("order_id")
.agg(F.sum(F.col("qty") * F.col("price")).alias("total"),
F.collect_list(F.struct("sku", "qty")).alias("lines"),
F.sort_array(F.collect_set("sku")).alias("skus")))
regrouped.orderBy("order_id").show(truncate=False)
+--------+-----+------------------+--------+
|order_id|total|lines |skus |
+--------+-----+------------------+--------+
|1 |45.0 |[{A1, 2}, {B7, 1}]|[A1, B7]|
|2 |10.0 |[{A1, 1}] |[A1] |
+--------+-----+------------------+--------+
Order 3 is gone again because the explode dropped it; collect_list order is not guaranteed after a shuffle, so sort the array (or keep posexplode positions and sort by them) when order matters.
In SQL: LATERAL VIEW
orders.createOrReplaceTempView("orders")
spark.sql("""
SELECT o.order_id, i.sku, i.qty
FROM orders o
LATERAL VIEW OUTER explode(o.items) t AS i
ORDER BY o.order_id, i.sku
""").show()
spark.sql("""
SELECT order_id, pos, tag FROM orders
LATERAL VIEW posexplode(tags) t AS pos, tag
""").show()
+--------+----+----+
|order_id| sku| qty|
+--------+----+----+
| 1| A1| 2|
| 1| B7| 1|
| 2| A1| 1|
| 3|NULL|NULL|
+--------+----+----+
+--------+---+-------+
|order_id|pos| tag|
+--------+---+-------+
| 1| 0| gift|
| 1| 1|express|
+--------+---+-------+
LATERAL VIEW OUTER is the SQL form of explode_outer. Generators can also be used directly in the SELECT list, as in the DataFrame examples.
Pitfalls
- Losing rows:
explodedrops rows with NULL or empty arrays. Use the_outervariants when every parent must survive. - Row explosion: very long arrays (thousands of events per user) multiply data volume and shuffle cost. Filter the array with
filterbefore exploding, or aggregate with higher-order functions instead. - Double counting after exploding: order-level values (shipping fee, order total) are repeated on every line. Aggregate them at order level, not line level.
- Multiple explodes produce Cartesian products.
- Lost order: if position matters, use
posexplode.
In interviews
“What is the difference between explode and explode_outer?”, “how do you keep the array index?” and “after exploding line items your revenue doubled, why?” are common. A strong answer mentions grain, NULL and empty arrays, posexplode, Cartesian products from multiple explodes, and higher-order functions as the alternative.
Practice questions
An order has items ARRAY<STRUCT<sku, qty, price>>. How do you compute the order total without exploding?
F.aggregate("items", F.lit(0.0), lambda acc, i: acc + i["qty"] * i["price"]). It stays at order grain, runs in the JVM and avoids the explode-then-group shuffle. Use coalesce if NULL arrays should give 0.
Why did some orders disappear after explode, and how do you keep them?
explode emits no rows when the array is NULL or empty. Use explode_outer (or posexplode_outer, inline_outer, LATERAL VIEW OUTER) to keep one row with NULLs for those parents.
Your job fails on Spark 4 with INVALID_ARRAY_INDEX. What changed?
ANSI mode is on by default, so accessing an array position that does not exist raises an error instead of returning NULL. Use get(col, i) or try_element_at(col, i) where missing elements are expected.
You exploded items and tags in the same select and quantities doubled. Why?
Two generators over independent arrays produce the Cartesian product of their elements for each row. Explode each array into its own table at its own grain, or use arrays_zip if elements are meant to pair up by position.
When would you model a field as a map rather than a struct?
When the keys vary between records or are open-ended (custom attributes, labels, query parameters). A struct is better when field names are known and stable, because each field has its own type, can be pruned on read and is visible in the schema.
How do you change one field deep inside a struct?
Use withField with a dotted path, for example F.col("customer").withField("address.country", F.lit("IN")), and dropFields to remove fields. It avoids rebuilding the struct by hand. It does not create a parent struct that is NULL.
Key takeaways
- Structs hold fixed named fields, arrays hold repeated values, maps hold variable keys; declare them explicitly in schemas.
- Dot paths and
col["key"]read nested values; under ANSI, usegetortry_element_atfor array positions that may not exist. - Higher-order functions (
transform,filter,aggregate,exists) keep the original grain and avoid shuffles. withFieldanddropFieldsedit structs;struct,to_jsonandfrom_jsonbuild and parse nested payloads.explodechanges the grain and drops NULL or empty arrays; use_outervariants to keep parents andposexplodeto keep positions.- Two explodes on one row multiply rows; aggregate order-level values before exploding line items.
Progress is saved in this browser only. No account needed.

