Skip to content

Commit 32abd31

Browse files
committed
[AURON #2351] Reduce native broadcast task serialization size
1 parent ffc6a27 commit 32abd31

2 files changed

Lines changed: 31 additions & 13 deletions

File tree

spark-extension-shims-spark/src/test/scala/org/apache/auron/AuronCheckConvertBroadcastExchangeSuite.scala

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,13 +16,43 @@
1616
*/
1717
package org.apache.auron
1818

19+
import org.apache.spark.SparkEnv
1920
import org.apache.spark.sql.{AuronQueryTest, Row}
2021
import org.apache.spark.sql.execution.auron.plan.NativeBroadcastExchangeExec
2122
import org.apache.spark.sql.execution.exchange.BroadcastExchangeExec
2223

2324
class AuronCheckConvertBroadcastExchangeSuite extends AuronQueryTest with BaseAuronSQLSuite {
2425
import testImplicits._
2526

27+
test("do not serialize the broadcast relation with Spark tasks") {
28+
withSQLConf(
29+
"spark.auron.enable.broadcastExchange" -> "true",
30+
"spark.auron.enable.bhj" -> "false") {
31+
val payload = "x" * 4096
32+
(0 until 256)
33+
.map(i => (i, s"$i$payload"))
34+
.toDF("key", "payload")
35+
.createOrReplaceTempView("broad_cast_table1")
36+
Seq(0, 255).toDF("key").createOrReplaceTempView("broad_cast_table2")
37+
38+
val df = spark.sql(
39+
"select /*+ broadcast(a)*/ b.key from broad_cast_table1 a " +
40+
"inner join broad_cast_table2 b on a.key = b.key")
41+
42+
checkAnswer(df, Seq(Row(0), Row(255)))
43+
val exchange = collectFirst(df.queryExecution.executedPlan) {
44+
case broadcastExchangeExec: NativeBroadcastExchangeExec => broadcastExchangeExec
45+
}.get
46+
val broadcast = exchange.executeBroadcast[Any]()
47+
try {
48+
val serialized = SparkEnv.get.closureSerializer.newInstance().serialize(broadcast)
49+
assert(serialized.remaining() < 64 * 1024)
50+
} finally {
51+
broadcast.destroy()
52+
}
53+
}
54+
}
55+
2656
test(
2757
"test bhj broadcastExchange to native where spark.auron.enable.broadcastExchange is true") {
2858
withSQLConf("spark.auron.enable.broadcastExchange" -> "true") {

spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeBroadcastExchangeBase.scala

Lines changed: 1 addition & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,6 @@ import scala.collection.immutable.SortedMap
2727
import scala.concurrent.Promise
2828
import scala.jdk.CollectionConverters._
2929

30-
import org.apache.commons.lang3.reflect.MethodUtils
3130
import org.apache.spark.OneToOneDependency
3231
import org.apache.spark.Partition
3332
import org.apache.spark.SparkException
@@ -145,18 +144,7 @@ abstract class NativeBroadcastExchangeBase(mode: BroadcastMode, override val chi
145144
.map(_.copy())
146145
.toArray
147146

148-
val broadcast = relationFuture.get // broadcast must be resolved
149-
val v = mode.transform(dataRows)
150-
val dummyBroadcasted = new Broadcast[Any](-1) {
151-
override protected def getValue(): Any = v
152-
override protected def doUnpersist(blocking: Boolean): Unit = {
153-
MethodUtils.invokeMethod(broadcast, true, "doUnpersist", Array(blocking))
154-
}
155-
override protected def doDestroy(blocking: Boolean): Unit = {
156-
MethodUtils.invokeMethod(broadcast, true, "doDestroy", Array(blocking))
157-
}
158-
}
159-
dummyBroadcasted.asInstanceOf[Broadcast[T]]
147+
sparkContext.broadcast(mode.transform(dataRows)).asInstanceOf[Broadcast[T]]
160148
}
161149

162150
def doExecuteBroadcastNative[T](): broadcast.Broadcast[T] = {

0 commit comments

Comments
 (0)