comphead commented on code in PR #6243:
URL: https://github.com/apache/datafusion-comet/pull/6243#discussion_r4146847329
##########
spark/src/main/scala/org/apache/comet/CometExecIterator.scala:
##########
@@ -225,6 +226,15 @@ class CometExecIterator(
CometExecIterator.startMemoryUsageLog()
+ /**
+ * What the producer behind one of this plan's Arrow stream inputs threw
while native pulled a
+ * batch from it, if anything. See `CometArrowStream.inputFailure`.
+ */
+ private def inputFailure: Option[Throwable] =
+ inputObjects.iterator
+ .collect { case stream: ArrowArrayStream =>
CometArrowStream.inputFailure(stream) }
Review Comment:
After #6372, `main` no longer imports `org.apache.arrow.c.ArrowArrayStream`
in this file. `git` merges the file with `main` without a conflict, so I expect
this line to stop compiling until the import is added back. I haven't built the
merge.
##########
spark/src/main/scala/org/apache/spark/sql/comet/execution/arrow/CometNativeArrowSource.scala:
##########
@@ -248,14 +255,69 @@ object CometArrowStream extends Logging {
}
if (context != null) {
val streamRef = arrowStream
+ exportedReaders.put(streamRef, reader)
context.addTaskCompletionListener[Unit] { _ =>
+ exportedReaders.remove(streamRef)
Review Comment:
On `main`, #6372 adds `streamRef.release()` to this same listener, so the
rebase needs a manual merge here. Would it make sense to keep
`exportedReaders.remove(streamRef)` first, so a throwing `release()` cannot
leave the entry behind in the static map?
##########
spark/src/main/scala/org/apache/comet/CometExecIterator.scala:
##########
@@ -247,6 +257,12 @@ class CometExecIterator(
result
} catch {
+ // A JVM input threw while native pulled a batch from it. Native saw
only the text of that
+ // exception and failed with a CometNativeException built from it, so
rethrow the exception
+ // itself, which is what the task would have thrown without Comet.
+ case _: Throwable if inputFailure.isDefined =>
Review Comment:
Nit: `inputFailure` is evaluated twice here, in the guard and in the body. A
single match (for example a small `unapply` that returns `inputFailure`) would
avoid the repeat. The same explanation of why Arrow Java loses the exception
also appears in `CometArrowStream.inputFailure`, `ffi.md` and each new test, so
would a one-line comment be enough here?
##########
spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala:
##########
@@ -1494,6 +1494,62 @@ class CometIcebergWriteActionSuite
}
}
+ // https://github.com/apache/datafusion-comet/issues/6234. The native writer
pulls its input
+ // through an Arrow C stream, and Arrow Java hands native only the text of
an exception thrown
+ // while producing a batch, so the user got a CometNativeException instead
of Spark's exception.
+ // The first batch is read on the JVM to derive the stream's schema, which
is why the overflow
+ // has to land past it: row 90000 of a single 100000-row data file, so one
task reads it all.
+ test("native acceleration: an input error past the first batch fails like
the JVM writer") {
+ assumeNativeAcceleration()
+ withIcebergCatalog { _ =>
+ withSQLConf(SQLConf.ANSI_ENABLED.key -> "true") {
+ val (nativeTable, jvmTable) = ("overflow_native", "overflow_jvm")
+ Seq("overflow_src", nativeTable, jvmTable).foreach { table =>
+ spark.sql(s"CREATE TABLE $catalog.$ns.$table (v INT) USING iceberg")
+ }
+ withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
+ spark
+ .range(0, 100000, 1, 1)
+ .selectExpr("CAST(id AS INT) AS v")
+ .writeTo(s"$catalog.$ns.overflow_src")
+ .append()
+ }
+ def insert(table: String): Unit =
+ spark.sql(
+ s"INSERT INTO $catalog.$ns.$table " +
+ s"SELECT v + ${Int.MaxValue - 90000} FROM
$catalog.$ns.overflow_src")
+
+ var nativeFailure: Throwable = null
Review Comment:
Nit: this is the third copy of the native-versus-JVM failed-write
scaffolding in this suite (`capturePlans(includeFailures = true)`, the
`CometIcebergWriteExec` check and the JVM `intercept`). Would a small helper
returning `(nativeFailure, jvmFailure)` be worth adding? The
`TakeOrderedAndProjectExec` test in `CometExecSuite` also reaches
`CometArrowStream.inputObjects` with the same overflow, so is this one mainly
here as the end-to-end repro from #6234?
--
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]