【问题标题】:What do foreachBatches contain in a streaming query from multiple Kafka topics?foreachBatches 在来自多个 Kafka 主题的流式查询中包含什么?
【发布时间】:2019-11-21 03:57:35
【问题描述】:

假设DataStreamReader 配置为订阅多个这样的主题(参见here):

// Subscribe to multiple topics
spark
  .readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "host1:port1,host2:port2")
  .option("subscribe", "topic1,topic2,topic3")

当我在此基础上使用foreachBatch 时,批次将包含什么?

  • 每批只包含来自一个主题的消息?
  • 或者一个批次可以包含来自不同主题的消息吗?

在我的用例中,我希望批量处理仅来自一个主题的消息。可以这样配置吗?

【问题讨论】:

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


    【解决方案1】:

    引用Structured Streaming + Kafka Integration Guide (Kafka broker version 0.10.0 or higher)中的官方文档:

    //订阅多个主题

    ...
    .option("subscribe", "topic1,topic2")
    

    上面的代码是底层 Kafka 消费者(流式查询)订阅的内容。

    当我在此基础上使用 foreachBatch 时,批次将包含什么?

    • 每批只包含一个主题的消息?

    这是正确的答案。

    我希望批量处理仅来自一个主题的消息。可以这样配置吗?

    这也记录在Structured Streaming + Kafka Integration Guide (Kafka broker version 0.10.0 or higher):

    源中的每一行都有以下架构:

    ...

    主题

    换句话说,输入数据集将具有topic 列,其中包含给定行(记录)来自的主题名称。

    为了获得“仅来自一个主题的消息的批次”,您只需 filterwhere 使用一个主题,例如

    val messages: DataFrame = ...
    assert(messages.isStreaming)
    
    messages
      .writeStream
      .foreachBatch { case (df, batchId) =>
        val topic1Only = df.where($"topic" === "topic1")
        val topic2Only = df.where($"topic" === "topic2")
        ...
      }
    

    【讨论】:

    • 至于我的问题的第二部分:通过配置,我的意思是避免过滤,因为批次需要完全处理或根本不处理。但是,当批次仅包含来自一个主题的消息时,一切都很好。谢谢你的回答。
    • @jacek 有没有办法以编程方式编写它而不是为每个主题指定一个 val?我正在考虑循环遍历似乎效率低下的主题列表
    • @collarblind 只需使用topic 字段,此行就会“路由”到该主题。
    • @JacekLaskowski 我的问题令人困惑。如果我有主题val topics=Seq("t1","t2") 而我的foreachBatch 有这个,topics.map(t => df.where($"topic" === "t1").write()。和你的代码一样吗?
    • 既然你有内置的数据源,为什么你在写 Kafka 的时候foreachBatch
    【解决方案2】:

    批处理将包含来自您的消费者订阅的所有主题(我会说是分区)的消息。

    【讨论】:

    • 感谢您的回答。它是基于观察还是记录在某处?问题是关于多个主题(而不是关于分区)。
    • @Beryllium 在幕后,消费者订阅给定主题的特定分区。如果一个消费者组中只有一个消费者,那么它会订阅所有的分区。
    猜你喜欢
    • 1970-01-01
    • 2018-12-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-03-19
    • 2019-05-05
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多