【发布时间】:2020-05-16 20:09:23
【问题描述】:
我在数据流中使用 KafkaIO 来读取来自一个主题的消息。我使用以下代码。
KafkaIO.<String, String>read()
.withReadCommitted()
.withBootstrapServers(endPoint)
.withConsumerConfigUpdates(new ImmutableMap.Builder<String, Object>()
.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true)
.put(ConsumerConfig.GROUP_ID_CONFIG, groupName)
.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 8000).put(ConsumerConfig.REQUEST_TIMEOUT_MS_CONFIG, 2000)
.build())
// .commitOffsetsInFinalize()
.withTopics(Collections.singletonList(topicNames))
.withKeyDeserializer(StringDeserializer.class)
.withValueDeserializer(StringDeserializer.class)
.withoutMetadata();
我使用直接运行器在本地运行数据流程序。一切运行良好。我并行运行同一程序的另一个实例,即另一个消费者。现在我在管道处理中看到重复的消息。
虽然我提供了消费者组 id,但以相同的消费者组 id(同一程序的不同实例)启动另一个消费者不应该处理另一个消费者处理的相同元素,对吧?
使用 dataflow runner 的结果如何?
【问题讨论】:
标签: google-cloud-dataflow apache-beam-io apache-beam apache-beam-kafkaio