【问题标题】:Can I set different autoBroadcastJoinThreshold value in sparkConf for different sql?我可以在 sparkConf 中为不同的 sql 设置不同的 autoBroadcastJoinThreshold 值吗?
【发布时间】:2018-03-21 08:05:50
【问题描述】:

我有很大的DataFrames:A(200g), B(20m), C(15m), D(10m), E(12m),我想把它们连接在一起: A join B, C join D and E 在同一个 SparkSession** 中使用 spark sql。就像:

absql:sql("select * from A a inner join B b on a.id=b.id").write.csv("/path/for/ab")
cdesql:sql("select * from C c inner join D d on c.id=d.id inner join E e on c.id=e.id").write.csv("/path/for/cde")

问题:

当我使用默认spark.sql.autoBroadcastJoinThreshold=10m

  • absql 需要很长时间,原因是absql skew。
  • cdesql正常

当我设置spark.sql.autoBroadcastJoinThreshold=20m

  • C、D、E会被广播,所有的任务都会在同一个executor中执行,还是需要很长的时间。
  • 如果设置num-executors=200,广播时间长
  • absql正常

【问题讨论】:

  • @Shaido 的回答可以解决我的问题。这个question让我知道函数broadcast不受参数autoBroadcastJoinThreshold的影响

标签: apache-spark broadcast skew


【解决方案1】:

您可以标记要广播的数据帧,而不是更改autoBroadcastJoinThreshold。通过这种方式,很容易决定应该广播哪些数据帧。

在 Scala 中它可能如下所示:

import org.apache.spark.sql.functions.broadcast
val B2 = broadcast(B)
B2.createOrReplaceTempView("B")

这里的数据帧 B 已被标记为广播,然后被注册为要与 Spark SQL 一起使用的表。


或者,这可以直接使用dataframe API来完成,第一个join可以写成:

A.join(broadcast(B), Seq("id"), "inner")

【讨论】:

  • 谢谢@Shaido!这对我很有帮助。在这个问题link
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2015-11-22
  • 1970-01-01
  • 1970-01-01
  • 2014-12-05
  • 1970-01-01
  • 2018-05-31
  • 1970-01-01
相关资源
最近更新 更多