This is an automated email from the ASF dual-hosted git repository.
acvictor pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gluten.git
The following commit(s) were added to refs/heads/main by this push:
new 8b85863408 [CORE] Share taskIdMapsForShuffle between
ColumnarShuffleManager and its block resolver (#13072)
8b85863408 is described below
commit 8b8586340840d7d588819b2cfb4b8e53565ca2db
Author: Ankita Victor <[email protected]>
AuthorDate: Wed Sep 23 15:05:14 2026 +0530
[CORE] Share taskIdMapsForShuffle between ColumnarShuffleManager and its
block resolver (#13072)
ColumnarShuffleManager built its IndexShuffleBlockResolver with the
single-arg constructor, so the resolver allocated its own taskIdMapsForShuffle
while the manager kept a second, separate map. The resolver records blocks
migrated in during executor decommissioning into its map, but unregisterShuffle
reads the manager's map to delete map output. With two maps, migrated blocks
were recorded where nothing read them, so their files were never deleted and
leaked disk on decommissioned exe [...]
taskIdMapsForShuffle is now declared before shuffleBlockResolver -- order
matters, since Scala initializes vals in declaration order and the previous
ordering would capture null -- and passed to the resolver. unregisterShuffle
now iterates under mapTaskIds.synchronized, matching Spark, because the
block-migration path mutates the same set under that lock. stop() now defers to
super.stop(), which stops the resolver.
Adds a ColumnarShuffleManagerSuite case that inserts an entry through the
resolver's map and asserts unregisterShuffle clears it, which fails if the two
maps are ever unshared again.
---
.../shuffle/sort/ColumnarShuffleManager.scala | 37 +++++++++++++++++++---
.../shuffle/sort/ColumnarShuffleManagerSuite.scala | 20 ++++++++++++
2 files changed, 52 insertions(+), 5 deletions(-)
diff --git
a/gluten-substrait/src/main/scala/org/apache/spark/shuffle/sort/ColumnarShuffleManager.scala
b/gluten-substrait/src/main/scala/org/apache/spark/shuffle/sort/ColumnarShuffleManager.scala
index 85c860bc0a..b86b4f9398 100644
---
a/gluten-substrait/src/main/scala/org/apache/spark/shuffle/sort/ColumnarShuffleManager.scala
+++
b/gluten-substrait/src/main/scala/org/apache/spark/shuffle/sort/ColumnarShuffleManager.scala
@@ -42,6 +42,12 @@ import scala.collection.JavaConverters._
* buffers deserialized rows; the row-based branches below produce the same
handles and the same
* writers as `SortShuffleManager`, so the copy is pure overhead here. Keeping
the subtype
* relationship lets Spark take the zero-copy path.
+ *
+ * WARNING: the `shuffleBlockResolver` wiring and the `registerShuffle` /
`getWriter` /
+ * `unregisterShuffle` bodies below are copied from Spark's own
`SortShuffleManager` rather than
+ * inherited, and Spark has changed them before without this copy being
updated. They must be
+ * re-checked against `SortShuffleManager` whenever a new Spark version is
supported, and any
+ * divergence either mirrored here or handled through a shim.
*/
class ColumnarShuffleManager(conf: SparkConf)
extends SortShuffleManager(conf)
@@ -51,11 +57,27 @@ class ColumnarShuffleManager(conf: SparkConf)
import ColumnarShuffleManager._
private lazy val shuffleExecutorComponents =
loadShuffleExecutorComponents(conf)
- override val shuffleBlockResolver = new IndexShuffleBlockResolver(conf)
- /** A mapping from shuffle ids to the number of mappers producing output for
those shuffles. */
+ /**
+ * A mapping from shuffle ids to the task ids of mappers producing output
for those shuffles.
+ *
+ * Must be declared before `shuffleBlockResolver`: Scala initializes vals in
declaration order, so
+ * the resolver would otherwise capture `null`.
+ */
private[this] val taskIdMapsForShuffle = new ConcurrentHashMap[Int,
OpenHashSet[Long]]()
+ // Mirrors SortShuffleManager: the resolver must share this map rather than
allocate its own. It
+ // records blocks migrated in during executor decommissioning, and
`unregisterShuffle` reads the
+ // same map to delete the corresponding map output.
+ //
+ // The argument is positional rather than named because the constructor
signature differs across
+ // supported Spark versions: 3.4 and 3.5 default both `_blockManager` and
`taskIdMapsForShuffle`,
+ // 4.0 drops the defaults, and 4.1 narrows the map type from `java.util.Map`
to
+ // `java.util.concurrent.ConcurrentMap`. The positional form compiles
against all of them, but it
+ // is signature-sensitive -- re-verify it when adding a new Spark version.
+ override val shuffleBlockResolver =
+ new IndexShuffleBlockResolver(conf, null, taskIdMapsForShuffle)
+
/** Obtains a [[ShuffleHandle]] to pass to tasks. */
override def registerShuffle[K, V, C](
shuffleId: Int,
@@ -179,8 +201,12 @@ class ColumnarShuffleManager(conf: SparkConf)
override def unregisterShuffle(shuffleId: Int): Boolean = {
Option(taskIdMapsForShuffle.remove(shuffleId)).foreach {
mapTaskIds =>
- mapTaskIds.iterator.foreach {
- mapId => shuffleBlockResolver.removeDataByMap(shuffleId, mapId)
+ // The block-migration path mutates this set under the same lock;
iterating without it
+ // risks a ConcurrentModificationException during decommissioning.
+ mapTaskIds.synchronized {
+ mapTaskIds.iterator.foreach {
+ mapId => shuffleBlockResolver.removeDataByMap(shuffleId, mapId)
+ }
}
}
true
@@ -188,7 +214,8 @@ class ColumnarShuffleManager(conf: SparkConf)
/** Shut down this ShuffleManager. */
override def stop(): Unit = {
- shuffleBlockResolver.stop()
+ // SortShuffleManager.stop() stops shuffleBlockResolver.
+ super.stop()
}
}
diff --git
a/gluten-substrait/src/test/scala/org/apache/spark/shuffle/sort/ColumnarShuffleManagerSuite.scala
b/gluten-substrait/src/test/scala/org/apache/spark/shuffle/sort/ColumnarShuffleManagerSuite.scala
index 2baaa8ee15..6aaf35fcef 100644
---
a/gluten-substrait/src/test/scala/org/apache/spark/shuffle/sort/ColumnarShuffleManagerSuite.scala
+++
b/gluten-substrait/src/test/scala/org/apache/spark/shuffle/sort/ColumnarShuffleManagerSuite.scala
@@ -17,6 +17,7 @@
package org.apache.spark.shuffle.sort
import org.apache.spark.SparkConf
+import org.apache.spark.util.collection.OpenHashSet
import org.scalatest.funsuite.AnyFunSuiteLike
@@ -36,4 +37,23 @@ class ColumnarShuffleManagerSuite extends AnyFunSuiteLike {
shuffleManager.stop()
}
}
+
+ // IndexShuffleBlockResolver records blocks migrated in during executor
decommissioning into its
+ // own taskIdMapsForShuffle, and unregisterShuffle deletes map output by
reading the manager's
+ // map. If the two are separate instances, migrated blocks are recorded
where nothing reads them
+ // and their files are never deleted, leaking disk on decommissioned
executors.
+ test("shares taskIdMapsForShuffle with its block resolver") {
+ val conf = new
SparkConf().setMaster("local[2]").setAppName("ColumnarShuffleManagerSuite")
+ val shuffleManager = new ColumnarShuffleManager(conf)
+ try {
+ val shuffleId = 0
+ shuffleManager.shuffleBlockResolver.taskIdMapsForShuffle
+ .put(shuffleId, new OpenHashSet[Long](16))
+
+ assert(shuffleManager.unregisterShuffle(shuffleId))
+ assert(shuffleManager.shuffleBlockResolver.taskIdMapsForShuffle.isEmpty)
+ } finally {
+ shuffleManager.stop()
+ }
+ }
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]