【问题标题】:Reading the same message several times from Kafka从 Kafka 多次读取同一条消息
【发布时间】:2017-01-24 22:47:49
【问题描述】:

我使用Spring Kafka API来实现Kafka消费者手动偏移管理:

@KafkaListener(topics = "some_topic")
public void onMessage(@Payload Message message, Acknowledgment acknowledgment) {
    if (someCondition) {
        acknowledgment.acknowledge();
    }
}

在这里,我希望消费者仅在 someCondition 成立时提交偏移量。否则消费者应该休眠一段时间并再次阅读相同的消息

卡夫卡配置:

@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
    ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(consumerConfig());
    factory.getContainerProperties().setAckMode(MANUAL);
    return factory;
}

private Map<String, Object> consumerConfig() {
    Map<String, Object> props = new HashMap<>();
    ...
    props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
    ...
    return props;
}

在当前配置下,如果someCondition == false,消费者不会提交偏移量,但仍会读取下一条消息。如果没有执行 Kafka acknowledgement,有没有办法让消费者重新阅读消息?

【问题讨论】:

    标签: java spring apache-kafka kafka-consumer-api spring-kafka


    【解决方案1】:

    您可以停止并重新启动容器,它将被重新发送。

    随着即将发布的 1.1 版本,您可以seek to the required offset,它将被重新发送。

    但是如果它们已经被检索到,你仍然会首先看到后面的消息,因此你也必须丢弃它们。

    second milestone 具有该功能,我们预计它将在下周发布。

    【讨论】:

    • 谢谢。现在我怀疑我是否将手动偏移管理用于其预期目的。似乎应该在消费者失败的情况下使用功能重新读取消息,而不是在使用消息的组件失败时。
    • 嗨@Gary Russel,你知道是否有一些关于在Spring Boot中使用seek的例子。我正在尝试一些类似的东西来收听一个主题,但要查看分区中消息的响应:stackoverflow.com/questions/48349993/…
    • Seek 不是该用例的正确解决方案;看看我的答案。
    【解决方案2】:

    正如@Gary 已经指出的那样,您的方向是正确的,seek() 是这样做的方法。当我遇到这个问题时,我今天找不到它的代码示例。这是任何想要解决问题的人的代码。

    public class Receiver implements AcknowledgingMessageListener<Integer, String>, ConsumerSeekAware {
    
        private ConsumerSeekCallback consumerSeekCallback;
    
    
        @Override
        public void onMessage(ConsumerRecord<Integer, String> record, Acknowledgment acknowledgment) {
    
            if (/*some condition*/) {
                //process
                acknowledgment.acknowledge(); //send ack
            } else {
    
                consumerSeekCallback.seek("your.topic", record.partition(), record.offset());
    
            }
        }
    
        @Override
        public void registerSeekCallback(ConsumerSeekCallback consumerSeekCallback) {
            this.consumerSeekCallback = consumerSeekCallback;
        }
    
        @Override
        public void onPartitionsAssigned(Map<TopicPartition, Long> map, ConsumerSeekCallback consumerSeekCallback) {
    
            // nothing is needed here for this program
        }
    
        @Override
        public void onIdleContainer(Map<TopicPartition, Long> map, ConsumerSeekCallback consumerSeekCallback) {
    
            // nothing is needed here for this program
        }
    
    }
    

    【讨论】:

    • 谢谢它帮助了我
    【解决方案3】:

    您可以尝试使用nack(long sleep),其中唯一的参数代表sleep interval ms,以实现上述行为。

    来自Spring for Apache Kafka documentation

    2.3版开始,确认界面有两个 附加方法 nack(long sleep) 和 nack(int index, long sleep)。 第一个用于记录侦听器,第二个用于批处理 听众。为您的侦听器类型调用错误的方法将抛出 一个 IllegalStateException。

    将上述信息应用到我们得到的代码示例中:

    @Component
    @Slf4j
    public class ExampleConsumer {
        private boolean nonError = false;
        
        @KafkaListener(topics = "topic_name")
        private void consumeSelectingMsgFromMailbox(ConsumerRecord<String, KafkaEventPojo> record, Acknowledgment ack) {
            log.info("Received record topic:{} partition:{} offset:{}", record.topic(), record.partition(), record.offset());
            
            if (nonError) {
                log.info("ACK: {}", offset);
                ack.acknowledge(); //send ack
                if (offset % 2 == 0)
                    nonError = false;
            } else {
                ack.nack(0); // immediate seek - no sleep time for consumer
                nonError = true;
            }
        }
    }
    
    

    配置如下:

    @Configuration
    @EnableKafka
    public class KafkaConsumerConfig {
        private ConcurrentKafkaListenerContainerFactory<String, KafkaEventPojo> factory;
    
        @Bean
        public Map<String, Object> consumerConfigs() {
            Map<String, Object> props = new HashMap<>();
            // ...
            props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
            // ...
    
            return props;
        }
    
        @Bean
        public ConsumerFactory<String, KafkaEventPojo> consumerFactory() {
            return new DefaultKafkaConsumerFactory<>(consumerConfigs());
        }
    
        @Bean
        public ConcurrentKafkaListenerContainerFactory<String, KafkaEventPojo> kafkaListenerContainerFactory() {
            if (this.factory == null) {
                this.factory =
                        new ConcurrentKafkaListenerContainerFactory<>();
                factory.setConsumerFactory(consumerFactory());
                factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
            }
    
            return this.factory;
        }
    

    示例产生:

    2020-07-31 17:05:19.275  INFO       17560 --- [ntainer#0-0-C-1] .d.g.a.a.k.c.InboxRetrievalEventConsumer : consumeSelectingMsgFromMailbox Received record topic:event_log_selectingMsgFromMailbox partition:1 offset:15
    2020-07-31 17:05:19.792  INFO       17560 --- [ntainer#0-0-C-1] .d.g.a.a.k.c.InboxRetrievalEventConsumer : consumeSelectingMsgFromMailbox Received record topic:event_log_selectingMsgFromMailbox partition:1 offset:15
    2020-07-31 17:05:19.793  INFO       17560 --- [ntainer#0-0-C-1] .d.g.a.a.k.c.InboxRetrievalEventConsumer : consumeSelectingMsgFromMailbox ACK: 15
    2020-07-31 17:05:19.805  INFO       17560 --- [ntainer#0-0-C-1] .d.g.a.a.k.c.InboxRetrievalEventConsumer : consumeSelectingMsgFromMailbox Received record topic:event_log_selectingMsgFromMailbox partition:1 offset:16
    2020-07-31 17:05:19.805  INFO       17560 --- [ntainer#0-0-C-1] .d.g.a.a.k.c.InboxRetrievalEventConsumer : consumeSelectingMsgFromMailbox ACK: 16
    2020-07-31 17:05:19.810  INFO       17560 --- [ntainer#0-0-C-1] .d.g.a.a.k.c.InboxRetrievalEventConsumer : consumeSelectingMsgFromMailbox Received record topic:event_log_selectingMsgFromMailbox partition:1 offset:17
    2020-07-31 17:05:20.313  INFO       17560 --- [ntainer#0-0-C-1] .d.g.a.a.k.c.InboxRetrievalEventConsumer : consumeSelectingMsgFromMailbox Received record topic:event_log_selectingMsgFromMailbox partition:1 offset:17
    2020-07-31 17:05:20.313  INFO       17560 --- [ntainer#0-0-C-1] .d.g.a.a.k.c.InboxRetrievalEventConsumer : consumeSelectingMsgFromMailbox ACK: 17
    2020-07-31 17:05:20.318  INFO       17560 --- [ntainer#0-0-C-1] .d.g.a.a.k.c.InboxRetrievalEventConsumer : consumeSelectingMsgFromMailbox Received record topic:event_log_selectingMsgFromMailbox partition:1 offset:18
    2020-07-31 17:05:20.318  INFO       17560 --- [ntainer#0-0-C-1] .d.g.a.a.k.c.InboxRetrievalEventConsumer : consumeSelectingMsgFromMailbox ACK: 18
    2020-07-31 17:05:20.322  INFO       17560 --- [ntainer#0-0-C-1] .d.g.a.a.k.c.InboxRetrievalEventConsumer : consumeSelectingMsgFromMailbox Received record topic:event_log_selectingMsgFromMailbox partition:1 offset:19
    2020-07-31 17:05:20.827  INFO       17560 --- [ntainer#0-0-C-1] .d.g.a.a.k.c.InboxRetrievalEventConsumer : consumeSelectingMsgFromMailbox Received record topic:event_log_selectingMsgFromMailbox partition:1 offset:19
    2020-07-31 17:05:20.828  INFO       17560 --- [ntainer#0-0-C-1] .d.g.a.a.k.c.InboxRetrievalEventConsumer : consumeSelectingMsgFromMailbox ACK: 19
    

    注意:KafkaEventPojo 是我的 POJO 实现,它按照我们的内部结构保存存储在 Kafka 中的记录数据 - 因此您可以根据需要进行更改。此外,上面的代码演示了将 nack 用于单个记录侦听器的用法。如果您需要批处理选项,您可以在提供的文档中找到如何执行此操作的示例。

    【讨论】:

      猜你喜欢
      • 2015-11-10
      • 2015-04-20
      • 1970-01-01
      • 2021-12-17
      • 1970-01-01
      • 2022-06-13
      • 1970-01-01
      • 2011-04-04
      • 1970-01-01
      相关资源
      最近更新 更多