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(); } } }
方案注意点
- 消息唯一标识要确保全局唯一,避免回调匹配错误;
- 通知Topic的消息要做幂等处理,防止消费者重复发送确认导致回调重复执行;
- 分布式部署时,必须用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
相关产品推荐
相关产品推荐

