【发布时间】:2020-10-21 07:27:10
【问题描述】:
我有一个 Spark UDF 来计算列的滚动计数,精确到时间。如果我需要计算 24 小时的滚动计数,例如对于时间为 2020-10-02 09:04:00 的条目,我需要回顾到 2020-10-01 09:04:00(非常精确)。
如果我在本地运行,滚动计数 UDF 可以正常工作并给出正确的计数,但是当我在集群上运行时,它给出的结果不正确。这是示例输入和输出
输入
+---------+-----------------------+
|OrderName|Time |
+---------+-----------------------+
|a |2020-07-11 23:58:45.538|
|a |2020-07-12 00:00:07.307|
|a |2020-07-12 00:01:08.817|
|a |2020-07-12 00:02:15.675|
|a |2020-07-12 00:05:48.277|
+---------+-----------------------+
预期输出
+---------+-----------------------+-----+
|OrderName|Time |Count|
+---------+-----------------------+-----+
|a |2020-07-11 23:58:45.538|1 |
|a |2020-07-12 00:00:07.307|2 |
|a |2020-07-12 00:01:08.817|3 |
|a |2020-07-12 00:02:15.675|1 |
|a |2020-07-12 00:05:48.277|1 |
+---------+-----------------------+-----+
最后两个条目值在本地是 4 和 5,但在集群上它们是不正确的。我最好的猜测是数据正在跨执行器分布,并且 udf 也在每个执行器上并行调用。由于 UDF 的参数之一是列(本示例中的分区键 - OrderName),如果是这种情况,我如何控制/纠正集群的行为。以便它以正确的方式计算每个分区的正确计数。有什么好的建议
【问题讨论】:
-
你能显示你的UDF代码吗?
-
我不能完全共享它,它类似于 udf {(ordername: Partition, time: Range, Long) process:{}},UDF 的初始要求,同一分区内的所有记录都已排序按日期。它的作用是针对每个分区(此处为订单名称),如果新记录用于现有分区,则将记录添加到队列中,增加计数,然后检查当前时间是否为当前时间,队列中到目前为止的所有记录是否在 24 小时内, 如果不是,则从开头删除记录(因为它是队列)
-
如果你能显示这将非常有用:
input data/dataframe和expected output/expected dataframe。 -
我用输入和预期的输出数据框更新了问题
标签: scala apache-spark user-defined-functions