AbsaOSS / AbsaOSS/enceladus

Optimise or allow for broadcast

未关闭
#70 0 条评论 0 个 reaction 已指派 1 人 已被 @yruslan 认领 在 GitHub 查看
Conformance feature priority: medium under discussion
主要语言
Scala
星标
33
派生
16
PR 合并指标
30 天内没有已合并 PR

描述

"Automatic broadcast on join tables <10MB happens in spark 2.x + already, this covers a lot of the tables already"
References:
- https://stackoverflow.com/questions/43984068/does-spark-sql-autobroadcastjointhreshold-work-for-joins-using-datasets-join-op
- https://jaceklaskowski.gitbooks.io/mastering-spark-sql/spark-sql-joins-broadcast.html

In short:

> "Spark SQL uses broadcast join instead of hash join to optimize join queries when the size of one side data is below spark.sql.autoBroadcastJoinThreshold." "Spark will use autoBroadcastJoinThreshold and automatically broadcast data"

> "Broadcasting large objects is unlikely provide any performance boost, and in practice will often degrade performance and result in stability issue. Remember that broadcasted object has to be first fetch to driver, then send to each worker, and finally loaded into memory."

贡献指南

打开贡献指南

评估

这个 Issue 还没有评估数据。

把新 issue 发到你的邮箱

精选适合新手参与的 GitHub issue 摘要。