Skip to content

Native shuffle fails with "no more field nodes" on a struct column when it reads converted Spark rows directly #6685

Description

@andygrove

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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    area:shuffleShuffle (JVM and native)bugSomething isn't workingpriority:mediumFunctional bugs, performance regressions, broken featuresrequires-triage

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions