Menu

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.

  • Intermediate
  • 16 min read
  • Updated Oct 2026
On this page
  1. Sample data
  2. Structs, arrays and maps
  3. The three complex types
  4. Reading struct fields
  5. Arrays and maps, safely
  6. Higher-order functions: work inside arrays without exploding
  7. Changing structs
  8. Building nested output
  9. JSON strings inside columns
  10. Pitfalls
  11. In interviews
  12. explode and posexplode
  13. What they do
  14. Keeping the position: posexplode
  15. Maps and arrays of structs
  16. Two explodes multiply rows
  17. Re-nesting after exploding
  18. In SQL: LATERAL VIEW
  19. Pitfalls
  20. In interviews
  21. Practice questions
  22. 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.sku on 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 -1 in 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.b must 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_entries and 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: explode drops rows with NULL or empty arrays. Use the _outer variants 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 filter before 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, use get or try_element_at for array positions that may not exist.
  • Higher-order functions (transform, filter, aggregate, exists) keep the original grain and avoid shuffles.
  • withField and dropFields edit structs; struct, to_json and from_json build and parse nested payloads.
  • explode changes the grain and drops NULL or empty arrays; use _outer variants to keep parents and posexplode to keep positions.
  • Two explodes on one row multiply rows; aggregate order-level values before exploding line items.

By DataDank Editorial · Last reviewed Oct 2026 · All examples run on PySpark 4.2.0 in local mode. Behaviour notes call out Spark 4 ANSI-mode differences (out-of-bounds array access raises an error; size(NULL) returns NULL rather than -1).

Progress is saved in this browser only. No account needed.

Search
Filter by type