【发布时间】:2018-09-01 13:26:33
【问题描述】:
我有一个 parquet 数据框,结构如下:
- ID 字符串
- 日期日期
- 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