【RabbitMQ #11】 | 消费者可靠性

发布时间:2026/10/1 21:01:30
【RabbitMQ #11】 | 消费者可靠性 一、消费者确认机制 Consumer Acknowledgement为了确认消费者是否成功处理消息RabbitMQ 提供了消费者确认机制Consumer Acknowledgement。 当消费者处理消息结束后向 RabbitMQ 发送回执告知消息处理状态。回执分为 3 种ack成功处理消息RabbitMQ 从队列删除该消息nack消息处理失败RabbitMQ 可再次投递消息reject消息处理失败并拒绝该消息RabbitMQ 直接删除消息Spring AMQP 封装了确认功能配置文件指定 ack 模式共 3 种none自动 ack。消息投递到消费者立刻 ackMQ 直接删除消息。不推荐业务处理失败消息已经丢失无法补救。manual手动 ack。业务代码手动调用 ack/nack/reject灵活需要开发者自己处理异常。auto自动模式推荐。Spring 利用 AOP 环绕拦截消息处理逻辑业务正常执行 → 自动返回ack业务异常 → 自动返回nack消息校验 / 解析异常 → 自动返回reject重点Go amqp091-go SDK 没有 Spring 那种开箱即用的 auto 自动 AOP 模式Spring 的auto是框架封装能力原生 AMQP 协议底层只有两种autoAcktrue≈ Springnone消息投递直接 ackMQ 删消息不安全autoAckfalse≈ Springmanual手动 ack所有 ack/nack/reject 全部手写代码自己捕获 panic、区分异常类型模拟 Spring auto 逻辑Go 代码模拟 Spring auto 自动确认package main import ( log github.com/rabbitmq/amqp091-go ) func main() { conn, err : amqp091.Dial(amqp://guest:guest127.0.0.1:5672/) if err ! nil { log.Fatal(err) } defer conn.Close() ch, err : conn.Channel() if err ! nil { log.Fatal(err) } defer ch.Close() // autoAckfalse 开启手动ack模拟Spring auto模式 msgs, err : ch.Consume( test_queue, // queue , // consumer false, // autoAckfalse【核心】 false, // exclusive false, // noLocal false, // noWait nil, // args ) if err ! nil { log.Fatal(err) } forever : make(chan struct{}) go func() { for d : range msgs { // 模拟Spring AOP环绕捕获panic func() { defer func() { // 捕获panic if r : recover(); r ! nil { log.Printf(panic异常: %v, r) // 业务panic → Nack消息重新入队 _ d.Nack(false, true) } }() // 业务处理逻辑 err : handleMsg(d.Body) if err ! nil { // 判断异常类型 if isMsgCheckErr(err) { // 消息校验异常 → reject丢弃消息不重投 log.Println(消息校验失败reject丢弃) _ d.Reject(false) } else { // 业务异常 → nack消息重新投递 log.Println(业务异常nack重入队列) _ d.Nack(false, true) } } else { // 正常处理成功 → ackMQ删除消息 log.Println(处理成功ack) _ d.Ack(false) } }() } }() log.Println(消费者启动等待消息) -forever } // handleMsg 业务处理函数 func handleMsg(body []byte) error { // 业务逻辑 return nil } // isMsgCheckErr 判断是否为消息校验类异常自行实现 // 例如json解析失败、参数非法这类不可恢复的消息异常 func isMsgCheckErr(err error) bool { return false }二、消费失败处理 本地重试问题消费失败后如果直接 nack (requeuetrue)消息会不断重新入队无限重试造成 MQ 消息堆积、压力飙升。 解决方案本地重试在当前 goroutine 内循环重试业务重试全部失败之后再决定消息去向不要把消息丢回 MQ 无限循环。Go 本地重试完整代码对齐 Spring 重试配置package main import ( fmt github.com/streadway/amqp log time ) // 本地重试配置对应spring配置 const ( maxAttempts 3 // 最大重试次数对应max-attempts initialInterval 1 * time.Second // 初始等待时间initial-interval multiplier 1.0 // 时间倍数 multiplier ) // 模拟业务消费函数 func handleMsg(msg *amqp.Delivery) error { // 模拟业务异常 return fmt.Errorf(业务处理失败) } // 判断是否是不可恢复异常这类不重试直接丢弃 func isMsgCheckErr(err error) bool { // json解析失败、参数非法这类不可恢复不重试 return false } func consume(queueName string, ch *amqp.Channel) { msgs, err : ch.Consume( queueName, , false, // autoAckfalse手动ack false, false, false, nil, ) if err ! nil { log.Fatal(err) } forever : make(chan bool) go func() { for d : range msgs { var err error waitTime : initialInterval success : false // 本地重试逻辑 for attempt : 0; attempt maxAttempts; attempt { err handleMsg(d) if err nil { success true break } // 不可恢复异常直接跳出不再重试 if isMsgCheckErr(err) { break } log.Printf(消费失败第%d次重试等待%v, err:%v, attempt1, waitTime, err) time.Sleep(waitTime) waitTime time.Duration(float64(waitTime) * multiplier) } if success { // 消费成功ACK err d.Ack(false) if err ! nil { log.Println(ack失败, err) } } else { // 本地重试全部耗尽Nack(false)不重新入队 // 注意第二个参数 requeuefalse千万不要写true err d.Nack(false, false) if err ! nil { log.Println(nack失败, err) } log.Println(本地重试全部失败消息丢弃/转入死信) } } }() log.Printf( [*] Waiting for messages. To exit press CTRLC) -forever } func main() { conn, err : amqp.Dial(amqp://guest:guestlocalhost:5672/) if err ! nil { log.Fatal(err) } defer conn.Close() ch, err : conn.Channel() if err ! nil { log.Fatal(err) } defer ch.Close() queueName : test_queue consume(queueName, ch) }三、MessageRecoverer 失败消息处理策略Spring AMQP 在重试耗尽后使用MessageRecoverer接口处理失败消息共 3 种实现表格Spring AMQP 类含义Go 代码等价操作RejectAndDontRequeueRecoverer默认重试耗尽拒绝消息不重新入队队列配置死信就进 DLQ没配置直接丢弃d.Nack(false, false)第二个参数requeuefalseImmediateRequeueMessageRecoverer重试耗尽消息放回原队列重新消费会无限循环生产严禁使用d.Nack(false, true)第二个参数requeuetrueRepublishMessageRecoverer重试耗尽重新发布一条新消息到指定交换机异常交换机 error.direct消息头附加异常信息手动调用ch.Publish()body 异常信息放入 header 投递到目标交换机区分 原生死信 DLQMQ Broker 自动转发原消息 RepublishMessageRecoverer代码手动发布一条全新消息原消息 ack 移除不属于原生 DLQ常用来做异常告警投递到 error.queue 人工排查。Go 实现 RepublishMessageRecoverer复刻 Springpackage main import ( github.com/streadway/amqp log ) // 模仿Spring RepublishMessageRecoverer type RepublishMessageRecoverer struct { ch *amqp.Channel targetExchange string routingKey string } func NewRepublishMessageRecoverer(ch *amqp.Channel, exchange, rk string) *RepublishMessageRecoverer { return RepublishMessageRecoverer{ ch: ch, targetExchange: exchange, routingKey: rk, } } // Recover 重试耗尽后调用 func (r *RepublishMessageRecoverer) Recover(delivery *amqp.Delivery, err error) error { // 复制header追加异常信息 headers : delivery.Headers if headers nil { headers make(amqp.Table) } headers[x-exception] err.Error() headers[x-original-exchange] delivery.Exchange headers[x-original-rk] delivery.RoutingKey // 手动发布一条新消息到 error.direct和Spring逻辑一致 errPub : r.ch.Publish( r.targetExchange, r.routingKey, false, false, amqp.Publishing{ Body: delivery.Body, Headers: headers, ContentType: delivery.ContentType, }, ) if errPub ! nil { log.Printf(republish失败:%v, errPub) return errPub } // 原消息ack从原队列移除 errAck : delivery.Ack(false) if errAck ! nil { log.Printf(ack原消息失败:%v, errAck) } return nil }四、核心总结面试背诵版消费者如何保证消息一定被消费开启消费者确认机制为 auto由 Spring 确认消息处理成功后返回 ack异常时返回 nack。开启消费者失败重试机制并设置MessageRecoverer多次重试失败后将消息投递到异常交换机交由人工处理。

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询