【问题标题】:Join computed twice on pyspark, maybe I dont understand lazyness?在 pyspark 上加入两次计算,也许我不理解懒惰?
【发布时间】: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


【解决方案1】:

更新答案(2021 年 5 月 9 日):

我认为您可以尝试使用withColumn 方法在连接后指定when(.. ,.. ).otherwise(..) 中的值来指定数据中的分区列(您可以为4 个不同的值嵌套多个when/otherwise 块)。不仅仅是使用partitionBy 写入数据。在这种情况下,您不需要重新计算 4 次。一次计算就足够了。

旧答案:

我认为您可能想使用df.cache() 函数来防止相同的计算。

join_df = (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').cache()

Spark 将计算所有结果并将其保存到内存和磁盘中。它将为新过滤器重用预先计算的join_df

【讨论】:

  • 我不认为我有 100+TB 用于缓存,也许我可以添加一些磁盘(比 AWS 上的 ram 便宜)。缓存不是很多吗?谢谢!
  • 是的,磁盘会便宜得多,但内存要快得多。不幸的是,在 spark 中,为了防止计算,我们需要缓存它,否则它将尝试重新计算每个动作。如果您的一张桌子很小,您也可以尝试像df_a.join(broadcast(df_b), ... 一样广播它,这将防止随机播放。
  • Gb 数据两边,广播减慢了工作,必须看看保存在磁盘上的时间是否少于进行 4 次连接
  • 使用内存缓存,它比执行 4 个连接要慢...不知道为什么,主连接比 for 循环中的 4 个单独连接花费更多时间
  • 嗯,在这种情况下,您的基本读取数据总量远低于连接值,因为在您的情况下缓存速度较慢。我认为您可以尝试通过在连接后在when(.. ,.. ).otherwise) 中指定值来使用withColumn mathob 在数据中指定一个分区列。不仅仅是用partitionBy 写你的数据。在这种情况下,您不需要重新计算 4 次。一次计算就足够了。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-06-15
  • 1970-01-01
  • 2015-09-06
  • 1970-01-01
相关资源
最近更新 更多