【问题标题】:Spring kafka Batch Listener- commit offsets manually in BatchSpring kafka Batch Listener - 在 Batch 中手动提交偏移量
【发布时间】:2017-12-13 16:22:16
【问题描述】:

我正在实现 spring kafka 批处理侦听器,它从 Kafka 主题读取消息列表并将数据发布到 REST 服务。 我想了解在 REST 服务出现故障的情况下的偏移管理,不应提交批处理的偏移,并且应为下一次轮询处理消息。我已经阅读了 spring kafka 文档,但是在理解 Listener Error Handler 和 Seek to current container error handlers in batch之间的区别时存在混淆。我使用的是 spring-boot-2.0.0.M7 版本,下面是我的代码。

Listener Config:

@Bean
    KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());

        factory.setConcurrency(Integer.parseInt(env.getProperty("spring.kafka.listener.concurrency")));
        // factory.getContainerProperties().setPollTimeout(3000);
        factory.getContainerProperties().setBatchErrorHandler(kafkaErrorHandler());

        factory.getContainerProperties().setAckMode(AckMode.BATCH);
        factory.setBatchListener(true);
        return factory;
    }
@Bean
    public Map<String, Object> consumerConfigs() {
        Map<String, Object> propsMap = new HashMap<>();
        propsMap.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, env.getProperty("spring.kafka.bootstrap-servers"));
        propsMap.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG,
                env.getProperty("spring.kafka.consumer.enable-auto-commit"));
        propsMap.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG,
                env.getProperty("spring.kafka.consumer.auto-commit-interval"));
        propsMap.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, env.getProperty("spring.kafka.session.timeout"));
        propsMap.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        propsMap.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        propsMap.put(ConsumerConfig.GROUP_ID_CONFIG, env.getProperty("spring.kafka.consumer.group-id"));
        return propsMap;
    }

Listener Class:

@KafkaListener(topics = "${spring.kafka.consumer.topic}", containerFactory = "kafkaListenerContainerFactory")
    public void listen(List<String> payloadList) throws Exception {
        if (payloadList.size() > 0)
            //Post to the service
    }

Kafka Error Handler:

public class KafkaErrorHandler implements BatchErrorHandler {

    private static Logger LOGGER = LoggerFactory.getLogger(KafkaErrorHandler.class);

    @Override
    public void handle(Exception thrownException, ConsumerRecords<?, ?> data) {
        LOGGER.info("Exception occured while processing::" + thrownException.getMessage());

            }

}

如何处理 Kafka 侦听器,以便在处理一批记录期间发生某些事情,我不会丢失数据。

【问题讨论】:

  • 嗨。你能找到方法吗?

标签: spring spring-boot spring-kafka


【解决方案1】:

使用 Apache Kafka,我们永远不会丢失数据。分区日志中确实有一个偏移量来寻找任意位置。

另一方面,当我们从一个分区消费记录时,不需要提交它们的偏移量——当前消费者将状态保存在内存中。当当前消费者死亡时,我们只需要为同一组中的其他新消费者提交。与错误无关,当前消费者总是继续轮询其当前内存偏移后面的新数据。

因此,要重新处理同一消费者中的相同数据,我们必须使用seek 操作将消费者移回所需位置。这就是为什么 Spring Kafka 引入了SeekToCurrentErrorHandler

这允许实现查找所有未处理的主题/分区,以便下一次轮询检索当前记录(以及剩余的其他记录)。 SeekToCurrentErrorHandler 正是这样做的。

https://docs.spring.io/spring-kafka/reference/htmlsingle/#_seek_to_current_container_error_handlers

【讨论】:

  • SeekToCurrentErrorHandler 是否也处理批处理操作?因为我在包 org.springframework.kafka.listener 中看不到 SeekToCurrentBatchErrorHandler 类。另外,如果我将错误处理程序设置为 seektocurrenterrorhandler,请告诉我日志是如何发生的。
  • 有一个SeekToCurrentBatchErrorHandler,但自从2.1 已经存在:docs.spring.io/spring-kafka/docs/2.1.0.RELEASE/reference/html/…。您可以将其代码升级或复制/粘贴到您自己的课程中
猜你喜欢
  • 1970-01-01
  • 2018-10-26
  • 2020-12-28
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-04-06
相关资源
最近更新 更多