RecordBatch field 0 should be integer when using HashPartitioning
- Dominant language
- Scala
- Stars
- 1.6k
- Forks
- 657
- Avg merge
- 2d 14h
- Merged PRs (30d)
- 80
Description
### Backend
VL (Velox)
### Bug description
When I use ColumnarShuffleExchangeExec and ShuffleExchangeExec with HashPartitioning I receive the error below. I tried to isolate the issue and created this repro code. It only fails when partitioning is HashPartitioning. Single and RoundRobin works fines.
I created it based on examples I found in the Gluten source, like https://github.com/oap-project/gluten/blob/81bb6c9b0652ec4df39e6f50c0405756b59d5a3d/gluten-core/src/main/scala/io/glutenproject/execution/TakeOrderedAndProjectExecTransformer.scala#L103
Repro code:
```
val sourcePath = "/tmp/test/source"
val outputPath = "/tmp/test/target"
val pColumn = "colC"
val isTablePartitioned = true
val targetFileSize = 134217728
val numShufflePartitions = spark.sessionState.conf.numShufflePartitions
val dfSource = spark
.range(5000)
.map { _ =>
(10L,
11,
scala.util.Random.nextInt(2))
}
.repartition(100)
.toDF("colA", "colB", "colC")
dfSource.write.partitionBy(pColumn).format("parquet").mode("overwrite").save(sourcePath)
val df = spark.read.format("parquet").load(sourcePath)
val physicalPlan = df.queryExecution.executedPlan
println("\n\n physicalPlan: " + physicalPlan)
val originalPlan = physicalPlan
//Remove top C2R
val planWithoutC2R = physicalPlan.children(0)
val partitionColumnsExpr = Array(planWithoutC2R.output.find(c => c.name.equals(pColumn)).get)
val partitioning = HashPartitioning(partitionColumnsExpr, numShufflePartitions)
// These two works
// SinglePartition
// RoundRobinPartitioning(numShufflePartitions)
val shuffleExec = ShuffleExchangeExec(partitioning, planWithoutC2R)
val transformedShuffleExec =
ColumnarShuffleExchangeExec(shuffleExec, planWithoutC2R, shuffleExec.child.output)
// TransformHints.tag(shuffleExec, transformedShuffleExec.doValidate().toTransformHint)
// Add C2R back
val planWithC2R = VeloxColumnarToRowExec(transformedShuffleExec)
val planToWrite = planWithC2R
println("\n\n planToWrite: " + planToWrite)
println("X: " + planToWrite.execute().first().getLong(0))
```
### Spark version
Spark-3.3.x
### Spark configurations
Running using unit tests
### System information
_No response_
### Relevant logs
```bash
[info] org.apache.spark.SparkException: Job aborted due to stage failure: Task 1 in stage 4.0 failed 1 times, most recent failure: Lost task 1.0 in stage 4.0 (TID 104) (172.31.117.175 executor driver): java.lang.RuntimeException: Exception: VeloxRuntimeError
[info] Error Source: RUNTIME
[info] Error Code: INVALID_STATE
[info] Reason: RecordBatch field 0 should be integer
[info] Retriable: False
[info] Expression: firstChild->type()->isInteger()
[info] Function: getFirstColumn
[info] File: /__w/1/s/Gluten/cpp/velox/shuffle/VeloxShuffleWriter.cc
[info] Line: 96
[info] Stack trace:
[info] # 0 _ZN8facebook5velox7process10StackTraceC1Ei
[info] # 1 _ZN8facebook5velox14VeloxExceptionC1EPKcmS3_St17basic_string_viewIcSt11char_traitsIcEES7_S7_S7_bNS1_4TypeES7_
[info] # 2 _ZN8facebook5velox6detail14veloxCheckFailINS0_17VeloxRuntimeErrorEPKcEEvRKNS1_18VeloxCheckFailArgsET0_
[info] # 3 0x0000000000000000
[info] # 4 _ZN6gluten18VeloxShuffleWriter5splitESt10shared_ptrINS_13ColumnarBatchEEl
[info] # 5 Java_io_glutenproject_vectorized_ShuffleWriterJniWrapper_split
[info] # 6 0x00007f1855018427
[info]
[info] at io.glutenproject.vectorized.ShuffleWriterJniWrapper.split(Native Method)
[info] at org.apache.spark.shuffle.ColumnarShuffleWriter.internalWrite(ColumnarShuffleWriter.scala:163)
[info] at org.apache.spark.shuffle.ColumnarShuffleWriter.write(ColumnarShuffleWriter.scala:218)
[info] at org.apache.spark.shuffle.ShuffleWriteProcessor.write(ShuffleWriteProcessor.scala:59)
[info] at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:99)
[info] at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:52)
[info] at org.apache.spark.scheduler.Task.run(Task.scala:136)
[info] at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:548)
[info] at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1504)
[info] at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:551)
[info] at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
[info] at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
[info] at java.lang.Thread.run(Thread.java:750)
[info]
[info] Driver stacktrace:
[info] at org.apache.spark.scheduler.DAGScheduler.failJobAndIndependentStages(DAGScheduler.scala:2682)
[info] at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2(DAGScheduler.scala:2618)
[info] at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2$adapted(DAGScheduler.scala:2617)
[info] at scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62)
[info] at scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55)
[info] at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49)
[info] at org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:2617)
[info] at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1(DAGScheduler.scala:1190)
[info] at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1$adapted(DAGScheduler.scala:1190)
[info] at scala.Option.foreach(Option.scala:407)
[info] at org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:1190)
[info] at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:2870)
[info] at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2812)
[info] at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2801)
[info] at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49)
[info] at org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:958)
[info] at org.apache.spark.SparkContext.runJob(SparkContext.scala:2344)
[info] at org.apache.spark.SparkContext.runJob(SparkContext.scala:2365)
[info] at org.apache.spark.SparkContext.runJob(SparkContext.scala:2384)
[info] at org.apache.spark.rdd.RDD.$anonfun$take$1(RDD.scala:1477)
[info] at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
[info] at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
[info] at org.apache.spark.rdd.RDD.withScope(RDD.scala:407)
[info] at org.apache.spark.rdd.RDD.take(RDD.scala:1450)
[info] at org.apache.spark.rdd.RDD.$anonfun$first$1(RDD.scala:1491)
[info] at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
[info] at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
[info] at org.apache.spark.rdd.RDD.withScope(RDD.scala:407)
[info] at org.apache.spark.rdd.RDD.first(RDD.scala:1491)
[info] at org.apache.spark.sql.TestGluten2.$anonfun$new$1(TestGluten2.scala:95)
[info] at org.scalatest.OutcomeOf.outcomeOf(OutcomeOf.scala:85)
[info] at org.scalatest.OutcomeOf.outcomeOf$(OutcomeOf.scala:83)
[info] at org.scalatest.OutcomeOf$.outcomeOf(OutcomeOf.scala:104)
[info] at org.scalatest.Transformer.apply(Transformer.scala:22)
[info] at org.scalatest.Transformer.apply(Transformer.scala:20)
[info] at org.scalatest.funsuite.AnyFunSuiteLike$$anon$1.apply(AnyFunSuiteLike.scala:226)
[info] at org.apache.spark.SparkFunSuite.withFixture(SparkFunSuite.scala:203)
[info] at org.scalatest.funsuite.AnyFunSuiteLike.invokeWithFixture$1(AnyFunSuiteLike.scala:224)
[info] at org.scalatest.funsuite.AnyFunSuiteLike.$anonfun$runTest$1(AnyFunSuiteLike.scala:236)
[info] at org.scalatest.SuperEngine.runTestImpl(Engine.scala:306)
[info] at org.scalatest.funsuite.AnyFunSuiteLike.runTest(AnyFunSuiteLike.scala:236)
[info] at org.scalatest.funsuite.AnyFunSuiteLike.runTest$(AnyFunSuiteLike.scala:218)
[info] at org.apache.spark.SparkFunSuite.org$scalatest$BeforeAndAfterEach$$super$runTest(SparkFunSuite.scala:64)
[info] at org.scalatest.BeforeAndAfterEach.runTest(BeforeAndAfterEach.scala:234)
[info] at org.scalatest.BeforeAndAfterEach.runTest$(BeforeAndAfterEach.scala:227)
[info] at org.apache.spark.sql.TestGluten2.org$scalatest$BeforeAndAfterEachTestData$$super$runTest(TestGluten2.scala:41)
[info] at org.scalatest.BeforeAndAfterEachTestData.runTest(BeforeAndAfterEachTestData.scala:213)
[info] at org.scalatest.BeforeAndAfterEachTestData.runTest$(BeforeAndAfterEachTestData.scala:206)
[info] at org.apache.spark.sql.TestGluten2.runTest(TestGluten2.scala:41)
[info] at org.scalatest.funsuite.AnyFunSuiteLike.$anonfun$runTests$1(AnyFunSuiteLike.scala:269)
[info] at org.scalatest.SuperEngine.$anonfun$runTestsInBranch$1(Engine.scala:413)
[info] at scala.collection.immutable.List.foreach(List.scala:431)
[info] at org.scalatest.SuperEngine.traverseSubNodes$1(Engine.scala:401)
[info] at org.scalatest.SuperEngine.runTestsInBranch(Engine.scala:396)
[info] at org.scalatest.SuperEngine.runTestsImpl(Engine.scala:475)
[info] at org.scalatest.funsuite.AnyFunSuiteLike.runTests(AnyFunSuiteLike.scala:269)
[info] at org.scalatest.funsuite.AnyFunSuiteLike.runTests$(AnyFunSuiteLike.scala:268)
[info] at org.scalatest.funsuite.AnyFunSuite.runTests(AnyFunSuite.scala:1563)
[info] at org.scalatest.Suite.run(Suite.scala:1112)
[info] at org.scalatest.Suite.run$(Suite.scala:1094)
[info] at org.scalatest.funsuite.AnyFunSuite.org$scalatest$funsuite$AnyFunSuiteLike$$super$run(AnyFunSuite.scala:1563)
[info] at org.scalatest.funsuite.AnyFunSuiteLike.$anonfun$run$1(AnyFunSuiteLike.scala:273)
[info] at org.scalatest.SuperEngine.runImpl(Engine.scala:535)
[info] at org.scalatest.funsuite.AnyFunSuiteLike.run(AnyFunSuiteLike.scala:273)
[info] at org.scalatest.funsuite.AnyFunSuiteLike.run$(AnyFunSuiteLike.scala:272)
[info] at org.apache.spark.SparkFunSuite.org$scalatest$BeforeAndAfterAll$$super$run(SparkFunSuite.scala:64)
[info] at org.scalatest.BeforeAndAfterAll.liftedTree1$1(BeforeAndAfterAll.scala:213)
[info] at org.scalatest.BeforeAndAfterAll.run(BeforeAndAfterAll.scala:210)
[info] at org.scalatest.BeforeAndAfterAll.run$(BeforeAndAfterAll.scala:208)
[info] at org.apache.spark.SparkFunSuite.run(SparkFunSuite.scala:64)
[info] at org.scalatest.tools.Framework.org$scalatest$tools$Framework$$runSuite(Framework.scala:318)
[info] at org.scalatest.tools.Framework$ScalaTestTask.execute(Framework.scala:513)
[info] at sbt.ForkMain$Run.lambda$runTest$1(ForkMain.java:413)
[info] at java.util.concurrent.FutureTask.run(FutureTask.java:266)
[info] at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
[info] at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
[info] at java.lang.Thread.run(Thread.java:750)
[info] Cause: java.lang.RuntimeException: Exception: VeloxRuntimeError
[info] Error Source: RUNTIME
[info] Error Code: INVALID_STATE
[info] Reason: RecordBatch field 0 should be integer
[info] Retriable: False
[info] Expression: firstChild->type()->isInteger()
[info] Function: getFirstColumn
[info] File: /__w/1/s/Gluten/cpp/velox/shuffle/VeloxShuffleWriter.cc
[info] Line: 96
[info] Stack trace:
[info] # 0 _ZN8facebook5velox7process10StackTraceC1Ei
[info] # 1 _ZN8facebook5velox14VeloxExceptionC1EPKcmS3_St17basic_string_viewIcSt11char_traitsIcEES7_S7_S7_bNS1_4TypeES7_
[info] # 2 _ZN8facebook5velox6detail14veloxCheckFailINS0_17VeloxRuntimeErrorEPKcEEvRKNS1_18VeloxCheckFailArgsET0_
[info] # 3 0x0000000000000000
[info] # 4 _ZN6gluten18VeloxShuffleWriter5splitESt10shared_ptrINS_13ColumnarBatchEEl
[info] # 5 Java_io_glutenproject_vectorized_ShuffleWriterJniWrapper_split
[info] # 6 0x00007f1855018427
[info] at io.glutenproject.vectorized.ShuffleWriterJniWrapper.split(Native Method)
[info] at org.apache.spark.shuffle.ColumnarShuffleWriter.internalWrite(ColumnarShuffleWriter.scala:163)
[info] at org.apache.spark.shuffle.ColumnarShuffleWriter.write(ColumnarShuffleWriter.scala:218)
[info] at org.apache.spark.shuffle.ShuffleWriteProcessor.write(ShuffleWriteProcessor.scala:59)
[info] at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:99)
[info] at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:52)
[info] at org.apache.spark.scheduler.Task.run(Task.scala:136)
[info] at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:548)
[info] at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1504)
[info] at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:551)
[info] at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
[info] at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
[info] at java.lang.Thread.run(Thread.java:750)
```
Contributor guide
Assessment
This issue has not been assessed yet.