Skip to content

Commit 1081a3f

Browse files
committed
Fix test error
1 parent 52c7616 commit 1081a3f

File tree

1 file changed

+5
-3
lines changed

1 file changed

+5
-3
lines changed

sql/core/src/test/scala/org/apache/spark/sql/execution/ExchangeCoordinatorSuite.scala

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -264,7 +264,7 @@ class ExchangeCoordinatorSuite extends SparkFunSuite with BeforeAndAfterAll {
264264
.setMaster("local[*]")
265265
.setAppName("test")
266266
.set("spark.ui.enabled", "false")
267-
.set(SQLConf.SHUFFLE_PARTITIONS.key, "5")
267+
.set(SQLConf.SHUFFLE_MAX_NUM_POSTSHUFFLE_PARTITIONS.key, "5")
268268
.set(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key, "true")
269269
.set(SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key, "-1")
270270
.set(
@@ -484,8 +484,10 @@ class ExchangeCoordinatorSuite extends SparkFunSuite with BeforeAndAfterAll {
484484
val df = spark.range(1).selectExpr("id AS key", "id AS value")
485485
val resultDf = df.join(df, "key").join(df, "key")
486486
val sparkPlan = resultDf.queryExecution.executedPlan
487-
assert(sparkPlan.collect { case p: ReusedExchangeExec => p }.length == 1)
488-
assert(sparkPlan.collect { case p @ ShuffleExchangeExec(_, _, Some(c)) => p }.length == 3)
487+
val queryStageInputs = sparkPlan.collect { case p: ShuffleQueryStageInput => p }
488+
assert(queryStageInputs.length === 3)
489+
assert(queryStageInputs(0).childStage === queryStageInputs(1).childStage)
490+
assert(queryStageInputs(1).childStage === queryStageInputs(2).childStage)
489491
checkAnswer(resultDf, Row(0, 0, 0, 0) :: Nil)
490492
}
491493
withSparkSession(test, 4, None)

0 commit comments

Comments
 (0)