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]