[
https://issues.apache.org/jira/browse/SPARK-58735?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Eyck Troschke updated SPARK-58735:
----------------------------------
Attachment: test_size_vs_explode.py
> "size(array)" forces full read of nested array columns and prevents column
> pruning for Parquet/Delta arrays
> -----------------------------------------------------------------------------------------------------------
>
> Key: SPARK-58735
> URL: https://issues.apache.org/jira/browse/SPARK-58735
> Project: Spark
> Issue Type: Improvement
> Components: Optimizer
> Affects Versions: 4.1.2
> Environment: Apache Spark 4.1.2
> Databricks Runtime (Photon)
> Reporter: Eyck Troschke
> Priority: Major
> Attachments: test_size_vs_explode.py
>
>
> When applying `size()` to a nested array column stored in Parquet/Delta,
> Spark reads the *entire* array column including all child struct fields, even
> though only array offsets / repetition levels are required to compute the
> array length. This results in significantly higher I/O compared to
> semantically equivalent queries using `explode()` or accessing smaller
> sub‑arrays.
> This behavior occurs in both Apache Spark and Databricks runtimes. It appears
> to be caused by conservative Catalyst planning for nested arrays, which
> prevents column pruning and page skipping.
> *Example:*
> Table size: ~100GB, ~18B rows
> Column: `events` (ARRAY<STRUCT<...>>)
> Sub‑column: `events.flag_a` (small ARRAY)
> ```
> select sum(size(events)); -- 49s, ~100GB read
> select sum(size(events.flag_a)); -- 8s, ~6GB read
> select count(*) from t lateral view explode(events); -- 9s, ~10GB
> read
> select count(*) from t lateral view explode(events.flag_a); -- 8s, ~6GB read
> ```
> Although these queries are semantically equivalent, `size(array)` triggers a
> full scan of all nested array fields, while `explode(array)` allows column
> pruning.
> *Expected Behavior:*
> `size(array)` should only require reading array offsets / repetition levels
> and should not force full read of all nested struct fields.
> *Actual Behavior:*
> Spark reads all nested array child columns, causing large and unnecessary I/O.
> *Impact:*
> Severe performance degradation for large nested arrays in Parquet/Delta
> tables.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]