【发布时间】:2021-08-28 20:46:37
【问题描述】:
自从我上次使用 spark 太久了,我再次使用 Spark 3.1,这是我的问题: 我有 20M 行左加入 400M 行,原代码是:
times= [50000,20000,10000,1000]
for time in times:
join = (df_a.join(df_b,
[
df_a["a"] == df_b["a"],
(unix_timestamp(events["date"]) - unix_timestamp(details["date"])) / 3600
> 5,
(df_a["task"]) = (df_b["task"]-time))
], 'left')
知道每次迭代(时间变量)都包含下一个我在与每个值进行比较之前使 DataFrame 更轻的想法,因此编码如下:
times= [50000,20000,10000,1000]
join = (df_a.join(df_b,
[
df_a["a"] == df_b["a"],
(unix_timestamp(events["date"]) - unix_timestamp(details["date"])) / 3600
> 5,
(df_a["task"]) = (df_b["task"]-50000))
], 'left')
join.checkpoint() # Save current state, and cleaned dataframe
for time in times:
step_join = join_df.where((join_df["task"]) = (join_df["task"]-time)))
# Make calculations and Store result for the iteration...
在查看 Spark 历史服务器上的可视化 SQL 图时,似乎没有使用我在第二次连接上改进的解决方案(?),它使整个左连接在每次迭代时再次出现,没有使用清洁器,更轻的 DataFrame。
我的最终想法是在下一次迭代中使用新的 df,这样每个过滤器都会更轻。我的想法正确吗?我错过了什么吗?
它的样子,这是一个仍在运行的代码,中间的SortMergeJoin是解耦过滤器,第二个“过滤器”只是过滤了一点,但在左右你可以看到它是再次计算 SortMergeJoin 而不是重用先前计算的。
上次必须删除检查点,因为连接上有 55B 行,很难存储数据 (>100TB)
我的集群配置为 30 个实例 64vcore 488GB RAM + 驱动程序
"spark.executor.instances", "249").config("spark.executor.memoryOverhead", "10240").config(
"spark.executor.memory", "87g").config("spark.executor.cores", "12").config("spark.driver.cores", "12").config(
"spark.default.parallelism", "5976").config("spark.sql.adaptive.enabled", "true").config(
"spark.sql.adaptive.skewJoin.enabled", "true").config("spark.sql.shuffle.partitions", "3100").config(
"spark.yarn.driver.memoryOverhead", "10240").config("spark.sql.autoBroadcastJoinThreshold", "2100").config(
"spark.sql.legacy.timeParserPolicy", "LEGACY").getOrCreate()
我在这个网站上使用 excel 计算器来调整除了 spark.sql.shuffle.partitions https://www.c2fo.io/c2fo/spark/aws/emr/2016/07/06/apache-spark-config-cheatsheet/ 现在每个节点使用 10 个执行程序之外的所有内容
尝试在连接上使用 .cache() 它仍然比 4 个并行连接慢,第一个连接更慢。 请注意,.cache() 对子集有好处,但对于 100TB 的连接结果,它会更慢,因为它会缓存到磁盘。 谢谢!
【问题讨论】:
-
很难从这个问题中理解你想要做什么,但我猜想在相同的数据集上循环 4 个连接可能不是正确的解决方案。也许你解释更多你试图解决的问题,它是一个归因模型吗?
-
我将多个组绑定到时间然后进行多个聚合,这些组使用任务“a”和任务“b”中的时间进行过滤,这就是我需要加入它的原因,然后对于每个组,我进行了一些计算,例如 sum/avg/etc 并创建一个数据框,其中包含每个时间窗口的所有聚合
标签: apache-spark pyspark apache-spark-sql