【问题标题】:SPARK - Joining two data streams - maintenance of cacheSPARK - 加入两个数据流 - 缓存的维护
【发布时间】:2016-04-19 07:57:36
【问题描述】:

很明显,火花流中的开箱即用连接功能并不能保证很多现实生活中的用例。原因是它只加入了微批处理 RDD 中包含的数据。

用例是将来自两个 kafka 流的数据连接起来,并将 stream1 中的每个对象用它在 spark 中的 stream2 中的相应对象丰富,并将其保存到 HBase。

实施将

  • 在内存中维护来自 stream2 对象的数据集,在收到对象时添加或替换对象

  • 对于stream1中的每个元素,访问缓存以从stream2中找到匹配的对象,如果找到匹配则保存到HBase,否则将其放回kafka流中。

这个问题是关于 Spark 流的探索,它的 API 是为了找到实现上述方法的方法。

【问题讨论】:

  • 问题是……?
  • 把查询放在最后一行,看看现在对你有没有意义。

标签: join apache-spark apache-kafka spark-streaming


【解决方案1】:

您可以将传入的RDDs 加入到其他RDDs - 而不仅仅是那个微批次中的那些。基本上,您会保留一个“运行总数”RDD,您可以填写如下内容:

var globalRDD1: RDD[...] = sc.emptyRDD
var globalRDD2: RDD[...] = sc.emptyRDD

dstream1.foreachRDD(rdd => if (!rdd.isEmpty) globalRDD1 = globalRDD1.union(rdd))
dstream2.foreachRDD(rdd => if (!rdd.isEmpty) {
  globalRDD2 = globalRDD2.union(rdd))
  globalRDD1.join(globalRDD2).foreach(...) // etc, etc
}

【讨论】:

  • 感谢您的回复。您提到的内容将使“缓存”RDD 保持最新。但是,每次都继续广播 RDD 会很昂贵。对此有什么想法吗?可以有不同的方法吗?
  • 广播是什么意思?我的代码中没有broadcast
  • 应该使用不同的词来表示它,并不意味着像“共享变量”中那样广播。当您每次都执行 globalRDD1.join(globalRDD2) 时,您会将所有数据从驱动程序发送到执行程序 - 甚至是在之前的微批处理中更早发送的数据,并且可能已经处理过。再次感谢您的回复。
  • 不,我不相信它是这样工作的。否则,跨越两个大的 RDDsunion 运算符将非常昂贵,而且它不是(不像将所有数据来回传送给驱动程序那样昂贵)。
  • 在这里,看看这个答案:“联合是一种非常有效的操作,因为它不会移动任何数据”stackoverflow.com/questions/29977526/…
【解决方案2】:

一个好的开始是查看mapWithState。这是对updateStateByKey 的更有效替代。这些是在PairDStreamFunction 上定义的,因此假设stream2 中V 类型的对象由K 类型的某个键标识,您的第一点将如下所示:

def stream2: DStream[(K, V)] = ???

def maintainStream2Objects(key: K, value: Option[V], state: State[V]): (K, V) = {
  value.foreach(state.update(_))
  (key, state.get())
}

val spec = StateSpec.function(maintainStream2Objects)

val stream2State = stream2.mapWithState(spec)

stream2State 现在是一个流,其中每个批次都包含 (K, V) 对以及每个键的最新值。您可以在此流和 stream1 上执行连接,以执行您的第二点的进一步逻辑。

【讨论】:

  • 那么,您将如何在 24 小时后“过期”过时数据?
猜你喜欢
  • 2017-07-16
  • 1970-01-01
  • 2018-12-19
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-02-21
  • 2018-02-09
  • 2013-03-29
相关资源
最近更新 更多