【发布时间】: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