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


##########
spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala:
##########
@@ -1256,6 +1268,35 @@ case class CometExecRule(session: SparkSession)
   private def hasEnabledHandler(op: SparkPlan): Boolean =
     allExecs.get(op.getClass).exists(_.enabledConfig.forall(_.get(op.conf)))
 
+  /**
+   * Converts the rows a typed Dataset operation produces to Arrow, so the 
operators above it can
+   * run natively. See [[CometConf.COMET_CONVERT_FROM_TYPED_DATASET_ENABLED]].
+   *
+   * Spark inserts the columnar transitions after this rule, but it does not 
look below a
+   * `RowToColumnarTransition` such as `CometSparkToColumnarExec`. That is 
harmless above a leaf.
+   * Here the typed operation's own operators sit below the conversion, and 
without a transition
+   * they would read a Comet child through `CometExec.doExecute`, Spark's 
interpreted
+   * columnar-to-row path. So the subtree gets its transitions now, from 
Spark's own rule, and
+   * `EliminateRedundantTransitions` later replaces each one over a Comet 
child with Comet's own.
+   * Spark's rule leaves existing transitions alone, which matters because 
this rule runs over the
+   * same plan twice under AQE.
+   */
+  private def convertTypedDatasetOutput(op: SerializeFromObjectExec): 
SparkPlan = {
+    val unsupported = op.output.filterNot(a =>
+      CometSparkToColumnarExec.isTypeSupported(a.dataType, a.name, 
ListBuffer.empty))
+    if (unsupported.nonEmpty) {
+      withFallbackReason(
+        op,
+        "Comet cannot convert the output of a typed Dataset operation to Arrow 
because it does " +
+          "not support the type of these columns: " +
+          unsupported.map(a => s"${a.name}: 
${a.dataType.simpleString}").mkString(", "))
+    } else {
+      val withTransitions =
+        ApplyColumnarRulesAndInsertTransitions(Seq.empty, outputsColumnar = 
false).apply(op)
+      convertToComet(withTransitions, 
CometSparkToColumnarExec).getOrElse(withTransitions)

Review Comment:
   [P2] Preserve short-circuiting when a limit consumes typed output. With 
`spark.comet.convert.typedDataset.enabled=true`, `map(...).limit(1)` fills an 
Arrow batch before returning its first row. A user function that throws on row 
30 therefore fails the query, although Spark and conversion-disabled Comet 
return the first row successfully. The resulting plan is `CometCollectLimit -> 
CometSparkRowToColumnar -> SerializeFromObject`. Please preserve row-level 
limiting before batching where valid, or decline conversion for this pipeline, 
and add a regression test.
   
   Evidence: Reproduced at the reviewed head in CometTestBase on Spark 
4.1.3/JDK 17, with AQE both false and true: `spark.range(0, 100, 1, 1).map { i 
=> if (i == 30L) throw new IllegalArgumentException("unexpected evaluation of 
row 30"); i + 1L }.toDF().limit(1).collect()`. Spark and Comet with typed 
conversion disabled return `[Row(1)]`. Enabling conversion throws 
`SparkException` caused by that `IllegalArgumentException`, with 
`RowArrowReader.loadNextBatch` in the stack. Spark’s 
`CollectLimitExec.executeCollect` calls `child.executeTake(limit)`, whereas the 
inserted Arrow reader consumes a batch before the limit can stop it. 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