This is an automated email from the ASF dual-hosted git repository.
marin-ma 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 f47f027220 [CORE] Avoid a per-row copy on row-based shuffle writes
(#13039)
f47f027220 is described below
commit f47f027220e9f73c6f279f442c793c1b93e14425
Author: Ankita Victor <[email protected]>
AuthorDate: Mon Sep 21 13:59:16 2026 +0530
[CORE] Avoid a per-row copy on row-based shuffle writes (#13039)
---
.../shuffle/sort/ColumnarShuffleManager.scala | 13 +++++++-
.../shuffle/sort/ColumnarShuffleManagerSuite.scala | 39 ++++++++++++++++++++++
2 files changed, 51 insertions(+), 1 deletion(-)
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 3b4126dc7c..85c860bc0a 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
@@ -32,8 +32,19 @@ import java.util.concurrent.ConcurrentHashMap
import scala.collection.JavaConverters._
+/**
+ * Extends [[SortShuffleManager]] rather than [[ShuffleManager]] on purpose.
+ *
+ * `spark.shuffle.manager` is process-wide, so this manager also serves the
row-based exchanges that
+ * Gluten does not offload.
`ShuffleExchangeExec.needToCopyObjectsBeforeShuffle` gates the per-row
+ * defensive `UnsafeRow.copy()` on `isInstanceOf[SortShuffleManager]`, falling
through to a
+ * catch-all `true` for any other implementation. The copy only protects
`ExternalSorter`, which
+ * 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.
+ */
class ColumnarShuffleManager(conf: SparkConf)
- extends ShuffleManager
+ extends SortShuffleManager(conf)
with SupportsColumnarShuffle
with Logging {
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
new file mode 100644
index 0000000000..2baaa8ee15
--- /dev/null
+++
b/gluten-substrait/src/test/scala/org/apache/spark/shuffle/sort/ColumnarShuffleManagerSuite.scala
@@ -0,0 +1,39 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.spark.shuffle.sort
+
+import org.apache.spark.SparkConf
+
+import org.scalatest.funsuite.AnyFunSuiteLike
+
+class ColumnarShuffleManagerSuite extends AnyFunSuiteLike {
+
+ // ShuffleExchangeExec.needToCopyObjectsBeforeShuffle checks
+ // `SparkEnv.get.shuffleManager.isInstanceOf[SortShuffleManager]` and, when
that is false, falls
+ // through to a catch-all `true` that makes every map task run
UnsafeRow.copy() per row. Since
+ // spark.shuffle.manager is process-wide, breaking this subtype relationship
silently regresses
+ // every row-based exchange: no exception, no test failure, only lost
throughput.
+ test("is a SortShuffleManager so row-based exchanges keep Spark's zero-copy
write path") {
+ val conf = new
SparkConf().setMaster("local[2]").setAppName("ColumnarShuffleManagerSuite")
+ val shuffleManager = new ColumnarShuffleManager(conf)
+ try {
+ assert(shuffleManager.isInstanceOf[SortShuffleManager])
+ } finally {
+ shuffleManager.stop()
+ }
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]