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]

Reply via email to