【发布时间】:2020-10-08 14:29:14
【问题描述】:
给定一个数据框,我试图计算过去 30 天内我看到 emailId 的次数。我的函数中的主要逻辑如下:
val new_df = df
.withColumn("transaction_timestamp", unix_timestamp($"timestamp").cast(LongType))
val winSpec = Window
.partitionBy("email")
.orderBy(col("transaction_timestamp"))
.rangeBetween(-NumberOfSecondsIn30Days, Window.currentRow)
val resultDF = new_df
.filter(col("condition"))
.withColumn("count", count(col("email")).over(winSpec))
配置:
spark.executor.cores=5
所以,我可以看到 5 个具有窗口功能的阶段,其中一些阶段完成得非常快(在几秒钟内),还有 2 个甚至在 3 小时内都没有完成,最后被卡住了少量任务(进展非常缓慢):
这对我来说是一个数据倾斜的问题,如果我从数据集中删除所有包含 5 个最高频率 email ids 的行,工作很快就会完成(不到 5 分钟)。
如果我尝试在 window partitionBy 中使用其他键,作业将在几分钟内完成:
Window.partitionBy("email", "date")
但显然,如果我这样做,它会执行错误的计数计算,这不是一个可接受的解决方案。
我尝试了各种其他的 spark 设置,以增加内存、内核、并行度等,但似乎都没有帮助。
Spark 版本:2.2
当前 Spark 配置:
-执行器-内存:100G
-executor-cores: 5
-驱动内存:80G
-spark.executor.memory=100g
使用每台具有 16 核、128 GB 内存的机器。最大节点数最多 500 个。
解决这个问题的正确方法是什么?
更新:只是为了提供更多上下文,这里是原始数据帧和相应的计算数据帧:
val df = Seq(
("a@gmail.com", "2019-10-01 00:04:00"),
("a@gmail.com", "2019-11-02 01:04:00"),
("a@gmail.com", "2019-11-22 02:04:00"),
("a@gmail.com", "2019-11-22 05:04:00"),
("a@gmail.com", "2019-12-02 03:04:00"),
("a@gmail.com", "2020-01-01 04:04:00"),
("a@gmail.com", "2020-03-11 05:04:00"),
("a@gmail.com", "2020-04-05 12:04:00"),
("b@gmail.com", "2020-05-03 03:04:00")
).toDF("email", "transaction_timestamp")
val expectedDF = Seq(
("a@gmail.com", "2019-10-01 00:04:00", 1),
("a@gmail.com", "2019-11-02 01:04:00", 1), // prev one falls outside of last 30 days win
("a@gmail.com", "2019-11-22 02:04:00", 2),
("a@gmail.com", "2019-11-22 05:04:00", 3),
("a@gmail.com", "2019-12-02 03:04:00", 3),
("a@gmail.com", "2020-01-01 04:04:00", 1),
("a@gmail.com", "2020-03-11 05:04:00", 1),
("a@gmail.com", "2020-04-05 12:04:00", 2),
("b@gmail.com", "2020-05-03 03:04:00", 1) // new email
).toDF("email", "transaction_timestamp", count")
【问题讨论】:
-
您能否提及您的机器规格。 ?总内核数和内存?
标签: scala performance dataframe apache-spark apache-spark-sql