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]