Impala can query nested data: ARRAY, MAP and STRUCT columns stored in Parquet or ORC files. That lets you keep a session and its events, or an order and its line items, in one row, rather than splitting them into two tables and joining them on every query. The SQL looks unusual at first. You do not call an explode function. Instead you write the nested collection in the FROM clause as if it were a table, and Impala joins each parent row to its own elements.

This page is about Impala's side of nested data: what it supports, the query syntax, how the planner executes it, the limits that catch people out (Impala cannot write these columns at all), and how to design tables that Spark or Hive produce and Impala serves. For the type system, LATERAL VIEW and building nested values on the Hive side, see Hive complex types. For the file format underneath, see the Parquet format guide.

Parquet / ORC filenested schema, one column per leafHDFS / Ozone / S3 scanonly referenced leaves readRow batchparent row + collection slotSUBPLAN noderuns once per parent rowSINGULAR ROW SRCthe current parent rowUNNESTitems of s.eventsNESTED LOOP JOINparent x items, then filtersrows to aggregation / exchangeper row
A nested query in Impala: the scan carries each row's collection, and a SUBPLAN joins the parent row to its own elements without a shuffle.
Advertisement

What Impala supports, and the hard limits

Complex types arrived in Impala 2.3. The supported surface is narrower than Hive's, and every limit below has caused a production incident somewhere:

RuleConsequence
Only Parquet and ORC tables or partitions can hold complex columnsText, Avro, SequenceFile and RCFile tables with complex columns cannot be queried by Impala
Impala cannot write data files with complex columnsINSERT ... SELECT and CREATE TABLE AS SELECT into complex columns fail; Spark or Hive must produce the files
Complex columns cannot be partition keysPartition on a scalar such as a date
Hive's UNIONTYPE is not supportedTables using it cannot be queried by Impala
Maximum nesting depth is 100 levelsRarely binding, but generated schemas can hit it
Kudu tables do not support complex columnsNested data stays in Parquet, Iceberg or ORC

Impala can run CREATE TABLE with complex columns, so the DDL can live in Impala even though the data cannot be written there. Iceberg tables stored as Parquet carry nested columns too; see Impala with Iceberg for the table-format side.

What Impala supports, and the hard limits

Complex types arrived in Impala 2.3. The supported surface is narrower than Hive's, and every limit below has caused a production incident somewhere:

RuleConsequence
Only Parquet and ORC tables or partitions can hold complex columnsText, Avro, SequenceFile and RCFile tables with complex columns cannot be queried by Impala
Impala cannot write data files with complex columnsINSERT ... SELECT and CREATE TABLE AS SELECT into complex columns fail; Spark or Hive must produce the files
Complex columns cannot be partition keysPartition on a scalar such as a date
Hive's UNIONTYPE is not supportedTables using it cannot be queried by Impala
Maximum nesting depth is 100 levelsRarely binding, but generated schemas can hit it
Kudu tables do not support complex columnsNested data stays in Parquet, Iceberg or ORC

Impala can run CREATE TABLE with complex columns, so the DDL can live in Impala even though the data cannot be written there. Iceberg tables stored as Parquet carry nested columns too; see Impala with Iceberg for the table-format side.

Advertisement

The mental model: collections as joined child tables

The key idea: a collection column behaves like a child table that is already joined to its parent. Referencing it in the FROM clause next to its parent produces one output row per element, with an implicit join on "belongs to this parent row". There is no ON clause because the relationship is physical.

Each kind of complex type exposes its contents through pseudocolumns:

TypeHow you reach the contentsPseudocolumns
STRUCTDot notation directly on the column; no join neededField names
ARRAY of scalarsPut the array in FROMITEM (the value), POS (zero-based position)
ARRAY of STRUCTPut the array in FROM; read fields directlyField names, plus POS
MAPPut the map in FROMKEY, VALUE

Here is the table used through the rest of this page. Spark writes it; Impala only reads it.

CREATE TABLE sessions (
  session_id STRING,
  user_id    BIGINT,
  started_at TIMESTAMP,
  device     STRUCT<os: STRING, app_version: STRING>,
  events     ARRAY<STRUCT<ts: TIMESTAMP, page: STRING, dwell_ms: INT>>,
  attrs      MAP<STRING, STRING>,
  tags       ARRAY<STRING>
)
PARTITIONED BY (dt STRING)
STORED AS PARQUET;

Querying each type

The four basic shapes, one per type:

-- STRUCT: plain dot notation, still one row per session
SELECT s.session_id, s.device.os
FROM sessions s
WHERE s.dt = '2026-10-01' AND s.device.app_version = '5.2.0';

-- ARRAY of scalars: ITEM and POS
SELECT s.session_id, t.pos, t.item
FROM sessions s, s.tags t
WHERE s.dt = '2026-10-01' AND t.item = 'beta';

-- ARRAY of STRUCT: one row per event, fields read directly
SELECT s.session_id, e.pos, e.page, e.dwell_ms
FROM sessions s, s.events e
WHERE s.dt = '2026-10-01' AND e.dwell_ms > 30000;

-- MAP: KEY and VALUE
SELECT s.session_id, a.value AS campaign
FROM sessions s, s.attrs a
WHERE s.dt = '2026-10-01' AND a.key = 'campaign';

Two rules govern where complex columns may appear. In WHERE, GROUP BY, ORDER BY and HAVING you cannot use a complex column by itself; you refer to a scalar inside it, such as e.page or a.key. And the comma join above is an inner join: a session with an empty events array produces no rows at all. When you need those parents, use an outer join against the collection:

-- Keep sessions whose tags array is empty or NULL
SELECT s.session_id, t.item
FROM sessions s LEFT OUTER JOIN s.tags t
WHERE s.dt = '2026-10-01';

Per-parent aggregates are where nested data is strongest. A correlated subquery in the FROM clause runs over just the current row's collection, so there is no GROUP BY across the whole table and no shuffle:

SELECT s.session_id, v.n_events, v.total_dwell_ms, v.n_checkout
FROM sessions s,
     (SELECT count(*)                    AS n_events,
             sum(dwell_ms)               AS total_dwell_ms,
             count(CASE WHEN page = '/checkout' THEN 1 END) AS n_checkout
      FROM s.events) v
WHERE s.dt = '2026-10-01' AND v.n_events >= 5;

Impala 4.1 and later: select-list collections and UNNEST

Until Impala 4.1, only scalars could appear in the select list, so SELECT tags FROM sessions was an error. Impala 4.1 began relaxing that: arrays in the select list for Parquet tables (IMPALA-9498) and structs for ORC tables (IMPALA-9495), with later releases extending the coverage. Clients receive the collection rendered as JSON-like text, so treat it as display output, not as a typed value your application parses for logic. Later releases also allow collections in the select list beside an ORDER BY on other columns, but sorting is rejected when the select list contains a collection nested inside a struct. Check your exact version's documentation before relying on any of this; Cloudera runtimes backport selectively.

Impala 4.1 also added UNNEST, which brings a behaviour the join syntax cannot express: zipping. Joining two arrays of one row gives the cross product. UNNEST over several arrays pairs the i-th elements and pads the shorter arrays with NULL:

-- Two parallel arrays, zipped: 3 rows for arrays of length 3 and 2 (last b is NULL)
SELECT a1.item AS a, a2.item AS b
FROM complextypes_arrays t, UNNEST(t.arr1, t.arr2) AS (a1, a2);

-- Postgres-style form in the select list, same zipping result
SELECT id, UNNEST(arr1), UNNEST(arr2) FROM complextypes_arrays;

Use zipping for the "parallel arrays" schemas some producers emit (timestamps in one array, values in another). Better still, fix the producer to emit an array of structs, which keeps the pairing in the data.

How Impala executes nested queries

Run EXPLAIN on a nested query and you will see a pattern that tells you how it executes. The scan produces parent rows, each carrying a slot that points at its collection. A SUBPLAN node then runs a small plan once per parent row: a SINGULAR ROW SRC that yields the current parent, an UNNEST that yields that parent's elements, and a NESTED LOOP JOIN that combines them and applies predicates. Correlated FROM subqueries add an aggregation inside the subplan.

Three performance facts follow. First, the join is local: no exchange, no hash table and no shuffle, because parent and children are already in the same row batch. That is why nested queries often beat the equivalent two-table join on large data. Second, Parquet stores each scalar leaf as its own column, with repetition and definition levels encoding the nesting, so a query touching only events.page and events.dwell_ms reads those two leaf columns and skips events.ts, the map and everything else. Third, a whole collection is materialized in memory per row, so a few giant arrays (a bot session with two million events) inflate memory per row batch and can trip memory limits. That risk has to be capped upstream; see Impala's architecture for how executors and memory limits interact.

Producing nested tables for Impala

Since Impala cannot write nested columns, the production path is: a batch or streaming job in Spark builds the nested rows and writes Parquet, the table is declared in the metastore, and Impala refreshes its metadata.

from pyspark.sql import functions as F

sessions = (
    events_df                                    # one row per raw event
    .groupBy("dt", "session_id", "user_id")
    .agg(
        F.min("ts").alias("started_at"),
        F.first(F.struct("os", "app_version")).alias("device"),
        F.slice(                                 # cap pathological sessions
            F.sort_array(F.collect_list(F.struct("ts", "page", "dwell_ms"))),
            1, 5000).alias("events"),
        F.map_from_entries(F.collect_set(F.struct("attr_key", "attr_value"))).alias("attrs"),
        F.array_distinct(F.flatten(F.collect_list("tags"))).alias("tags"),
    )
)
(sessions.select("session_id", "user_id", "started_at", "device", "events", "attrs", "tags", "dt")
    .write.mode("overwrite").partitionBy("dt").parquet("s3a://lake/sessions/"))

Then, in Impala, ALTER TABLE sessions RECOVER PARTITIONS picks up new partition directories and REFRESH sessions PARTITION (dt='2026-10-01') reloads the file list for one partition. Two schema traps matter here. Struct field order is part of the contract: by default Impala resolves Parquet columns by position (the PARQUET_FALLBACK_SCHEMA_RESOLUTION query option defaults to POSITION), so inserting a new struct field in the middle shifts every later field and silently returns the wrong data. Append new fields at the end, or set the option to NAME for tables whose producers do not guarantee order. Also, map_from_entries fails on duplicate keys under Spark's default settings, so deduplicate keys before building maps.

Worked example: a session funnel without a shuffle

Suppose the app produces 200 million sessions a day with an average of 25 events each. Flattened, that is five billion event rows a day, each repeating the session id, user id, device and attributes. Nested, it is 200 million rows, and the session-level columns are stored once per session. Parquet's encodings already compress repeated values well, so the storage saving is smaller than 25 times, but the scan and join savings are real: the funnel question below never shuffles.

-- Share of sessions that reached checkout within 10 minutes of start, by OS
SELECT s.device.os,
       count(*)                                    AS sessions,
       sum(CASE WHEN v.reached > 0 THEN 1 ELSE 0 END) / count(*) AS checkout_rate
FROM sessions s,
     (SELECT count(*) AS reached
      FROM s.events e
      WHERE e.page = '/checkout'
        AND e.ts < s.started_at + INTERVAL 10 MINUTES) v
WHERE s.dt = '2026-10-01'
GROUP BY s.device.os;

The subquery references the parent's started_at, which is allowed because the subplan sees the current parent row. The only shuffle is the final GROUP BY on a handful of OS values. The flat design would need a join of five billion events to 200 million sessions, or a window function over the event table. For BI tools that cannot write the collection syntax, publish a view that flattens one collection (CREATE VIEW session_events AS SELECT s.session_id, s.dt, e.ts, e.page FROM sessions s, s.events e), and keep the nested table for engineers.

Failure modes

SymptomCauseFix
AnalysisException naming a complex column in WHERE or GROUP BYThe whole column was referenced instead of a scalar inside itReference a field, ITEM, KEY or VALUE
CTAS or INSERT fails on a nested tableImpala cannot write complex columnsProduce the data with Spark or Hive
Sessions vanish from a countInner comma join drops parents with empty or NULL collectionsLEFT OUTER JOIN the collection, or use a correlated subquery
Fields return values from the wrong field after a schema changePosition-based Parquet resolution and a struct field inserted mid-listAppend fields at the end, or set resolution to NAME
Query exceeds memory limit on a few partitionsA handful of huge collections materialized per rowCap collection sizes in the producer; isolate outliers
Two arrays produce far more rows than expectedTwo collection joins on one row give the cross productUse zipping UNNEST, or redesign as an array of structs

Trade-offs: nested or flat

Nest when the child is owned by one parent, is read together with it, and has a bounded size: events of a session, line items of an order, labels of a resource. Stay flat when children are queried on their own across parents, updated independently, or unbounded, or when the main consumers are BI tools that only understand flat tables. Remember the asymmetry too: Hive and Spark can read and write the nested table, while Impala can only read it, so every correction to nested data goes back through the producer job.

What to do next

  1. Confirm your Impala version and which select-list and UNNEST features it actually has; test on a small table.
  2. Pick one parent-child pair that you currently join on every query and model it as an array of structs.
  3. Write the producer in Spark with a cap on collection size and explicit, append-only struct field order.
  4. Rewrite the top three queries using collection joins and correlated FROM subqueries; compare the plans with EXPLAIN.
  5. Decide on POSITION or NAME schema resolution per table and document it beside the DDL.
  6. Publish flattened views for BI users, and alert on partitions whose maximum collection size spikes.
Key takeaway: Impala reads ARRAY, MAP and STRUCT columns from Parquet and ORC but never writes them, so Spark or Hive must produce the data. Query collections by placing them in the FROM clause, reach contents through fields, ITEM, POS, KEY and VALUE, use outer joins to keep empty parents, and use correlated FROM subqueries for per-row aggregates that need no shuffle. Cap collection sizes upstream, append struct fields at the end, and check your version before relying on select-list collections or UNNEST.