【发布时间】:2016-08-03 04:42:56
【问题描述】:
我正在尝试阅读另一个团队设置的 Kafka 主题。该主题在多个分区之间保持平衡。我的意思是每个新行都发送到一个单独的主题。一条消息是多行,因此消息在两个分区之间拆分。
例如:
分区 1:
“消息 1:详细信息 1 详细信息 1”
"message2: details2 details2"
分区 2:
“详情 1 详情 1”
“细节2细节2”
当我使用createDirectStream(ssc, kafkaparams, fromoffsets, messagehandler) 阅读该主题时,我会按照上面显示的顺序获得 RDD。
我想做的是:
“消息 1:详细信息 1 详细信息 1”
“详情 1 详情 1”
“消息 2:详细信息 2 详细信息 2”
“详情2详情2”
感谢我收到的任何帮助。
【问题讨论】:
-
问题是什么?有什么问题,你试过什么?
-
不能轻松地做到这一点——而不是使用您的数据。它们需要在同一个分区中,否则无法保证排序,或者即使在
foreachRDD中的同一个RDD中也能看到它们。理论上你可以做到这一点的唯一方法是拥有一个 seqID 或其他东西,你可以在记录被无序发送后使用它来重新排序记录。但您似乎没有任何此类数据。 -
对不起,问题是如何按顺序处理一个主题中的多个分区。当我使用 createStream 或 createDirectStream 读取数据时,它并没有按照添加到 kafka 的顺序出现。如何按添加顺序从分区中提取?我已经尝试了这两个流命令,但我无法从 Kafka 中找到有关 spark 流的相关信息。
-
@DavidGriffin 感谢您的回复。真不幸……我将不得不联系将数据发送到 kafka 的团队。
标签: scala apache-spark streaming apache-kafka spark-streaming