andygrove commented on code in PR #6814:
URL: https://github.com/apache/datafusion-comet/pull/6814#discussion_r4237757953
##########
spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala:
##########
@@ -1283,14 +1283,19 @@ object CometIcebergNativeScan extends
CometOperatorSerde[CometBatchScanExec] wit
// Spark's BatchScanExec.inputRDD returns
sparkContext.parallelize(empty, 1) when
// DPP filtering removes all input partitions. That
ParallelCollectionRDD is the only
Review Comment:
Since this comment is being rewritten anyway, could the first sentence say
exactly when Spark takes this path? `BatchScanExec.inputRDD` only returns
`parallelize(empty, 1)` when DPP removes every split of a scan that reports
`SinglePartition`. Only Spark 3.4 reports that, and only for a scan with one
input partition. On 3.5+, and on 3.4 with more than one split, a fully pruned
scan comes back as a `DataSourceRDD` with no partitions and goes through the
branch above. Wording like "when DPP removes the only split of a scan that
reports SinglePartition, which only Spark 3.4 does" would make it clear this
branch is the 3.4 case.
##########
spark/src/test/scala/org/apache/comet/CometIcebergNativeSuite.scala:
##########
@@ -5150,6 +5152,144 @@ class CometIcebergNativeSuite
}
}
+ Seq(false, true).foreach { aqeEnabled =>
Review Comment:
This also fixes a second symptom of #6667 that the tests don't cover yet. On
Spark 3.4 a sort-merge join between two single-split Iceberg tables has no
exchange either, since both scans report `SinglePartition`, so the two scans
run in one native block. When DPP prunes one side, main fails the query instead
of returning wrong rows. With the pruned scan as the first input it fails with
`All per-partition arrays must have length 0`, and as the second input with
`Missing planning data for key`. I checked locally on 3.4.3 with AQE on and
off, and this PR fixes both. Could you add a test for that shape next to this
one? A second single-file table partitioned by `store_id` and a query like the
one below would do it, run once with the DPP key on `a.store_id` and once on
`b.store_id`.
```sql
SELECT /*+ BROADCAST(d), MERGE(a, b) */ COUNT(*), SUM(b.v)
FROM fact a JOIN fact2 b ON a.id = b.id
JOIN single_split_dim d ON a.store_id = d.store_id
WHERE d.country = 'X'
```
--
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]