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]

Reply via email to