【问题标题】:Forcing pyspark join to occur sooner强制 pyspark join 更快发生
【发布时间】:2018-05-05 12:15:49
【问题描述】:

问题:我有两张大小相差很大的桌子。我想通过左外连接来加入一些 id。不幸的是,由于某种原因,即使在对所有记录执行连接后缓存我的操作之后,即使我只想要左表中的那些。见下文:

我的问题: 1. 我该如何设置,以便只有与左表匹配的记录通过代价高昂的争吵步骤得到处理?

LARGE_TABLE => ~900M 记录

SMALL_TABLE => 50 万条记录

代码:

combined = SMALL_TABLE.join(LARGE_TABLE SMALL_TABLE.id==LARGE_TABLE.id, 'left-outer')
print(combined.count())
...
...
# EXPENSIVE STUFF!
w = Window().partitionBy("id").orderBy(col("date_time"))
data = data.withColumn('diff_id_flag', when(lag('id').over(w) != col('id'), lit(1)).otherwise(lit(0)))

不幸的是,我的执行计划显示上述昂贵的转换操作是在大约 9 亿条记录上完成的。我觉得这很奇怪,因为我运行 df.count() 来强制连接急切地执行而不是懒惰地执行。

有什么想法吗?

附加信息: - 请注意,我的代码流中的昂贵转换发生在联接之后(至少我是这样解释的),但我的 DAG 显示作为联接的一部分发生的昂贵转换。这正是我想要避免的,因为转换成本很高。我希望连接执行,然后连接的结果通过昂贵的转换运行。 - 假设较小的表无法放入内存。

【问题讨论】:

  • 对不起。只是在里面放了一些东西来表明我正在做一些昂贵的改造。添加的语句更具说明性。
  • 您可能想尝试在加入后保留数据帧。然后运行df.count(),然后执行昂贵的争吵操作。

标签: apache-spark pyspark lazy-evaluation


【解决方案1】:

最好的方法是broadcast 微型数据框。缓存适用于多个操作,这似乎不适用于您的特定用例。

【讨论】:

  • 假设我的小数据框现在有 50 万条记录?显然仍然不想对所有 900M 大表记录进行复杂的争论。
  • @DonaldVetal 它实际上取决于较小数据框的实际大小。 (这是否适合内存)我不太确定你所说的争吵是什么意思,但在连接中你必须遍历每一行。除非你在加入之前的转换成本太高,在这种情况下你应该先加入然后再转换。
  • 在我的帖子底部添加了其他信息。
【解决方案2】:
  • df.count 对执行计划完全没有影响。这只是在没有任何充分理由的情况下执行的昂贵操作。
  • 其中的窗口函数应用程序需要与join 相同的逻辑。因为您通过idpartitionBy id 加入,所以两个阶段都需要对双方进行相同的哈希分区和完整数据扫描。没有可接受的理由将这两者分开。

    实际上join逻辑应该在窗口之前应用,作为同一阶段下游转换的过滤器。

【讨论】:

    猜你喜欢
    • 2011-07-19
    • 1970-01-01
    • 2010-12-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-10-07
    • 2010-12-19
    • 2014-11-16
    相关资源
    最近更新 更多