dwsmith1983 opened a new issue, #6278:
URL: https://github.com/apache/datafusion-comet/issues/6278
### Describe the bug
A storage-partitioned self-join of an Iceberg table returns duplicate rows
with the native Iceberg scan when partially clustered distribution is on. No
rows are lost; some come back several times.
Both sides of the self-join become `CometIcebergNativeScanExec` nodes in one
native block. `findAllPlanData` collects each scan's per-partition data under
`IcebergPlanDataInjector.getKey`, which is the metadata location plus
`scanHashCode`. That hash comes from the Iceberg `SparkScan` (pushed filters,
snapshot, branch, read schema) and does not include how Spark grouped the
scan's partitions for the join. When both sides read the same columns they get
the same key, and `results.flatMap(_._2).toMap` in `findAllPlanData` keeps only
one side's data.
With partially clustered distribution, Spark splits one side (one partition
per file) and replicates the other (every file of the key in each copy). Both
sides still have the same number of partitions, so the length check in
`CometExecRDD` passes, but the replicated side's files are injected into both
scans. Each key's full cross product is then produced once per split.
Other storage-partitioned join shapes I tried match Spark: joins of two
different tables (inner, left, full), identity and bucket partitioning,
disjoint partition values, a filter pruning one side, `bucket(4)` against
`bucket(8)`, and a self-join without partial clustering. A self-join whose two
sides read different columns also matches, because their hashes differ.
### Steps to reproduce
Comet on main 605051ad2 with `spark.comet.exec.enabled=true` and
`spark.comet.scan.icebergNative.enabled=true`, on Spark 3.5, 4.0 or 4.1.
```sql
CREATE TABLE cat.db.l4 (id INT, v STRING) USING iceberg PARTITIONED BY
(bucket(4, id))
TBLPROPERTIES ('read.split.target-size'='1',
'read.split.open-file-cost'='1', 'format-version'='2');
INSERT INTO cat.db.l4 SELECT CAST(id AS INT), concat('l', id) FROM range(0,
40);
INSERT INTO cat.db.l4 SELECT CAST(id AS INT), concat('l2_', id) FROM
range(0, 40, 3);
INSERT INTO cat.db.l4 SELECT CAST(id AS INT), concat('l3_', id) FROM
range(0, 12);
SET spark.sql.sources.v2.bucketing.enabled=true;
SET spark.sql.sources.v2.bucketing.pushPartValues.enabled=true;
SET
spark.sql.sources.v2.bucketing.partiallyClusteredDistribution.enabled=true;
SET spark.sql.requireAllClusterKeysForCoPartition=false;
SET spark.sql.iceberg.planning.preserve-data-grouping=true;
SET spark.sql.autoBroadcastJoinThreshold=-1;
SELECT a.id, a.v, b.v FROM cat.db.l4 a JOIN cat.db.l4 b ON a.id = b.id;
```
The split size properties make each data file its own task, so a bucket has
several input partitions for partial clustering to split. Both plans have no
shuffle: Spark runs `SortMergeJoin` over two `BatchScan`s, and Comet runs
`CometSortMergeJoin` over two `CometIcebergNativeScan`s.
Spark returns 126 rows and Comet returns 378. Every row comes back 3 times,
for example `[8, l8, l3_8]`. A self-join of an identity-partitioned table
(`PARTITIONED BY (k)`, `ON a.k = b.k`) gives 429 rows where Spark gives 252.
The results are the same with AQE on and off. Turning partially clustered
distribution off makes both match.
### Expected behavior
Same rows as Spark. Each scan in a native block should get its own
per-partition data.
### Additional context
A fix would give each scan a key that also identifies the plan node, for
example by including Spark's storage-partitioned join parameters, or by keying
per-partition data per exec rather than by `(metadataLocation, scanHashCode)`.
`CometNativeScanExec` builds its `sourceKey` from the source and a hash of
`NativeScanCommon` in the same way, so a Parquet self-join in one native block
may deserve the same check.
#5342 plans self-join tests for storage-partitioned joins before a grouping
report is turned on by default. This bug happens on main without that feature,
through Spark's own storage-partitioned join planning.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]