From 6e42143179fe8f69fd557034c68129e7c7cadc01 Mon Sep 17 00:00:00 2001 From: linfeng Date: Sun, 9 Aug 2026 20:34:58 +0800 Subject: [PATCH 1/4] fix per-partition job submission in NativeCollectLimit --- .../apache/auron/exec/AuronExecSuite.scala | 29 +++++++++++++++++++ .../auron/plan/NativeCollectLimitBase.scala | 11 +------ 2 files changed, 30 insertions(+), 10 deletions(-) diff --git a/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala b/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala index 69de4d834..5bae58b37 100644 --- a/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala +++ b/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala @@ -24,6 +24,7 @@ import org.apache.spark.sql.catalyst.expressions.WindowExpression import org.apache.spark.sql.catalyst.expressions.WindowSpecDefinition import org.apache.spark.sql.execution.auron.plan.{NativeCollectLimitExec, NativeGlobalLimitExec, NativeLocalLimitExec, NativeTakeOrderedExec} import org.apache.spark.sql.execution.auron.plan.NativeWindowExec +import org.apache.spark.sql.internal.SQLConf import org.apache.auron.BaseAuronSQLSuite import org.apache.auron.util.AuronTestUtils @@ -43,6 +44,34 @@ class AuronExecSuite extends AuronQueryTest with BaseAuronSQLSuite { } } + test("CollectLimit batches partition scans") { + withTempPath { path => + spark.range(8).repartition(8).write.parquet(path.getCanonicalPath) + + withSQLConf( + SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false", + "spark.sql.files.maxPartitionBytes" -> "4096", + "spark.sql.limit.initialNumPartitions" -> "100") { + val df = spark.read.parquet(path.getCanonicalPath).where("id < 0").limit(1) + val collectLimit = collectFirst(df.queryExecution.executedPlan) { + case exec: NativeCollectLimitExec => exec + }.get + val numPartitions = collectLimit.child.execute().getNumPartitions + assert(numPartitions > 1 && numPartitions <= 100) + + val jobGroup = s"collect-limit-${System.nanoTime()}" + spark.sparkContext.setJobGroup(jobGroup, "test CollectLimit job count") + try { + assert(collectLimit.executeCollect().isEmpty) + val jobCount = spark.sparkContext.statusTracker.getJobIdsForGroup(jobGroup).length + assert(jobCount < numPartitions) + } finally { + spark.sparkContext.clearJobGroup() + } + } + } + } + test("CollectLimit with offset") { if (AuronTestUtils.isSparkV34OrGreater) { withTempView("t1") { diff --git a/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeCollectLimitBase.scala b/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeCollectLimitBase.scala index d33661030..4a226fc65 100644 --- a/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeCollectLimitBase.scala +++ b/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeCollectLimitBase.scala @@ -17,7 +17,6 @@ package org.apache.spark.sql.execution.auron.plan import scala.collection.mutable -import scala.collection.mutable.ArrayBuffer import org.apache.spark.OneToOneDependency import org.apache.spark.sql.auron.{NativeHelper, NativeRDD, NativeSupports, Shims} @@ -43,15 +42,7 @@ abstract class NativeCollectLimitBase(limit: Int, offset: Int, override val chil override def executeCollect(): Array[InternalRow] = { val partial = Shims.get.createNativeLocalLimitExec(limit, child) - val buf = new ArrayBuffer[InternalRow] - - // collect rows partition-by-partition up to 'limit', avoiding full-partition collect. - val it = partial.execute().toLocalIterator - while (buf.size < limit && it.hasNext) { - val row = it.next().copy() - buf += row - } - val rows = buf.toArray + val rows = partial.executeTake(limit) if (offset > 0) rows.drop(offset) else rows } From 73ccf3bbf588fa208ff818a6be0f7cc9f2a10a3d Mon Sep 17 00:00:00 2001 From: linfeng Date: Sun, 9 Aug 2026 22:29:07 +0800 Subject: [PATCH 2/4] fix flaky CollectLimit suite on Spark 3.0 --- .../src/test/scala/org/apache/auron/exec/AuronExecSuite.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala b/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala index 5bae58b37..e0ba357f8 100644 --- a/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala +++ b/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala @@ -46,7 +46,7 @@ class AuronExecSuite extends AuronQueryTest with BaseAuronSQLSuite { test("CollectLimit batches partition scans") { withTempPath { path => - spark.range(8).repartition(8).write.parquet(path.getCanonicalPath) + spark.range(0, 80, 1, 8).write.parquet(path.getCanonicalPath) withSQLConf( SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false", From 73d8306a5445d72f427b67c7a9587594160d8430 Mon Sep 17 00:00:00 2001 From: linfeng Date: Mon, 10 Aug 2026 21:40:31 +0800 Subject: [PATCH 3/4] strengthen collect limit job submission suite --- .../test/scala/org/apache/auron/exec/AuronExecSuite.scala | 5 +++-- .../src/test/scala/org/apache/spark/sql/AuronQueryTest.scala | 4 ++++ 2 files changed, 7 insertions(+), 2 deletions(-) diff --git a/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala b/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala index e0ba357f8..c70de33e7 100644 --- a/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala +++ b/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala @@ -51,7 +51,7 @@ class AuronExecSuite extends AuronQueryTest with BaseAuronSQLSuite { withSQLConf( SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false", "spark.sql.files.maxPartitionBytes" -> "4096", - "spark.sql.limit.initialNumPartitions" -> "100") { + "spark.sql.limit.initialNumPartitions" -> "1") { val df = spark.read.parquet(path.getCanonicalPath).where("id < 0").limit(1) val collectLimit = collectFirst(df.queryExecution.executedPlan) { case exec: NativeCollectLimitExec => exec @@ -63,8 +63,9 @@ class AuronExecSuite extends AuronQueryTest with BaseAuronSQLSuite { spark.sparkContext.setJobGroup(jobGroup, "test CollectLimit job count") try { assert(collectLimit.executeCollect().isEmpty) + waitUntilListenerBusEmpty() val jobCount = spark.sparkContext.statusTracker.getJobIdsForGroup(jobGroup).length - assert(jobCount < numPartitions) + assert(jobCount > 0 && jobCount < numPartitions) } finally { spark.sparkContext.clearJobGroup() } diff --git a/spark-extension-shims-spark/src/test/scala/org/apache/spark/sql/AuronQueryTest.scala b/spark-extension-shims-spark/src/test/scala/org/apache/spark/sql/AuronQueryTest.scala index 263f06231..5540b38e3 100644 --- a/spark-extension-shims-spark/src/test/scala/org/apache/spark/sql/AuronQueryTest.scala +++ b/spark-extension-shims-spark/src/test/scala/org/apache/spark/sql/AuronQueryTest.scala @@ -93,4 +93,8 @@ abstract class AuronQueryTest true case _ => false } + + protected def waitUntilListenerBusEmpty(): Unit = { + spark.sparkContext.listenerBus.waitUntilEmpty() + } } From e97ec252ef5da4d3bf8a55acd5d6898250115df8 Mon Sep 17 00:00:00 2001 From: linfeng Date: Tue, 11 Aug 2026 21:19:51 +0800 Subject: [PATCH 4/4] address review comment --- .../src/test/scala/org/apache/auron/exec/AuronExecSuite.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala b/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala index c70de33e7..1251f7174 100644 --- a/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala +++ b/spark-extension-shims-spark/src/test/scala/org/apache/auron/exec/AuronExecSuite.scala @@ -57,7 +57,7 @@ class AuronExecSuite extends AuronQueryTest with BaseAuronSQLSuite { case exec: NativeCollectLimitExec => exec }.get val numPartitions = collectLimit.child.execute().getNumPartitions - assert(numPartitions > 1 && numPartitions <= 100) + assert(numPartitions > 1) val jobGroup = s"collect-limit-${System.nanoTime()}" spark.sparkContext.setJobGroup(jobGroup, "test CollectLimit job count")