dwsmith1983 opened a new pull request, #6297: URL: https://github.com/apache/datafusion-comet/pull/6297
## Which issue does this PR close? Closes #6278. ## Rationale for this change When several scans run in one native block, `findAllPlanData` collects each scan's common and per-partition planning data under a key and merges them with `.toMap`. `injectPlanData` and the native shuffle writer then look each scan's data up by the same key. If two scans share a key, both read the last scan's data. Their partition counts usually match, so the length check in `CometExecRDD` does not notice. Neither key identifies a scan node: - The Iceberg key is the metadata location plus `scanHashCode`, the hash of Iceberg's `SparkScan`. That hash leaves out how Spark grouped the scan's partitions for a storage-partitioned join. In a self-join with partially clustered distribution, one side is split and the other replicated, and Comet returned duplicate rows (378 where Spark returns 126 in the issue's example). - The Parquet key is the source string plus a hash of the required schema, data filters, projection and field types. Partition filters and DPP are not in it. Two sides of a bucketed self-join sit in one block with no exchange between them, and when neither side has a pushed data filter (a full outer join, or `spark.sql.parquet.filterPushdown=false`) they collide. Both read the right side's partitions: wrong rows, and 0 rows instead of 40 in a DPP case. With a pushed filter the keys differ only because each serialized filter gets a fresh expression id. ## What changes are included in this PR? - Both keys now include the scan operator's `plan_id`. `CometExecRule` sets it to `op.id`, the id of the `CometScanExec` or `BatchScanExec` being converted. Spark gives every plan-node instance its own `SparkPlan.id`, so two scans in one block never share it. The id is in the scan's proto, which parent blocks embed unchanged, so the driver, the executors and copies of the exec all compute the same key. `PlanDataInjector.withPlanId` is the one place that adds it. - `findAllPlanData` keeps one copy when the same key carries identical data and throws `CometRuntimeException` when the data differs, instead of silently keeping the last entry. - The scaladoc on `PlanDataInjector.getKey` and `CometScanWithPlanData.sourceKey` now says the key must be unique per scan node, for out-of-tree scans that build their own keys. The output attribute ids would not work as the key: `df.filter($"p" === 1).union(df.filter($"p" === 2))` puts two scans with the same output attributes and different files under one `CometUnionExec`. Two scans that really are identical in one block now ship their common data twice, and each common is parsed once per stage instead of being shared. Scans with pushed data filters already had distinct keys, so most plans see no difference. ## How are these changes tested? - `CometIcebergNativeSuite`: a storage-partitioned self-join of an Iceberg table, with partially clustered distribution on and off. It checks the answer against Spark and that the plan has two native Iceberg scans and no shuffle. On main it returns 378 rows where Spark returns 126. - `CometJoinSuite`: a bucketed, partitioned Parquet self-join with AQE on and off, as a full outer join on different partitions and as a DPP join with filter pushdown off. It checks the answer and that both scans sit under the join with no exchange between them, and each DPP scan carries its pruning filter. On main the answers are wrong. On Spark 3.4 the DPP shape runs with AQE off only, because 3.4 keeps V1 scans with AQE DPP in Spark. - `PlanDataInjectorSuite`: different `plan_id`s give different keys and the same `plan_id` gives the same key, for both scan types. - `CometScanWithPlanDataSuite`: `findAllPlanData` keeps one copy of identical data under a key and fails when the data differs. `CometExecSuite`, `CometJoinSuite`, `CometIcebergNativeSuite`, `PlanDataInjectorSuite`, `PlanDataInjectorShuffleLifecycleSuite`, `CometPlanEqualitySuite`, `CometScanWithPlanDataSuite`, `CometNativeShuffleSuite`, `CometCelebornShuffleManagerSuite`, `CometTopKSuite` and the TPC-DS plan stability suites pass on Spark 3.5. The new tests, `CometJoinSuite` and the Iceberg DPP tests pass on Spark 3.4, 4.0 and 4.1. -- 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]
