You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在不使用@KafkaListener时手动确认Kafka消息

动态Kafka容器实现手动消息确认

问题描述

之前使用@KafkaListener注解时,可通过消息header中的Acknowledgment对象执行手动确认,但现在通过动态创建ConcurrentMessageListenerContainer的方式消费消息,需要实现消息处理完成后手动确认的功能,现有核心代码如下:

public ConcurrentMessageListenerContainer<String, String> newContainer(String topic, int partition,
        int idlePeriod) {
    this.factory.setConsumerFactory(consumerFactory);
    this.factory.getContainerProperties().setIdleEventInterval(idlePeriod * 1000L);
    // this.factory.getContainerProperties().setAckMode(AckMode.MANUAL_IMMEDIATE);
    ConcurrentMessageListenerContainer<String, String> container = this.factory
            .createContainer(new TopicPartitionOffset(topic, partition));
    container.setupMessageListener((MessageListener<String, String>) record -> {
        // 消费消息
        kafkaService.proccessorConsumer(record);
    });
    this.containers.put("provisioning_group", container);
    container.start();
    return container;
}

解决方案

要实现手动确认,需完成以下3点修改:

1. 启用手动确认模式

取消注释setAckMode代码,将确认模式设置为AckMode.MANUAL_IMMEDIATE(立即提交偏移量)或AckMode.MANUAL(批量提交偏移量):

this.factory.getContainerProperties().setAckMode(AckMode.MANUAL_IMMEDIATE);

2. 替换监听器类型

MessageListener不提供确认对象,需改用AcknowledgingMessageListener,它的onMessage方法会传入ConsumerRecord和Acknowledgment两个参数,用于手动确认操作。

3. 消息处理完成后执行确认

在消息处理逻辑执行完成后,调用acknowledge()方法提交偏移量,同时可添加异常处理逻辑避免消息丢失。

修改后的完整newContainer方法

public ConcurrentMessageListenerContainer<String, String> newContainer(String topic, int partition,
        int idlePeriod) {
    this.factory.setConsumerFactory(consumerFactory);
    this.factory.getContainerProperties().setIdleEventInterval(idlePeriod * 1000L);
    // 启用手动立即确认模式
    this.factory.getContainerProperties().setAckMode(AckMode.MANUAL_IMMEDIATE);
    ConcurrentMessageListenerContainer<String, String> container = this.factory
            .createContainer(new TopicPartitionOffset(topic, partition));
    // 设置随机groupId,避免重复消费冲突
    String groupId = UUID.randomUUID().toString();
    container.getContainerProperties().setGroupId(groupId);
    // 使用AcknowledgingMessageListener替代MessageListener
    container.setupMessageListener((AcknowledgingMessageListener<String, String>) (record, acknowledgment) -> {
        try {
            // 执行消息处理逻辑
            kafkaService.proccessorConsumer(record);
            // 消息处理完成后手动确认提交偏移量
            acknowledgment.acknowledge();
        } catch (Exception e) {
            // 异常场景处理:可选择不确认,让Kafka重新投递消息
            log.error("消息处理失败", e);
        }
    });
    this.containers.put(groupId, container);
    container.start();
    return container;
}

关键说明

  • AckMode.MANUAL_IMMEDIATE:调用确认方法后立即将偏移量提交到Kafka,适合需确保消息处理完成后立即提交的场景。
  • AckMode.MANUAL:调用确认方法后不会立即提交,会等待容器的批量提交间隔统一提交,适合批量处理场景。
  • 异常处理:若消息处理失败,不调用acknowledge()可让Kafka重新投递该消息,但需配合消费者重试配置,避免无限循环。

内容的提问来源于stack exchange,提问作者kikicoder

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.16 00:10:36