【问题标题】:Is it possible to specify Kafka bootstrap servers for two different clusters for a single pipeline using KafkaIO.read?是否可以使用 KafkaIO.read 为单个管道的两个不同集群指定 Kafka 引导服务器?
【发布时间】:2020-08-28 22:27:12
【问题描述】:

我目前正在使用 Google Cloud Dataflow 和 Apache Beam 来使用来自 Kafka 主题的消息,该主题存在于两个不同的 Kafka 集群中,两个集群包含相同的主题名称但主题中的数据不同。 Kafka 集群是分开的,因为它们包含来自不同区域的数据。

我只是想知道是否可以通过在单个 KafkaIO.read 数据流管道步骤中列出两个集群的所有引导服务器来使用来自两个集群的数据?

.withBootstrapServers("CLUSTER1_SERVER:PORT,CLUSTER2_SERVER:PORT");

我正在阅读有关 Kafka 引导服务器的文档,我不清楚在连接到引导服务器后,是否只会从第一个成功的引导服务器连接集群中使用消息,或者它是否会尝试提供的所有引导服务器并从找到的所有集群中消费。如果是前者,那么我将需要创建第二个 Dataflow 管道来处理来自第二个集群的消息,但如果我可以在一个管道中处理来自两个集群的消息会容易得多。

任何信息将不胜感激。

【问题讨论】:

  • 您能分享您遵循的文档和Dataflow版本吗?谢谢!
  • @muscat 我关注了此页面上的文档:kafka.apache.org/documentation,我目前使用的 Dataflow/Apache Beam 版本是 2.18

标签: java google-cloud-platform apache-kafka google-cloud-dataflow apache-beam


【解决方案1】:

Beam KafkaIO 只是将此标志传递给 Kafka 的 ConsumerConfig 的 BOOTSTRAP_SERVERS_CONFIG flag。我认为这个参数是为了从同一个 Kafka 集群中传入多个代理进行故障转移。不适用于从不同的 Kafka 集群传入服务器。有关 Kafka 架构的详细信息,请参阅here。我怀疑当您从多个集群中指定服务器时,它只会选择第一个活动的。

【讨论】:

    【解决方案2】:

    我认为通过同一个 KafkaIO 实例从不同的集群中读取数据并不是一个好主意,因为在后台,它使用 KafkaConsumer 来读取消息,并且它只能从一个集群中读取,这是不打算用于故障转移情况。另外,实际上KafkaIO 中使用了两个 Kafka 消费者(一个用于消息,另一个用于偏移量),因此可能会更糟,结果将无法预测。

    同时,您可以为不同的集群拥有两个KafkaIO 源,然后通过键或下游的任何其他属性连接消息。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2017-05-06
      • 1970-01-01
      • 2018-08-19
      • 2019-11-08
      • 1970-01-01
      • 1970-01-01
      • 2020-04-06
      相关资源
      最近更新 更多