【问题标题】:How to join two Dstream in spark streaming如何在火花流中加入两个 Dstream
【发布时间】:2018-04-05 12:55:30
【问题描述】:

有一个系统同时为两个kafka topic生成数据。

例如:
第 1 步:系统创建一个数据,例如(id=1, main=A, detail=a, ...).

第 2 步:数据将分为两部分,例如(id=1, main=A ...)(id=1, detail=a, ...)

第 3 步:一个发送到 topic1,另一个发送到 topic2

所以我想使用火花流合并两个主题的数据:

data_main = KafkaUtils.createStream(ssc, zkQuorum='', groupId='', topics='topic1')
data_detail = KafkaUtils.createStream(ssc, zkQuorum='', groupId='', topics='topic2')

result = data_main.transformWith(lambda x, y: x.join(y), data_detail)
# outout:
# (id=1, main=A, detail=a, ...)

但是想想这种情况:

(id=1, main=A ...) 可能在data_main 的批次1 中,(id=1, detail=a, ...) 可能在data_detail 的批次2 中。 它们非常接近,但不在同一批次时间。

如何处理这种情况?非常感谢您的建议

【问题讨论】:

    标签: apache-spark spark-streaming


    【解决方案1】:

    你试过开窗吗? 因为窗口化有助于回顾并查看其他批处理间隔中的数据。

    窗口化允许您在数据的滑动窗口上应用转换。

    每次窗口在源 DStream 上滑动时,落在窗口内的源 RDD 会被组合并操作以生成窗口化 DStream 的 RDD。所以基本上你可以组合来自多个批处理间隔的数据

    【讨论】:

      猜你喜欢
      • 2019-04-02
      • 1970-01-01
      • 2023-03-13
      • 2016-08-14
      • 1970-01-01
      • 2019-07-07
      • 1970-01-01
      • 2019-10-11
      • 1970-01-01
      相关资源
      最近更新 更多