【发布时间】:2016-05-31 15:22:06
【问题描述】:
如果我有一个每分钟有卷的 RDD,例如
(("12:00" -> 124), ("12:01" -> 543), ("12:02" -> 102), ... )
我想将其映射到一个数据集,其中包含本分钟的音量、前一分钟的音量、前 5 分钟的平均音量。例如
(("12:00" -> (124, 300, 245.3)),
("12:01" -> (543, 124, 230.2)),
("12:02" -> (102, 543, 287.1)))
输入 RDD 可以是 RDD[(DateTime, Int)],输出可以是 RDD[(DateTime, (Int, Int, Float))]。
有什么好的方法可以做到这一点?
【问题讨论】:
-
您的数据是完整的还是可能有一些缺失的记录?
-
可能会有差距,我会默认为零。我不介意解决方案是否能解决这个问题。
-
在纯 scala 中,我会转换为 DateTime 并使用 SortedMap。你的数据集有多大?
-
我想这种方法可能会有所帮助,我会在时间序列上创建窗口:stackoverflow.com/questions/23402303/…。我认为创建窗口与您尝试做的很接近,也就是说,将一个值与其系列中的上一个/下一个值组合在一起
-
评论赞赏,有用的东西。
标签: scala apache-spark