andygrove opened a new issue, #6685:
URL: https://github.com/apache/datafusion-comet/issues/6685

   ### Describe the bug
   
   A native shuffle whose child is a Spark-to-Arrow conversion, 
`CometSparkToColumnarExec` or `CometLocalTableScanExec`, fails on a struct 
column from the second batch of a partition:
   
   ```
   org.apache.comet.CometNativeException: C Data interface error: 
java.lang.IllegalArgumentException: no more field nodes for field a: Int(32, 
true) and vector []
        at 
org.apache.arrow.util.Preconditions.checkArgument(Preconditions.java:365)
        at 
org.apache.arrow.vector.VectorLoader.loadBuffers(VectorLoader.java:109)
   ```
   
   The shuffle reads a child that is not a native operator through 
`executeColumnar()` and wraps the batches in `ColumnarBatchArrowReader`, which 
closes each batch once native has it. Both conversions write every batch into 
the same vectors (`RowArrowReader`), and closing a struct vector removes its 
children, so the next batch has none. A native operator over the same 
conversion, such as a partial aggregate, is not affected, because it reads the 
conversion with `doExecuteAsArrowStream()`. Flat and `map<string,string>` 
columns are not affected either.
   
   ### Steps to reproduce
   
   With `spark.comet.convert.rdd.enabled=true` 
(`spark.comet.sparkToColumnar.enabled=true` on 1.0 and 1.1):
   
   ```scala
   import org.apache.spark.sql.Row
   import org.apache.spark.sql.functions.col
   import org.apache.spark.sql.types._
   
   val schema = StructType(Seq(
     StructField("id", IntegerType),
     StructField("inner", StructType(Seq(StructField("a", IntegerType), 
StructField("b", StringType))))))
   val rdd = spark.sparkContext.parallelize((0 until 20000).map(i => Row(i, 
Row(i, s"s$i"))), 1)
   spark.createDataFrame(rdd, schema).repartition(4, col("id")).collect()
   ```
   
   The plan is `CometExchange ... CometNativeShuffle` over 
`CometSparkRowToColumnar` over `Scan ExistingRDD`. It fails with AQE on and 
off, at the default batch size, once a partition holds more than one batch. 
With `spark.comet.exec.localTableScan.enabled=true` and 
`spark.comet.batchSize=16`, `Seq.tabulate(200)(i => (i, (i, 
s"s$i"))).toDF("id", "inner").repartition(2, col("id"))` fails the same way. So 
does the typed Dataset conversion in #6564 
(`spark.comet.convert.typedDataset.enabled`) under a repartition, an orderBy or 
a shuffled join.
   
   ### Expected behavior
   
   The query returns the same rows as Spark.
   
   ### Additional context
   
   Reproduced on main at ba08acd81 with Spark 4.1. The code path dates from 
#4572, so 1.0.0 and branch-1.1 have it too, by reading the code (not run 
there). That makes it not a 1.1.0 regression. Every conversion that reaches it 
is off by default.
   
   #6607 fixes it by reading a `CometNativeArrowSource` child of a native 
shuffle as an Arrow stream, as a native operator does. With only that change 
applied on top of main and #6564, all of the cases above pass.
   


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