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):
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 ba08acd 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.
Describe the bug
A native shuffle whose child is a Spark-to-Arrow conversion,
CometSparkToColumnarExecorCometLocalTableScanExec, fails on a struct column from the second batch of a partition:The shuffle reads a child that is not a native operator through
executeColumnar()and wraps the batches inColumnarBatchArrowReader, 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 withdoExecuteAsArrowStream(). Flat andmap<string,string>columns are not affected either.Steps to reproduce
With
spark.comet.convert.rdd.enabled=true(spark.comet.sparkToColumnar.enabled=trueon 1.0 and 1.1):The plan is
CometExchange ... CometNativeShuffleoverCometSparkRowToColumnaroverScan ExistingRDD. It fails with AQE on and off, at the default batch size, once a partition holds more than one batch. Withspark.comet.exec.localTableScan.enabled=trueandspark.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 ba08acd 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
CometNativeArrowSourcechild 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.