【发布时间】:2017-05-07 12:19:57
【问题描述】:
我有两个 DataFrame A 和 B:
-
A的列(id, info1, info2)大约有 2 亿行 -
B只有列id有 100 万行
id 列在两个 DataFrame 中都是唯一的。
我想要一个新的 DataFrame,它过滤 A 以仅包含来自 B 的值。
如果 B 非常小,我知道我会这样做
A.filter($("id") isin B("id"))
但B 仍然很大,所以不是所有的都可以作为广播变量。
我知道我可以使用
A.join(B, Seq("id"))
但这不会利用独特性,恐怕会导致不必要的洗牌。
完成该任务的最佳方法是什么?
【问题讨论】:
-
是什么让你觉得“我怕会造成不必要的洗牌”?
-
我相信 Spark 不会将所有的小数据帧存储到所有节点,导致它在加入时会随机播放。此外,如果 spark 知道唯一性,如果找到一个值,它可能会停止发送值。如果我错了,请纠正我。
-
听起来是对的,但考虑到所有 Spark 优化都是相当新的/年轻的并且不一定经过实战考验,请根据具体情况猜测。
标签: apache-spark apache-spark-sql