【问题标题】:Fast Spark alternative to WHERE column IN other_column快速 Spark 替代 WHERE 列 IN other_column
【发布时间】:2020-09-03 21:17:20
【问题描述】:

我正在寻找一种快速的 PySpark 替代品

SELECT foo FROM bar
WHERE foo IN (SELECT baz FROM bar)

预先收集到 Python 列表中绝对不是一个选项,因为处理的数据帧非常大,并且相对于我提出的其他选项,收集占用了大量时间。所以我想不出一种使用原生 PySparkian where(col(bar).isin(baz)) 的方法,因为在这种情况下 baz 必须是一个列表。

我想出的一个选项是right JOIN 作为 IN 的替代品,left_semi JOIN 作为 NOT IN 的替代品,请考虑以下示例:

bar_where_foo_is_in_baz = bar.join(bar.select('baz').alias('baz_'), col('foo') == col('baz_'), 'right').drop('baz_')

然而,这非常冗长,一段时间后阅读时很难解释,并且当在WHERE 中处理大量条件时会导致相当多的头疼,所以我想避免这种情况。

还有其他选择吗?

编辑(请阅读):

由于我似乎误导了很多答案,我的具体要求是将“WHERE - IN”子句翻译成没有 .collect() 的 PySpark,或者一般来说,映射到 Pythonic 列表(作为内部函数 @987654332 @ 需要我)。

【问题讨论】:

  • Pyspark 有不错的 SQL 支持,所以相信你可以使用 pyspark SQL 库来做到这一点。 spark.apache.org/docs/2.2.0/… 应该让你继续前进。 databricks-prod-cloudfront.cloud.databricks.com/public/… 也用于参考。
  • 听起来不错,我去看看,谢谢。
  • 你知道“broadcast JOIN”的工作原理就是这样,即在所有执行程序中创建一个值的哈希图吗?所以这是一个简单的 SQL 查询,如果你想确保 Spark 优化器不会出错,还有一个额外的提示。不要重新发明轮子。自 25 年前第一个 MPP 数据库以来,此类问题就一直存在。
  • 在 spark 我认为你最好的选择是使用 left_semi 或 left_anti
  • stackoverflow.com/questions/42545788/… 似乎有你需要的答案。

标签: sql pyspark where-in


【解决方案1】:

来自 Jacek Laskowski 的 GitBook

Spark SQL 使用 broadcast join(又名广播哈希连接)而不是 当一侧数据的大小为时,哈希连接优化连接查询 下面spark.sql.autoBroadcastJoinThreshold

广播连接对于大表之间的连接非常有效 (事实)相对较小的表格(尺寸) 用于执行星型模式连接。它可以避免发送所有数据 网络上的大表。

还有also

Spark SQL 2.2 支持使用广播标准的 BROADCAST 提示 函数或 SQL cmets:

  • SELECT /*+ MAPJOIN(b) */ …​
  • SELECT /*+ BROADCASTJOIN(b) */ …​
  • SELECT /*+ BROADCAST(b) */ …​

其实Spark documentation的提示更全面:

BROADCAST 提示引导 Spark 在何时广播每个指定的表 将它们与另一个表或视图连接起来。当 Spark 决定加入时 方法,广播散列连接(即 BHJ)是首选,即使 统计信息高于配置 spark.sql.autoBroadcastJoinThreshold.

当连接的双方都是 指定时,Spark 会广播具有较低统计信息的那个。笔记 Spark 不保证总是选择 BHJ,因为并非所有情况 (例如完全外连接)支持 BHJ。


最后,如果您认真对待开发高效的 Spark 作业,您应该花一些时间来了解这头野兽是如何工作的。

例如presentation about joins from Databricks 应该会有所帮助

【讨论】:

  • 我将在这里回复评论以及答案。问题是,我完全了解广播连接及其用途,但是,我不太能够辨别您如何建议将其应用于手头的场景。我对优化我用作示例的 JOIN 并不特别感兴趣,而是我想以某种方式将“WHERE - IN”子句转换为 Spark,而无需收集,或者更一般地说,映射到列表(因为 .isin() 函数需要)。也许我应该更清楚我的要求,我会编辑这个问题。
  • 那么,您能向我解释一下您的SELECT a.* FROM a WHERE a.x IN (SELECT xx FROM b WHERE wtf)SELECT a.* FROM a JOIN (SELECT xx FROM b WHERE wtf) bb ON a.x =bb.xx 以及SELECT a.* FROM a WHERE EXISTS (SELECT 1 FROM b WHERE wtf AND b.xx =a.x) 之间的区别吗? AFAIK 一个体面的查询编译器会将所有三个转换为完全相同的执行计划。
  • 我认为没有。但是我已经在问题中声明了自己,JOIN,特别是 INNER ,可以用作“WHERE - IN”,这是我的临时解决方案,但它有点笨拙,我希望避免它。是否证明我无法避免它,非常欢迎对连接进行优化。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-10-04
  • 2016-06-16
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多