【发布时间】: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