【发布时间】: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