andygrove opened a new pull request, #6243:
URL: https://github.com/apache/datafusion-comet/pull/6243

   ## Which issue does this PR close?
   
   Closes #6234.
   
   ## Rationale for this change
   
   A native plan pulls each batch of a JVM input through the Arrow C stream 
that `CometArrowStream.stream` exports. When the JVM producer throws while 
building a batch, Arrow Java's exported stream 
(`ExportedArrayStreamPrivateData.getNext`) catches the throwable and hands 
native only its `printStackTrace` text. Native then fails the plan with a 
`CometNativeException` built from that text. So the exception class, error 
condition and SQLSTATE the user sees depend on which batch failed.
   
   The issue reports this for the native Iceberg writer, but every JVM input to 
a native plan goes through the same stream, and the bug reaches further than 
the issue says:
   
   - Inputs of Arrow-backed batches go through `wrapColumnarBatchRDD` or 
`inputObjects`. That covers the native Iceberg and Parquet writers, 
`CometTakeOrderedAndProjectExec`, the native limit, and broadcast and 
non-direct shuffle inputs. These are affected only past the first batch, 
because `reconcileStreamSchema` reads the first batch on the JVM.
   - Row inputs (`CometSparkToColumnarExec`, `CometLocalTableScanExec`) export 
their readers directly, so a failure in any batch was wrapped, including the 
first.
   
   ## What changes are included in this PR?
   
   - `CometArrowStream.stream` wraps every exported reader in a 
`FailureRecordingReader`. The wrapper forwards every call to the reader and 
keeps the first throwable that `loadNextBatch` throws. It is registered under 
its exported `ArrowArrayStream` and removed in the stream's existing 
task-completion listener.
   - When native execution fails, `CometExecIterator.getNextBatch` checks its 
input streams and rethrows a recorded throwable instead of the native error. 
For a nested plan, such as the writer's input, that is the typed exception the 
upstream `CometExecIterator` already produced.
   - `ffi.md` describes the path.
   
   There are no native changes.
   
   ## How are these changes tested?
   
   - `CometExecIteratorLifecycleSuite`: a native limit plan reads an input that 
throws on its second batch. The test checks that the plan fails with that exact 
exception instance, and that completing the task drops the recorded failure.
   - `CometExecSuite`:
     - `TakeOrderedAndProjectExec` over a single 100,000-row file, overflowing 
at row 90,000 under ANSI, matches Spark's exception class, error condition and 
SQLSTATE through `checkSparkError`.
     - `CometSparkToColumnarExec` over an RDD that throws in its first batch 
surfaces the same exception as Spark, at the same depth.
   - `CometIcebergWriteActionSuite`: the issue's repro, an ANSI overflow at row 
90,000 through the native writer, must fail like the JVM writer, with 
`ARITHMETIC_OVERFLOW` and nothing committed.
   - All four fail on main with `C Data interface error`. Removing either the 
rethrow or the registry removal makes the unit test fail.
   - The four tests pass on Spark 3.4, 3.5, 4.0, 4.1 and 4.2. The Iceberg test 
is canceled on 4.2, which has no Iceberg on its classpath.
   - The following full suites also pass:
     - On 4.1: `CometExecSuite`, `CometIcebergWriteActionSuite`, 
`CometNativeShuffleSuite`, `CometParquetWriterSuite`, `CometInMemoryCacheSuite` 
and `CometJoinSuite`.
     - On 3.4: `CometExecSuite` and `CometIcebergWriteActionSuite`.
   - scalafix (on 3.5 and 4.0), scalastyle, spotless and prettier are clean.
   


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