【问题标题】:Apache Beam KafkaIO consumers in consumer group reading same message消费者组中的 Apache Beam KafkaIO 消费者阅读相同的消息
【发布时间】: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


    【解决方案1】:

    我认为您设置的选项不能保证跨管道的消息传递不重复。

    • ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG:这是 flag for the Kafka consumer 不适用于 Beam 管道本身。似乎这是尽力而为且定期进行,因此您可能仍会在多个管道中看到重复项。

    • withReadCommitted():这只是意味着 Beam 不会读取未提交的消息。同样,它不会防止跨多个管道重复。

    请参阅here,了解 Beam 源用于确定 Kafka 源起点的协议。

    为了保证不重复交付,您可能必须阅读不同主题或不同订阅的内容。

    【讨论】:

    • 即使我面临这个问题,我也想运行多个消费者实例,从同一主题读取。但是消息正在传递到所有正在运行的实例。调试配置后我发现,组名附加了一些唯一的前缀,每个实例都有唯一的组名。 ex - group.id = Reader-0_offset_consumer_559337182_my_group group.id = Reader-0_offset_consumer_559337345_my_group 所以每个实例都有唯一的 group.id 分配,这就是为什么消息被传递到所有实例。
    • @Aditya 我正在使用数据流管道,在我的情况下 group.id 对于两个消费者来说也是相同的,但消息仍然由两个消费者处理。你有解决办法吗?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-04-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多