23人参与 • 2026-09-13 • C/C++
今天我们来聊聊一个让很多开发者头疼的话题——mq消息丢失问题
有些小伙伴在工作中,一提到消息队列就觉得很简单,但真正遇到线上消息丢失时,排查起来却让人抓狂
其实,我在实际工作中,也遇到过mq消息丢失的情况
今天这篇文章,专门跟大家一起聊聊这个话题,希望对你会有所帮助
在深入解决方案之前,我们先搞清楚消息在哪几个环节可能丢失:

理解了问题根源,接下来我们看5种实用的解决方案
生产者发送消息后等待broker确认,确保消息成功到达
这是防止消息丢失的第一道防线

// rabbitmq生产者确认配置
@bean
public rabbittemplate rabbittemplate() {
rabbittemplate template = new rabbittemplate(connectionfactory);
template.setconfirmcallback((correlationdata, ack, cause) -> {
if (ack) {
// 消息成功到达broker
messagestatusservice.markconfirmed(correlationdata.getid());
} else {
// 发送失败,触发重试
retryservice.scheduleretry(correlationdata.getid());
}
});
return template;
}
// 可靠发送方法
public void sendreliable(string exchange, string routingkey, object message) {
string messageid = generateid();
// 先落库保存发送状态
messagestatusservice.savesendingstatus(messageid, message);
// 发送持久化消息
rabbittemplate.convertandsend(exchange, routingkey, message, msg -> {
msg.getmessageproperties().setdeliverymode(messagedeliverymode.persistent);
msg.getmessageproperties().setmessageid(messageid);
return msg;
}, new correlationdata(messageid));
}
将消息保存到磁盘,确保broker重启后消息不丢失
这是防止broker端消息丢失的关键

// 持久化队列配置
@bean
public queue orderqueue() {
return queuebuilder.durable("order.queue") // 队列持久化
.deadletterexchange("order.dlx") // 死信交换机
.build();
}
// 发送持久化消息
public void sendpersistentmessage(object message) {
rabbittemplate.convertandsend("order.exchange", "order.create", message, msg -> {
msg.getmessageproperties().setdeliverymode(messagedeliverymode.persistent); // 消息持久化
return msg;
});
}
// kafka持久化配置
@bean
public producerfactory<string, object> producerfactory() {
map<string, object> props = new hashmap<>();
props.put(producerconfig.acks_config, "all"); // 所有副本确认
props.put(producerconfig.retries_config, 3); // 重试次数
props.put(producerconfig.enable_idempotence_config, true); // 幂等性
returnnew defaultkafkaproducerfactory<>(props);
}
优点:
缺点:
消费者处理完消息后手动向broker发送确认,broker收到确认后才删除消息
这是保证消息不丢失的最后一道防线

// 手动确认消费者
@rabbitlistener(queues = "order.queue")
public void handlemessage(order order, message message, channel channel) {
long deliverytag = message.getmessageproperties().getdeliverytag();
try {
// 业务处理
orderservice.processorder(order);
// 手动确认
channel.basicack(deliverytag, false);
log.info("消息处理完成: {}", order.getorderid());
} catch (exception e) {
log.error("消息处理失败: {}", order.getorderid(), e);
// 处理失败,重新入队
channel.basicnack(deliverytag, false, true);
}
}
// 消费者容器配置
@bean
public simplerabbitlistenercontainerfactory containerfactory() {
simplerabbitlistenercontainerfactory factory = new simplerabbitlistenercontainerfactory();
factory.setacknowledgemode(acknowledgemode.manual); // 手动确认
factory.setprefetchcount(10); // 预取数量
factory.setconcurrentconsumers(3); // 并发消费者
return factory;
}
通过事务保证本地业务操作和消息发送的原子性,要么都成功,要么都失败

// 本地事务表方案
@transactional
public void createorder(order order) {
// 1. 保存订单到数据库
orderrepository.save(order);
// 2. 保存消息到本地消息表
localmessage localmessage = new localmessage();
localmessage.setbusinessid(order.getorderid());
localmessage.setcontent(json.tojsonstring(order));
localmessage.setstatus(messagestatus.pending);
localmessagerepository.save(localmessage);
// 3. 事务提交,本地业务和消息存储保持一致性
}
// 定时任务扫描并发送消息
@scheduled(fixeddelay = 5000)
public void sendpendingmessages() {
list<localmessage> pendingmessages = localmessagerepository.findbystatus(messagestatus.pending);
for (localmessage message : pendingmessages) {
try {
// 发送消息到mq
rabbittemplate.convertandsend("order.exchange", "order.create", message.getcontent());
// 更新消息状态为已发送
message.setstatus(messagestatus.sent);
localmessagerepository.save(message);
} catch (exception e) {
log.error("发送消息失败: {}", message.getid(), e);
}
}
}
// rocketmq事务消息
public void sendtransactionmessage(order order) {
transactionmqproducer producer = new transactionmqproducer("order_producer");
// 发送事务消息
message msg = new message("order_topic", "create", json.tojsonbytes(order));
transactionsendresult result = producer.sendmessageintransaction(msg, null);
if (result.getlocaltransactionstate() == localtransactionstate.commit_message) {
log.info("事务消息提交成功");
}
}
通过重试机制处理临时故障,通过死信队列处理最终无法消费的消息

// 重试队列配置
@bean
public queue orderqueue() {
return queuebuilder.durable("order.queue")
.withargument("x-dead-letter-exchange", "order.dlx") // 死信交换机
.withargument("x-dead-letter-routing-key", "order.dead")
.withargument("x-message-ttl", 60000) // 60秒后进入死信
.build();
}
// 死信队列配置
@bean
public queue orderdeadletterqueue() {
return queuebuilder.durable("order.dead.queue").build();
}
// 消费者重试逻辑
@rabbitlistener(queues = "order.queue")
public void handlemessagewithretry(order order, message message, channel channel) {
long deliverytag = message.getmessageproperties().getdeliverytag();
try {
orderservice.processorder(order);
channel.basicack(deliverytag, false);
} catch (temporaryexception e) {
// 临时异常,重新入队重试
channel.basicnack(deliverytag, false, true);
} catch (permanentexception e) {
// 永久异常,直接确认进入死信队列
channel.basicack(deliverytag, false);
log.error("消息进入死信队列: {}", order.getorderid(), e);
}
}
// 死信队列消费者
@rabbitlistener(queues = "order.dead.queue")
public void handledeadlettermessage(order order) {
log.warn("处理死信消息: {}", order.getorderid());
// 发送告警、记录日志、人工处理等
alertservice.sendalert("死信消息告警", order.tostring());
}
为了帮助大家选择合适的方案,我整理了详细的对比表:
| 方案 | 可靠性 | 性能影响 | 复杂度 | 适用场景 |
|---|---|---|---|---|
| 生产者确认 | 高 | 中 | 低 | 所有需要可靠发送的场景 |
| 消息持久化 | 中 | 中 | 低 | broker重启保护 |
| 消费者确认 | 高 | 低 | 中 | 确保消息被成功处理 |
| 事务消息 | 最高 | 高 | 高 | 强一致性要求的业务 |
| 重试+死信 | 高 | 低 | 中 | 处理临时故障和最终死信 |
初创项目/简单业务:
电商/交易系统:
大数据/日志处理:
金融/支付系统:
消息丢失问题是消息队列使用中的常见挑战,通过今天介绍的5种方案,我们可以构建一个可靠的消息系统:
有些小伙伴可能会问:“我需要全部使用这些方案吗?”
我的建议是:根据业务需求选择合适的组合
对于关键业务,建议至少使用前三种方案;对于普通业务,可以根据实际情况适当简化
记住,没有完美的方案,只有最适合的方案
以上为个人经验,希望能给大家一个参考,也希望大家多多支持代码网。
您想发表意见!!点此发布评论
版权声明:本文内容由互联网用户贡献,该文观点仅代表作者本人。本站仅提供信息存储服务,不拥有所有权,不承担相关法律责任。 如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 2386932994@qq.com 举报,一经查实将立刻删除。
发表评论