Spark:如何按时间范围加入RDD

M. *_*irn 14 cassandra apache-spark rdd

我有一个微妙的Spark问题,我无法绕过头.

我们有两个RDD(来自Cassandra).RDD1包含Actions和RDD2包含Historic数据.两者都有一个id可以匹配/加入.但问题是这两个表有一个N:N关系.Actions包含多个具有相同ID的行,因此也是如此Historic.以下是两个表的一些示例日期.

Actions 时间实际上是一个时间戳

id  |  time  | valueX
1   |  12:05 | 500
1   |  12:30 | 500
2   |  12:30 | 125
Run Code Online (Sandbox Code Playgroud)

Historic set_at实际上是一个时间戳

id  |  set_at| valueY
1   |  11:00 | 400
1   |  12:15 | 450
2   |  12:20 | 50
2   |  12:25 | 75
Run Code Online (Sandbox Code Playgroud)

我们如何以某种方式加入这两个表,我们得到这样的结果

1   |  100  # 500 - 400 for Actions#1 with time 12:05 because Historic was in that time at 400
1   |  50   # 500 - 450 for Actions#2 with time 12:30 because H. was in that time at 450
2   |  50   # 125 - 75  for Actions#3 with time 12:30 because H. was in that time at 75
Run Code Online (Sandbox Code Playgroud)

我无法想出一个感觉正确的好解决方案,而无需对大型数据集进行大量迭代.我总是要考虑从Historic集合中制作一个范围,然后以某种方式检查是否Actions符合范围,例如(11:00 - 12:15)进行计算.但这对我来说似乎很慢.有没有更有效的方法呢?在我看来,这种问题可能很受欢迎,但我还没有找到任何暗示.你如何在火花中解决这个问题?

我目前的尝试到目前为止(完成代码的一半)

case class Historic(id: String, set_at: Long, valueY: Int)
val historicRDD = sc.cassandraTable[Historic](...)

historicRDD
.map( row => ( row.id, row ) )
.reduceByKey(...) 
// transforming to another case which results in something like this; code not finished yet
// (List((Range(0, 12:25), 400), (Range(12:25, NOW), 450)))

// From here we could join with Actions
// And then some .filter maybe to select the right Lists tuple
Run Code Online (Sandbox Code Playgroud)

maa*_*asg 4

这是一个有趣的问题。我还花了一些时间找出一种方法。这就是我想出的:

给定案例类别Action(id, time, x)和Historic(id, time, y)

  • 将动作与历史结合起来(这可能很重)
  • 过滤与给定操作不相关的所有历史数据
  • key the results by (id,time) - 区分不同时间的相同键
  • 将操作的历史记录减少到最大值,为我们留下给定操作的相关历史记录

在火花中:

val actionById = actions.keyBy(_.id)
val historyById = historic.keyBy(_.id)
val actionByHistory = actionById.join(historyById)
val filteredActionByidTime = actionByHistory.collect{ case (k,(action,historic)) if (action.time>historic.t) => ((action.id, action.time),(action,historic))}
val topHistoricByAction = filteredActionByidTime.reduceByKey{ case ((a1:Action,h1:Historic),(a2:Action, h2:Historic)) =>  (a1, if (h1.t>h2.t) h1 else h2)}

// we are done, let's produce a report now
val report = topHistoricByAction.map{case ((id,time),(action,historic)) => (id,time,action.X -historic.y)}
Run Code Online (Sandbox Code Playgroud)

使用上面提供的数据,报告如下所示:

report.collect
Array[(Int, Long, Int)] = Array((1,43500,100), (1,45000,50), (2,45000,50))
Run Code Online (Sandbox Code Playgroud)

(我将时间转换为秒以获得简单的时间戳)