Skip to content

Storage-partitioned self-join of an Iceberg table returns duplicate rows with partially clustered distribution #6278

Description

@dwsmith1983

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 605051a with spark.comet.exec.enabled=true and spark.comet.scan.icebergNative.enabled=true, on Spark 3.5, 4.0 or 4.1.

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 BatchScans, and Comet runs CometSortMergeJoin over two CometIcebergNativeScans.

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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

area:scanParquet scan / data readingbugSomething isn't workingcorrectnesspriority:criticalData corruption, silent wrong results, security issues

Type

No type

Projects

No projects

    Milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions