【发布时间】:2020-07-20 17:58:03
【问题描述】:
我正在使用 DirectRunner 运行多个 Apache Beam KafkaIO 实例,它们从同一主题中读取。但是消息正在传递到所有正在运行的实例。在看到我发现的 Kafka 配置后,组名会附加一些唯一的前缀,并且每个实例都有唯一的组名。
- group.id = Reader-0_offset_consumer_559337182_my_group
- 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