【问题标题】:Spark Streaming Kafka火花流卡夫卡
【发布时间】: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


【解决方案1】:

如果每个分区内的顺序得到保证,以便分区 1 中的元素 x 与分区 2 中的元素 x 相关,您可以根据分区号和每个分区迭代器 (zipWithIndex) 内的元素索引对 RDD 元素进行排序。

这将允许您跨分区“重新同步”

【讨论】:

    猜你喜欢
    • 2018-09-15
    • 2018-08-13
    • 2018-02-24
    • 2023-03-19
    • 2018-08-15
    • 1970-01-01
    • 2019-04-11
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多