【发布时间】:2016-08-22 07:43:51
【问题描述】:
我需要对具有以下格式的 rdd 进行一些处理:
RDD[(String,List[(String,String,String,String)])]
以下是来自 RDD 的示例条目:
(600,List((5,111,1,1), (15,111,1,5), (38,111,2,null))
(700,List((5,111,1,1), (35,111,1,5), (39,111,2,null))
我需要根据在列表中每个元组的第一个元素中找到的时间戳值将每个条目拆分为多个条目。每个条目应包含 20 分钟间隔内的时间戳。
例如,第一个条目应该分成2个条目:
List((5,111,1,1), (15,111,1,5))
List((38,111,2,null))
最终结果应该是RDD[(String,List[(String,String,String,String)])]:
(600,List((5,111,1,1), (15,111,1,5)))
(600,List((38,111,2,null)))
(700,List((5,111,1,1))
(700,List((35,111,1,5), (39,111,2,null))
任何提示如何做到这一点以及应用哪些功能?
【问题讨论】:
标签: scala apache-spark rdd