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]

Reply via email to