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]

Reply via email to