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

Spring Boot整合spring-kafka:如何实现生产者感知消费者消息确认?

实现Spring Kafka消费者确认后通知生产者的方案

Kafka本身没有原生支持消费者确认消息后主动通知生产者的机制,需要通过自定义逻辑来实现,以下是两种可行的落地方案:


方案一:基于额外通知Topic实现(实时性高)

核心思路是:消费者完成业务处理并确认消息后,向专门的通知Topic发送一条确认消息;生产者端监听这个Topic,根据消息唯一标识匹配对应的回调逻辑。

消费者端代码改造

@KafkaListener(topics = "chat-group-1")
public void listenChatGroup1(ConsumerRecord<String, String> record, 
                             Acknowledgment ack, 
                             KafkaTemplate<String, String> kafkaTemplate) {
    // 执行业务逻辑
    doWork();
    
    // 确认Kafka消息
    ack.acknowledge();
    
    // 发送确认通知到专门的Topic,携带原消息的唯一标识(这里用消息key,也可以用自定义的messageId)
    String messageUniqueId = record.key();
    kafkaTemplate.send("consumer-ack-notify", messageUniqueId, "ACK_SUCCESS");
}

生产者端实现回调感知

@Component
public class ProducerService {
    private final KafkaTemplate<String, String> kafkaTemplate;
    // 分布式场景下替换为Redis等分布式存储,避免本地Map无法跨实例共享
    private final ConcurrentHashMap<String, Runnable> ackCallbackMap = new ConcurrentHashMap<>();

    public ProducerService(KafkaTemplate<String, String> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    // 封装消息发送方法,注册回调逻辑
    public void sendMessage(String topic, String messageId, String content) {
        kafkaTemplate.send(topic, messageId, content);
        // 注册消费者确认后的回调逻辑
        ackCallbackMap.put(messageId, () -> {
            // 这里就是你需要的onConsumerAcknowledge逻辑
            System.out.println("消息[" + messageId + "]已被消费者确认处理");
            // 可执行后续操作:更新业务状态、触发下游流程等
        });
    }

    // 监听消费者确认通知Topic
    @KafkaListener(topics = "consumer-ack-notify")
    public void handleConsumerAck(ConsumerRecord<String, String> record) {
        String messageId = record.key();
        Runnable callback = ackCallbackMap.remove(messageId);
        if (callback != null) {
            callback.run();
        }
    }
}

方案注意点

  1. 消息唯一标识要确保全局唯一,避免回调匹配错误;
  2. 通知Topic的消息要做幂等处理,防止消费者重复发送确认导致回调重复执行;
  3. 分布式部署时,必须用Redis等分布式存储替换本地ConcurrentHashMap,保证回调逻辑能跨实例被触发。

方案二:基于数据库状态流转实现(可靠性高)

核心思路是:生产者发送消息时,在数据库插入一条状态为「待处理」的记录;消费者完成处理并确认消息后,更新该记录的状态为「已确认」;生产者通过轮询或数据库变更监听来触发回调。

生产者端发送消息

@Component
public class ProducerService {
    private final KafkaTemplate<String, String> kafkaTemplate;
    private final MessageRecordMapper messageRecordMapper;

    public ProducerService(KafkaTemplate<String, String> kafkaTemplate, MessageRecordMapper messageRecordMapper) {
        this.kafkaTemplate = kafkaTemplate;
        this.messageRecordMapper = messageRecordMapper;
    }

    public void sendMessage(String topic, String content) {
        String messageId = UUID.randomUUID().toString();
        // 插入数据库,标记状态为待处理
        MessageRecord record = new MessageRecord();
        record.setMessageId(messageId);
        record.setContent(content);
        record.setStatus("PENDING");
        record.setCallbackHandled(false);
        messageRecordMapper.insert(record);
        
        // 发送消息到业务Topic
        kafkaTemplate.send(topic, messageId, content);
    }

    // 定时轮询数据库,触发已确认消息的回调
    @Scheduled(fixedRate = 5000)
    public void checkAckedMessages() {
        List<MessageRecord> ackedRecords = messageRecordMapper.selectByStatusAndCallback("ACKED", false);
        for (MessageRecord record : ackedRecords) {
            // 执行回调逻辑
            System.out.println("消息[" + record.getMessageId() + "]已被消费者确认");
            // 标记回调已处理,避免重复触发
            record.setCallbackHandled(true);
            messageRecordMapper.updateById(record);
        }
    }
}

消费者端代码改造

@KafkaListener(topics = "chat-group-1")
public void listenChatGroup1(ConsumerRecord<String, String> record, 
                             Acknowledgment ack, 
                             MessageRecordMapper messageRecordMapper) {
    doWork();
    ack.acknowledge();
    
    // 更新数据库状态为已确认
    String messageId = record.key();
    MessageRecord updateRecord = new MessageRecord();
    updateRecord.setMessageId(messageId);
    updateRecord.setStatus("ACKED");
    messageRecordMapper.updateById(updateRecord);
}

方案优化

如果需要更高的实时性,可以用Canal等Binlog监听工具替代定时轮询,实时感知数据库状态变更并触发回调。


方案对比

方案类型优点缺点
额外通知Topic实时性高、无需依赖数据库需维护额外Topic、需处理幂等
数据库状态流转可靠性高、状态持久化实时性依赖轮询频率或Binlog、复杂度稍高

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 22:09:18