【问题标题】:Apache Beam KafkaIO consumers in consumer group getting assigned unique group id消费者组中的 Apache Beam KafkaIO 消费者被分配了唯一的组 ID
【发布时间】:2020-07-20 17:58:03
【问题描述】:

我正在使用 DirectRunner 运行多个 Apache Beam KafkaIO 实例,它们从同一主题中读取。但是消息正在传递到所有正在运行的实例。在看到我发现的 Kafka 配置后,组名会附加一些唯一的前缀,并且每个实例都有唯一的组名。

  1. group.id = Reader-0_offset_consumer_559337182_my_group
  2. group.id = Reader-0_offset_consumer_559337345_my_group

因此,每个实例都分配了唯一的 group.id,这就是消息被传递到所有实例的原因。

pipeline.apply("ReadFromKafka", KafkaIO.<String, String>read().withReadCommitted()
            .withConsumerConfigUpdates(
                    new ImmutableMap.Builder<String, Object>().put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true)
                            .put(ConsumerConfig.GROUP_ID_CONFIG, "my_group")
                            .put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 5).build())
            .withKeyDeserializer(StringDeserializer.class).withValueDeserializer(StringDeserializer.class)
            .withBootstrapServers(servers).withTopics(Collections.singletonList(topicName)).withoutMetadata()

那么我必须提供什么配置才能使组中的所有消费者都不会阅读相同的消息

【问题讨论】:

  • 使用 DirectRunner 运行多个 KafkaIO 实例并读取同一主题的原因是什么?
  • @AlexeyRomanenko,我们没有使用 GCP,而是在我们自己的裸机上运行它。所以我们不能使用数据流。所以我们想通过在 k8s pod 中部署来扩展并增加 pod 的数量。但是我在这里看到的问题是,因为每个实例都被分配了一个唯一的 groupId,所以每当我发送消息时,消息都会发送到所有组/实例。希望这可以澄清问题
  • 我不建议您在生产中使用 DirectRunner 来处理大量数据,因为该运行器应该主要用于测试,它在管道运行期间包含并执行许多额外的检查,这就是为什么与其他跑步者相比,它可能会很慢。您是否可以选择在分布式 Spark 或 Flink 集群上使用 Spark 或 Flink 运行器?
  • @AlexeyRomanenko 不,目前我们没有使用 Flink 的 Spark 的选项。另外,请恢复反对票,因为它是有效的场景
  • 我没有投反对票,但我给你的帖子 +1。我希望人们可以有不同的情况,我只是建议如何更好地使用它。

标签: apache-beam apache-beam-kafkaio


【解决方案1】:

是的,这是因为组名附加了一些唯一的前缀,并且每个实例都有唯一的组名。因为这个kafka不知道你是否再启动一个实例。因此,相同的消息会传递给所有消费者。

因此,我能想到的一种解决方法是,与其给出主题并让 beam 计算出所有分区的消费者数量,不如使用 DirectRunner 为 apache beam KafkaIO 的每个实例明确给出主题分区。

您必须将TopicPartition 类型的List 传递给方法withTopicPartitions

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())
                .withTopicPartitions(Arrays.asList(new TopicPartition(topicName, 0)))
                .withKeyDeserializer(StringDeserializer.class)
                .withValueDeserializer(StringDeserializer.class)
                .withoutMetadata();

上面的代码只会读取来自partition 0的消息。因此,通过这种方式,您可以启动同一程序的多个实例,而无需将相同的消息传递给所有消费者

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-05-14
    • 2017-05-13
    • 1970-01-01
    • 2020-03-10
    • 1970-01-01
    • 2020-09-19
    • 1970-01-01
    • 2022-10-18
    相关资源
    最近更新 更多