sunchao commented on code in PR #6564:
URL: https://github.com/apache/datafusion-comet/pull/6564#discussion_r4178740950


##########
spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala:
##########
@@ -473,6 +498,23 @@ case class CometExecRule(session: SparkSession)
       case op if shouldApplySparkToColumnar(conf, op) =>
         convertToComet(op, CometSparkToColumnarExec).getOrElse(op)
 
+      // Typed Dataset operations (`map`, `flatMap`, `mapPartitions`, 
`mapGroups`, ...) pass JVM
+      // objects between their operators, so those stay on Spark. Each of them 
ends in
+      // `SerializeFromObjectExec`, though, whose output is ordinary rows, and 
converting those to
+      // Arrow lets the operators above the typed operation run natively.
+      case op: SerializeFromObjectExec
+          if CometConf.COMET_CONVERT_FROM_TYPED_DATASET_ENABLED.get(conf) =>
+        if (op
+            
.getTagValue(CometExecRule.SKIP_TYPED_DATASET_CONVERSION_UNDER_LIMIT)
+            .isDefined) {
+          withFallbackReason(
+            op,
+            "Comet does not convert the output of a typed Dataset operation 
below a limit " +
+              "because Arrow batching could evaluate rows beyond Spark's 
row-level limit")
+        } else {
+          convertTypedDatasetOutput(op)

Review Comment:
   [P2] Preserve lazy consumption before downstream `MapPartitionsExec` 
consumers. The physical-limit guard does not cover iterator operations such as 
`mapPartitions(_.take(1))`. With typed conversion enabled, an intervening 
native filter batches the earlier typed output and evaluates row 30 before 
returning the first result. A user function throwing on row 30 therefore fails 
a query that succeeds in Spark and conversion-disabled Comet. Please decline 
conversion across these lazy iterator boundaries unless full consumption is 
guaranteed, or otherwise preserve row-level consumption, and add this 
regression case.
   
   Evidence: Reproduced at this head on Spark 4.1.3/JDK 17: `spark.range(0, 
100, 1, 1).map { i => if (i == 30L) throw new 
IllegalArgumentException("unexpected evaluation of row 30"); i + 1L 
}.filter(col("value") > 0L).mapPartitions(_.take(1)).toDF().collect()`. With 
AQE both false and true, Spark and Comet with 
`spark.comet.convert.typedDataset.enabled=false` return `[Row(1)]`. Setting it 
to true throws `SparkException` caused by `CometNativeException` wrapping that 
`IllegalArgumentException`. The failing plan contains `MapPartitions -> 
DeserializeToObject -> CometColumnarToRow -> CometFilter -> 
CometSparkRowToColumnar -> SerializeFromObject`, with no physical limit node. 
Spark’s `MapPartitionsExec` passes a lazy mapped iterator to the user function, 
whereas `RowArrowReader.loadNextBatch` consumes upstream rows to fill a batch. 
The six-configuration probe failed only in the two conversion-enabled cases.



-- 
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