简介:讲解 RabbitMQ 两种实现延迟消息方案:DLX 死信交换机、延迟消息插件;针对订单超时业务痛点,引入阶梯式延迟消息优化方案,附完整 Go 代码。
一、什么是延迟消息
延迟消息:生产者发送消息时指定延时时间,消费者不会立刻收到消息,等待设定时间到达后才会收到消息。
典型业务场景:下单后超时未支付自动取消订单、释放库存。 业务流程:下单 → 保存订单 → 扣减库存 → 设置倒计时,超时自动取消订单、恢复库存。
业务痛点
- 延迟时间设置太短:消息频繁触发,服务器压力大
- 延迟时间设置太长:业务时效性差
延迟任务:在指定一段时间之后才执行的任务。
RabbitMQ 实现延迟消息有两种主流方案:
- 死信交换机 DLX(原生支持,无需插件)
- 延迟消息插件
rabbitmq-delayed-message-exchange(社区插件)
二、方案 1:死信交换机 DLX 模拟延迟消息
1. 什么是死信
队列内消息满足下面任意条件,就会变成死信 (Dead Letter):
- 消费者调用
basic.reject/basic.nack消费失败,并且requeue=false(不重新放回队列) - 消息超时 TTL,超过时间无人消费
- 队列消息堆满,最早入队的消息被丢弃
默认情况下,死信会直接丢弃。 如果队列配置了
dead-letter-exchange属性,死信不会丢弃,会投递到这个指定交换机,这个交换机就叫做死信交换机 DLX。
2. DLX 实现延迟消息原理
- 业务交换机
simple.direct→ 业务队列simple.queue,业务队列绑定死信交换机dlx.direct - 死信交换机
dlx.direct→ 死信队列dlx.queue,消费者监听死信队列 - 生产者发送消息到业务队列,给消息设置 TTL(如 30s)
- 消息在业务队列等待 30s 超时 → 变成死信,自动转发到死信交换机,路由到死信队列
- 消费者监听死信队列,收到消息执行业务逻辑。
本质:利用消息过期转死信的特性,间接模拟延迟消息
DLX 方案缺点
如果队列中消息 TTL 各不相同,前面消息未过期,后面消息即使到时间也不会被处理(队列 FIFO)。适合所有消息延迟时长相同场景。
三、方案 2:RabbitMQ 延迟消息插件
原理
插件提供一种特殊类型交换机x-delayed-message。消息投递到此交换机后,先在交换机内部暂存,等待设定延迟时间到期,再投递到目标队列。
使用前提:安装插件,重启 RabbitMQ。
核心代码片段(Go)
// 声明延迟交换机,类型为x-delayed-message err := ch.ExchangeDeclare( "delay.exchange", "x-delayed-message", // 插件专属交换机类型,不是原生direct/topic true, false, false, false, // x-delayed-type:指定底层路由模式 direct/topic amqp.Table{ "x-delayed-type": "direct", }, ) // 普通队列,不需要配置死信参数 q, _ := ch.QueueDeclare("delay.queue",false,false,false,false,nil) ch.QueueBind("delay.queue", "delay.key", "delay.exchange", false, nil) // 发送消息:Header中设置x-delay,单位毫秒 headers := amqp.Table{} headers["x-delay"] = 30000 // 延迟30s err = ch.Publish( "delay.exchange", "delay.key", false, false, amqp.Publishing{ ContentType: "text/plain", Body: []byte("超时取消订单消息"), Headers: headers, // x-delay写在消息header }, )优缺点小结 优点:可以给每条消息自定义延迟时间,代码简单直观,无 DLX 的 FIFO 阻塞问题 缺点:社区第三方插件,非官方内置;大量延迟消息占用 MQ 内存;MQ 重启未到期消息会丢失。 适用场景:延迟时间较短的业务。
四、订单超时取消业务:原始方案的问题
需求:下单 30 分钟未支付,自动取消订单,恢复库存。 原始方案:下单直接发送一条 30 分钟延迟消息。 存在两大问题:
- 高并发场景下,大量消息堆积在 MQ,MQ 压力巨大。
- 绝大部分订单下单 1 分钟内就完成支付,但是这条消息依旧要在 MQ 等待满 30 分钟,浪费 MQ 存储资源。
优化方案:阶梯式延迟消息(分段检查)
核心思想:把长延迟拆分成多次短延迟,指数退避重试检查
- 创建订单,发送第一轮短延迟消息(例如 1 分钟)
- 消费者收到消息,查询数据库订单支付状态
- 订单已支付:直接结束流程,不再产生新消息
- 订单未支付:判断是否达到最大超时 30 分钟
- 未到 30 分钟:计算下一次延迟时间(1min → 2min →4min →8min,指数递增),重发延迟消息
- 已达到 30 分钟:执行取消订单、恢复库存,结束流程
优势:已支付订单直接终止,不会持续生成消息;只有长时间未支付的订单才持续产生少量延迟消息,极大降低 MQ 压力。技术选型:延迟消息插件,支持动态设置每次消息的延迟时长。
五、完整 Go 代码实现(阶梯延迟消息)
依赖:
go get github.com/rabbitmq/amqp091-go前置条件:RabbitMQ 安装rabbitmq-delayed-message-exchange延迟插件
package main import ( "context" "encoding/json" "fmt" "log" "time" "github.com/rabbitmq/amqp091-go" ) // ===================== 常量定义 ===================== const ( // RabbitMQ连接地址 MQURL = "amqp://guest:guest@127.0.0.1:5672/" // 延迟交换机名称 DelayExchange = "order.delay.exchange" // 延迟队列名称 DelayQueue = "order.delay.queue" // 路由key,用于交换机投递消息到队列 RoutingKey = "order.delay.key" // 订单最大超时时间:30分钟,超过该时间未支付则取消订单 MaxTimeout = 30 * time.Minute ) // DelayMsg 延迟消息载体结构体,每条消息携带订单信息 type DelayMsg struct { OrderId string `json:"order_id"` // 订单编号 CreateAt time.Time `json:"create_at"` // 订单创建时间,用于判断总超时 NextDelayMs int32 `json:"next_delay_ms"` // 本次消息延迟毫秒数 } // 模拟订单数据库,生产环境替换为MySQL/Redis var orderDB = map[string]bool{} // key:订单号 value:是否已支付 func main() { // 1. 建立MQ连接 conn, err := amqp091.Dial(MQURL) if err != nil { log.Fatalf("mq connect fail: %v", err) } defer conn.Close() ch, err := conn.Channel() if err != nil { log.Fatalf("open channel fail: %v", err) } defer ch.Close() // 声明延迟交换机 x-delayed-message err = ch.ExchangeDeclare( DelayExchange, "x-delayed-message", true, false, false, false, amqp091.Table{ "x-delayed-type": "direct", //底层路由类型direct }, ) if err != nil { log.Fatalf("declare exchange err: %v", err) } // 声明队列 q, err := ch.QueueDeclare( DelayQueue, true, false, false, false, nil, ) if err != nil { log.Fatalf("declare queue err: %v", err) } // 队列绑定交换机 err = ch.QueueBind(q.Name, RoutingKey, DelayExchange, false, nil) if err != nil { log.Fatalf("queue bind err: %v", err) } // ========== 生产者:模拟下单,发送第一轮延迟消息 ========== orderId := "ORDER_10001" orderDB[orderId] = false //新建订单标记未支付 firstDelay := int32(60 * 1000) //第一轮延迟1分钟 sendDelayMsg(ch, orderId, time.Now(), firstDelay) fmt.Printf("创建订单 %s,发送第一轮延迟消息,延迟%ds\n", orderId, firstDelay/1000) // ========== 消费者:监听延迟队列,处理订单状态检查 ========== msgs, err := ch.Consume(q.Name, "", false, false, false, false, nil) if err != nil { log.Fatalf("consume err: %v", err) } forever := make(chan struct{}) go func() { for d := range msgs { var msg DelayMsg err := json.Unmarshal(d.Body, &msg) if err != nil { log.Printf("json unmarshal fail: %v", err) d.Ack(false) continue } now := time.Now() orderCreateTime := msg.CreateAt orderId := msg.OrderId fmt.Printf("\n收到延迟消息,订单:%s\n", orderId) // 1. 查询订单支付状态 isPaid := orderDB[orderId] if isPaid { fmt.Printf("订单 %s 已支付,无需处理,结束流程\n", orderId) d.Ack(false) continue } // 2. 判断是否超过最大30分钟超时 if now.Sub(orderCreateTime) >= MaxTimeout { fmt.Printf("订单 %s 超时30分钟,执行取消订单,恢复库存!\n", orderId) // 业务逻辑:update订单状态、恢复库存 delete(orderDB, orderId) d.Ack(false) continue } // 3. 未支付 && 未超时:阶梯递增延迟时间 nextDelayMs := msg.NextDelayMs * 2 // 边界保护:下一次延迟不能超过剩余到期时间 leftTime := MaxTimeout - now.Sub(orderCreateTime) if int64(nextDelayMs) > leftTime.Milliseconds() { nextDelayMs = int32(leftTime.Milliseconds()) } fmt.Printf("订单未支付,未超时,继续投递下一轮延迟消息,下次延迟:%d秒\n", nextDelayMs/1000) // 重发延迟消息 sendDelayMsg(ch, orderId, orderCreateTime, nextDelayMs) d.Ack(false) } }() fmt.Println("消费者启动成功,等待消息...") <-forever } // sendDelayMsg 发送延迟消息函数 func sendDelayMsg(ch *amqp091.Channel, orderId string, createAt time.Time, delayMs int32) { msg := DelayMsg{ OrderId: orderId, CreateAt: createAt, NextDelayMs: delayMs, } body, _ := json.Marshal(msg) ctx := context.Background() headers := amqp091.Table{} headers["x-delay"] = delayMs // 插件核心:延迟毫秒数 err := ch.PublishWithContext(ctx, DelayExchange, RoutingKey, false, false, amqp091.Publishing{ ContentType: "application/json", Headers: headers, Body: body, }, ) if err != nil { log.Printf("publish delay msg err: %v", err) } }六、总结 & 面试要点
- RabbitMQ 延迟消息两种实现:DLX 死信交换机(原生)、延迟消息插件(第三方)
- DLX 利用消息 TTL 过期转死信模拟延迟;缺点是 FIFO 队列阻塞,不同 TTL 消息会互相等待。
- 延迟插件使用
x-delayed-message交换机,消息 header 设置x-delay,适合每条消息延时不同场景。 - 订单超时业务优化:不要一次性发送 30 分钟长延迟消息,采用阶梯式分段延迟,减少 MQ 消息堆积,节约资源。
- 阶梯延迟逻辑:每次消费检查订单状态,已支付直接终止;未支付则指数退避重发延迟消息,直到最大超时。