首页 / 资讯中心 / 文章详情

【RabbitMQ #13】 | 延迟消息

【RabbitMQ #13】 | 延迟消息 ★ FEATURED ARTICLE
简介讲解 RabbitMQ 两种实现延迟消息方案DLX 死信交换机、延迟消息插件针对订单超时业务痛点引入阶梯式延迟消息优化方案附完整 Go 代码。一、什么是延迟消息延迟消息生产者发送消息时指定延时时间消费者不会立刻收到消息等待设定时间到达后才会收到消息。典型业务场景下单后超时未支付自动取消订单、释放库存。 业务流程下单 → 保存订单 → 扣减库存 → 设置倒计时超时自动取消订单、恢复库存。业务痛点延迟时间设置太短消息频繁触发服务器压力大延迟时间设置太长业务时效性差延迟任务在指定一段时间之后才执行的任务。RabbitMQ 实现延迟消息有两种主流方案死信交换机 DLX原生支持无需插件延迟消息插件rabbitmq-delayed-message-exchange社区插件二、方案 1死信交换机 DLX 模拟延迟消息1. 什么是死信队列内消息满足下面任意条件就会变成死信 (Dead Letter)消费者调用basic.reject/basic.nack消费失败并且requeuefalse不重新放回队列消息超时 TTL超过时间无人消费队列消息堆满最早入队的消息被丢弃默认情况下死信会直接丢弃。 如果队列配置了dead-letter-exchange属性死信不会丢弃会投递到这个指定交换机这个交换机就叫做死信交换机 DLX。2. DLX 实现延迟消息原理业务交换机simple.direct→ 业务队列simple.queue业务队列绑定死信交换机dlx.direct死信交换机dlx.direct→ 死信队列dlx.queue消费者监听死信队列生产者发送消息到业务队列给消息设置 TTL如 30s消息在业务队列等待 30s 超时 → 变成死信自动转发到死信交换机路由到死信队列消费者监听死信队列收到消息执行业务逻辑。本质利用消息过期转死信的特性间接模拟延迟消息DLX 方案缺点如果队列中消息 TTL 各不相同前面消息未过期后面消息即使到时间也不会被处理队列 FIFO。适合所有消息延迟时长相同场景。三、方案 2RabbitMQ 延迟消息插件原理插件提供一种特殊类型交换机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 分钟延迟消息。 存在两大问题高并发场景下大量消息堆积在 MQMQ 压力巨大。绝大部分订单下单 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:guest127.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 消息堆积节约资源。阶梯延迟逻辑每次消费检查订单状态已支付直接终止未支付则指数退避重发延迟消息直到最大超时。
阅读完成 · 觉得有帮助?
咨询建站