[
https://issues.apache.org/jira/browse/SPARK-59609?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated SPARK-59609:
-----------------------------------
Labels: pull-request-available (was: )
> Parquet aggregate pushdown reads statistics from the wrong column after
> schema merging
> --------------------------------------------------------------------------------------
>
> Key: SPARK-59609
> URL: https://issues.apache.org/jira/browse/SPARK-59609
> Project: Spark
> Issue Type: Bug
> Components: Optimizer, SQL
> Affects Versions: 5.0.0
> Environment: {code}
> Spark: 5.0.0-SNAPSHOT
> Commit: 9b4e4f7547c07c5eba3d2b0822a6d3dbb023a150
> Java: OpenJDK 17.0.17
> OS: macOS 15.7.4
> {code}
> Reporter: Yibo Dong
> Priority: Major
> Labels: pull-request-available
>
> h2. What happened
> With Parquet schema merging, aggregate pushdown can read statistics from the
> wrong physical column when one file is missing the aggregated column.
> This causes silently incorrect aggregate results.
> h2. How to reproduce
> Run the following SQL:
> {code:sql}
> CREATE TABLE with_value
> USING PARQUET
> LOCATION '${spark.sql.warehouse.dir}/with_value'
> AS
> SELECT id, value, other
> FROM VALUES
> (1, 10, 100),
> (2, 20, 200)
> AS t(id, value, other);
> CREATE TABLE without_value
> USING PARQUET
> LOCATION '${spark.sql.warehouse.dir}/without_value'
> AS
> SELECT id, other
> FROM VALUES
> (3, 300),
> (4, 400)
> AS t(id, other);
> CREATE OR REPLACE TEMP VIEW t
> USING PARQUET
> OPTIONS (
> path '${spark.sql.warehouse.dir}',
> mergeSchema 'true',
> recursiveFileLookup 'true'
> );
> SELECT id, value, other
> FROM t
> ORDER BY id;
> SELECT count(*) AS rows_all FROM t;
> SELECT count(value) AS rows_value FROM t;
> SELECT min(value) AS lo FROM t;
> SELECT max(value) AS hi FROM t;
> EXPLAIN EXTENDED
> SELECT count(value), min(value), max(value)
> FROM t;
> {code}
> Run it once with Parquet aggregate pushdown enabled:
> {code:bash}
> spark-sql
> --master 'local[2]'
> --conf spark.ui.enabled=false
> --conf spark.sql.adaptive.enabled=false
> --conf spark.sql.shuffle.partitions=1
> --conf spark.sql.sources.useV1SourceList=
> --conf "spark.sql.warehouse.dir=$(mktemp -d)"
> --conf spark.sql.parquet.aggregatePushdown=true
> -f repro.sql
> {code}
> Then run it again with a fresh warehouse and aggregate pushdown disabled:
> {code:bash}
> spark-sql
> --master 'local[2]'
> --conf spark.ui.enabled=false
> --conf spark.sql.adaptive.enabled=false
> --conf spark.sql.shuffle.partitions=1
> --conf spark.sql.sources.useV1SourceList=
> --conf "spark.sql.warehouse.dir=$(mktemp -d)"
> --conf spark.sql.parquet.aggregatePushdown=false
> -f repro.sql
> {code}
> h2. Expected result
> After schema merging, Spark reads the rows as:
> || id || value || other ||
> | 1 | 10 | 100 |
> | 2 | 20 | 200 |
> | 3 | NULL | 300 |
> | 4 | NULL | 400 |
> The file created by {{without_value}} does not contain the {{value}} column,
> so the merged schema should supply {{NULL}} for {{value}} in those rows.
> Both configurations should therefore produce:
> || Aggregate || Result ||
> | COUNT(*) | 4 |
> | COUNT(value) | 2 |
> | MIN(value) | 10 |
> | MAX(value) | 20 |
> h2. Actual result
> With {{spark.sql.parquet.aggregatePushdown=true}}:
> || Aggregate || Result ||
> | COUNT(*) | 4 |
> | COUNT(value) | *4* |
> | MIN(value) | 10 |
> | MAX(value) | *400* |
> With {{spark.sql.parquet.aggregatePushdown=false}}:
> || Aggregate || Result ||
> | COUNT(*) | 4 |
> | COUNT(value) | 2 |
> | MIN(value) | 10 |
> | MAX(value) | 20 |
> With aggregate pushdown enabled, the physical plan reports:
> {code}
> PushedAggregation: [COUNT(value), MIN(value), MAX(value)]
> {code}
> With aggregate pushdown disabled, it reports:
> {code}
> PushedAggregation: []
> {code}
> The pushed-down result counts rows from the file where {{value}} is absent
> and returns {{400}} as {{MAX(value)}}. The value {{400}} belongs to the
> {{other}} column in that file.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]