【问题标题】:Apache Beam KafkaIO mention topic partition instead of topic nameApache Beam KafkaIO 提到主题分区而不是主题名称
【发布时间】:2020-07-05 13:18:01
【问题描述】:

Apache Beam KafkaIO 支持 kafka 消费者仅从指定分区读取。我有以下代码。

KafkaIO.<String, String>read()
                .withCreateTime(Duration.standardMinutes(1))
                .withReadCommitted()
                .withBootstrapServers(endPoint)
                .withConsumerConfigUpdates(new ImmutableMap.Builder<String, Object>()
                        .put(ConsumerConfig.GROUP_ID_CONFIG, groupName)
                        .put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 5)
                        .put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest")
                        .build())
                .commitOffsetsInFinalize()
                .withTopicPartitions(List<TopicPartitions>)

我有以下 2 个问题。

  1. 如何从 kafka 获取分区名称?怎么在 kafkaIO 中提及?
  2. Apache Beam 生成的 kafka 消费者数量是否等于创建 kafka 消费者时提到的分区列表?

【问题讨论】:

    标签: apache-beam apache-beam-io apache-beam-kafkaio


    【解决方案1】:

    我自己找到了答案。

    如何让 kafkaIO 从特定分区读取数据?

    kafkaIO 有 withTopicPartitions(List&lt;TopicPartitions&gt;) 方法,它接受 TopicPartition 对象的列表。

    主题分区被命名为从零开始的连续数字。因此,以下应该工作

    KafkaIO.<String, String>read()
                    .withCreateTime(Duration.standardMinutes(1))
                    .withReadCommitted()
                    .withBootstrapServers(endPoint)
                    .withConsumerConfigUpdates(new ImmutableMap.Builder<String, Object>()
                            .put(ConsumerConfig.GROUP_ID_CONFIG, groupName)
                            .put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 5)
                            .put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest")
                            .build())
                    .commitOffsetsInFinalize()
                    .withTopicPartitions(Arrays.asList(new TopicPartition(topicName, 0),new TopicPartition(topicName, 1),new TopicPartition(topicName, 2)))
    

    要测试它,请使用kafkacat 和以下命令

    kafkacat -P -b localhost:9092 -t sample -p 0 - 此命令生成到指定分区。

    Apache Beam 生成的 kafka 消费者数量是否等于创建 kafka 消费者时提到的分区列表?

    它将生成一个消费者组,其消费者数量与在构建 kafka Producer 对象期间明确提到的分区数量一样。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-08-31
      • 2022-01-10
      • 1970-01-01
      • 2015-03-05
      相关资源
      最近更新 更多