andygrove opened a new issue, #6234:
URL: https://github.com/apache/datafusion-comet/issues/6234
### Describe the bug
When a native plan reads its input from the JVM through
`CometNativeArrowSource.stream`, a Spark exception thrown while producing any
batch after the first reaches the user wrapped in a `CometNativeException`,
instead of as the exception Spark throws.
`reconcileStreamSchema` reads the first batch on the JVM to derive the
stream's schema, so an error in that batch propagates unchanged. Every later
batch is pulled by the native side through the exported stream's `get_next`.
Arrow Java reports the thrown exception to the C Data interface only as the
stream's error string, and arrow-rs turns that into
`ArrowError::CDataInterface`. The user gets:
```text
org.apache.spark.SparkException: Job aborted due to stage failure: ...
<- org.apache.comet.CometNativeException: C Data interface error:
org.apache.spark.SparkArithmeticException: [ARITHMETIC_OVERFLOW] integer
overflow. ...
```
Spark throws `SparkArithmeticException` with the `ARITHMETIC_OVERFLOW`
condition. So the exception class, the error condition and the SQLSTATE all
depend on which batch happens to fail.
### Steps to reproduce
On `main` (646ff181a) with the default Spark 4.1 profile:
1. Set `spark.comet.write.iceberg.splitOperator.enabled=true`,
`spark.comet.iceberg.write.enabled=true` and `spark.sql.ansi.enabled=true`.
2. Create an Iceberg table `src (id INT)` holding 0 to 99999 in a single
data file, for example written with Comet off and `.coalesce(1)`.
3. Create an Iceberg table `t (v INT)` and run `INSERT INTO t SELECT id +
(2147483647 - 90000) FROM src`. The plan has `CometIcebergWrite`, and the
overflow happens at row 90000, well past the first batch.
The error comes back wrapped as above. With `- 100` in place of `- 90000`
the overflow is in the first batch, and the error is the correct
`SparkArithmeticException`. With Comet off, both statements fail with
`SparkArithmeticException`.
### Expected behavior
The same exception class and error condition as Spark, whichever batch
fails. One option is for the JVM side to keep the `Throwable` it caught in the
stream callback and rethrow that when the native side reports the C Data
interface error, rather than relying on the error string.
### Additional context
Found while reviewing #5318. A `MERGE_CARDINALITY_VIOLATION` from its native
`MergeRows` operator under the native Iceberg writer surfaces the same way,
where Spark throws `SparkRuntimeException`. Nothing in the mechanism is
specific to the writer, so other native plans fed through
`CometNativeArrowSource.stream` are probably affected too, but I have only
reproduced it under the native Iceberg writer. #4517 is related but a different
path: there DataFusion wraps a typed error inside a single native plan.
--
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]