【问题标题】:Spark Window performance issuesSpark Window 性能问题
【发布时间】:2018-09-01 13:26:33
【问题描述】:

我有一个 parquet 数据框,结构如下:

  1. ID 字符串
  2. 日期日期
  3. 480 个 Double 类型的其他特征列

我必须用相应的加权移动平均值替换 480 个特征列中的每一个,窗口为 250。 最初,我尝试使用以下简单代码对单个列执行此操作:

var data = sparkSession.read.parquet("s3://data-location")
var window = Window.rowsBetween(-250, Window.currentRow - 1).partitionBy("ID").orderBy("DATE")
data.withColumn("Feature_1", col("Feature_1").divide(avg("Feature_1").over(window))).write.parquet("s3://data-out")

输入数据包含 2000 万行,每个 ID 关联大约 4-5000 个日期。 我已在 AWS EMR 集群(m4.xlarge 实例)上运行此程序,其中一列的结果如下:

  • 4 个执行器 X 4 个内核 X 10 GB + 1 GB 用于纱线开销(因此每个任务 2.5GB,16 个并发运行的任务),耗时 14 分钟
  • 8 个执行器 X 4 个内核 X 10GB + 1 GB 用于纱线开销(因此每个任务 2.5GB,32 个并发运行的任务),耗时 8 分钟

我已经调整了以下设置,希望能降低总时间:

  • spark.memory.storageFraction 0.02
  • spark.sql.windowExec.buffer.in.memory.threshold 100000
  • spark.sql.constraintPropagation.enabled false

第二个有助于防止日志中出现一些溢出,但对实际性能没有任何帮助。

我不明白为什么只有 2000 万条记录需要这么长时间。我知道,对于计算加权移动平均值,它需要进行 20 M X 250(窗口大小)的平均值和除法,但是对于 16 个内核(第一次运行),我不明白为什么需要这么长时间。我无法想象剩下的 479 个特征列需要多长时间!

我还尝试通过设置来增加默认的随机分区:

  • spark.sql.shuffle.partitions 1000

但即使有 1000 个分区,它也没有减少时间。 还尝试在调用窗口聚合之前按 ID 和 DATE 对数据进行排序,但没有任何好处。

有什么方法可以改善这一点,或者窗口函数在我的用例中通常运行缓慢?这只是 2000 万行,远不及 spark 可以处理其他类型的工作负载..

【问题讨论】:

    标签: apache-spark apache-spark-sql spark-dataframe apache-spark-mllib


    【解决方案1】:

    您的数据集大小约为 70 GB。 如果我对每个 id 的理解正确,它会按日期对所有记录进行排序,然后取前面的 250 条记录进行平均。由于您需要在 400 多列上应用它,我建议在创建镶木地板时尝试分桶以避免洗牌。写入分桶拼花文件需要相当长的时间,但所有 480 列的推导可能不需要 8 分钟 *480 执行时间。

    请在创建 parquet 文件时尝试分桶或重新分区和 sortwithin 并让我知道它是否有效。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-02-10
      • 1970-01-01
      • 2018-08-16
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-02-01
      相关资源
      最近更新 更多