【问题标题】:filter to join pyspark过滤加入 pyspark
【发布时间】:2018-04-07 20:26:48
【问题描述】:

现在我有以下代码:

df1 = df.filter((df.col1.isin(List1)) | (df.col2.isin(List2)))

我从 dataframe 中做collect() 得到 List1 所以我喜欢使用 join 我尝试了以下

df1=df.filter(df.col2.isin(List2))
df2=df.join(df_List1,'col1','leftsemi')
df3=df1.join(df2,'col1' ,'outer')

我有两个问题:

  1. 转换原始语句的正确方法是什么
  2. 在性能方面值得做吗

【问题讨论】:

  • “spark 无法完成工作”是什么意思?帖子中没有作业触发代码
  • 我正在收集以检查 df3 的结果
  • 您是否检查了 Web UI 以了解发生了什么?这个处理多少数据?初始数据源是什么?您是否有足够的资源以预期的速度运行?您是否使用任何缓存?
  • 第一个命令运行良好,使用外连接是否正确?
  • 使用外连接是可以的并且确实有效。这可能只是数据集大小和可用资源的问题。外部连接可以产生大量数据,集群可能正在为此苦苦挣扎。您可能需要查看 web ui

标签: python apache-spark join pyspark


【解决方案1】:

在性能方面值得做吗

与往常一样,在询问与性能相关的问题时,您应该同时使用真实数据(或真实反映数据真实分布的数据)和使用您可以支配的实际资源进行测试。

话虽这么说,如果List1List2 足够小,df.filter((df.col1.isin(List1)) | (df.col2.isin(List2))) 可以查询成功,则很难期望使用joins 有任何改进,因为基于ORJOINS 无法轻松优化.

可以表示为:

SELECT DISTINCT col1, col2 FROM ( 
   SELECT t.* FROM t JOIN r1 WHERE t.col1 = r1.col1
   UNION
   SELECT t.* FROM t JOIN r2 WHERE t.col2 = r2.col2
)

WITH 
  t1 AS (SELECT t.* FROM t JOIN r1 WHERE t.col1 = r1.col1),
  t2 AS (SELECT t.* FROM t JOIN r2 WHERE t.col2 = r2.col2) 
SELECT * FROM t1 FULL OUTER JOIN t2 ON t1.col1 = t2.col1  -- if col1 is an unique identifier

一般情况下(我们可以放心地忽略t 足够小以存储在单台机器的内存中的情况)两者都需要t 的完全洗牌,使t 的大小成为限制因素。

此外,当第一个组件为真时,与本地对象的逻辑析取可以短路评估。在一般的连接情况下不可能这样做。

未来 Spark 应该支持 isin (SPARK-23945) 中的单列 DataFrame,并且优化器应该能够为您在广播和哈希连接之间做出决定。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2011-06-23
    • 1970-01-01
    • 2017-03-16
    • 2021-04-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多