☰
rabbitmqdemo:从Docker部署到死信队列的RabbitMQ实践指南
2026/10/12 2:39:21 网站建设 项目流程

简介:RabbitMQ与MFC集成示例工程,面向需要在Windows C++桌面应用中接入消息队列的开发者。由于网上MFC环境结合RabbitMQ的现成资料较少,作者整理出这套支持消息发送与接收的可运行Demo,适合学习AMQP协议、异步通信与界面集成的中初级开发者。压缩包共63个文件,约47.48MB,包含C++头文件与源文件、Visual Studio工程文件、资源脚本、编译中间文件、可执行程序、动态链接库、静态库及调试符号等,既有完整源码也有可直接运行的exe。目前已有326人浏览学习,对于RabbitMQ与MFC这一少见组合,是值得参考的基础示例。资源内RabbitmqClient封装了连接与消息接口,rabbitmqDemoDlg提供对话框操作界面,并附带了rabbitmq.4.dll等依赖库,可帮助读者快速跑通发送接收流程,理解队列声明、绑定与收发逻辑;同时展示了MFC界面与消息处理线程的配合方式以及常见错误处理与日志记录思路,借助该示例还可以理解RabbitMQ的消息队列模型、AMQP协议流程、队列声明与绑定等关键环节。

1. rabbitmqdemo 是什么:一个能照单全收的消息中间件骨架

收到 rabbitmqdemo 这个项目需求时,我最初以为它只是把官方教程抄一遍的入门示例。真把代码从头过了一遍才发现,这个 demo 覆盖了连接生命周期、交换机与队列声明、消息持久化、消费者手动确认、失败重试与死信兜底这条完整链路——它不是教你怎么发一条消息,而是把「消息中间件上生产前必须做的那几个决策」一次性摆到桌面上。如果你正打算在一个没有队列的系统里引入异步处理,或者已经在用 RabbitMQ 但总在连接断开、消息丢失、重复消费之间反复翻车,这个 demo 值得你照着跑一遍再按需改造。它回答的核心问题只有一个:在 RabbitMQ 里,怎么做到发出去的消息不丢、收下来的消息能兜底。

2. 把 RabbitMQ 本地跑起来:Docker 命令、两个端口与最小拓扑

2.1 本地部署:Docker 一条命令,5672 与 15672 两个端口要分清

先把 broker 跑起来再做别的。我一般不会在本机直接装 RabbitMQ,因为不同系统版本不一致会带来一堆和代码无关的环境坑。常见做法是直接拉官方镜像,一条命令把 broker 和管理控制台同时启动:

docker run -d \ --name rabbitmq-demo \ -p 5672:5672 \ -p 15672:15672 \ -e RABBITMQ_DEFAULT_USER=demo \ -e RABBITMQ_DEFAULT_PASS=demo123456 \ rabbitmq:3-management

参数说明:-p 5672:5672把宿主机 5672 端口映射到容器内,这是 AMQP 协议端口,生产者和消费者客户端都走它;-p 15672:15672是管理界面和 HTTP API 的端口,浏览器打开http://localhost:15672就能看到控制台。rabbitmq:3-management是官方维护的带管理插件的镜像,省掉手动启用插件的步骤,本地 demo 最省事。

特别注意默认的 guest 账号:RabbitMQ 出于安全限制,guest 只允许从 loopback 地址连接。Docker 端口映射环境下,连接经过网桥到达容器,broker 从容器内视角看到的是网关 IP 而不是 127.0.0.1,经常出现「客户端连不上、管理台却能登录」的诡异现象。所以我在启动参数里直接通过RABBITMQ_DEFAULT_USER和RABBITMQ_DEFAULT_PASS创建了一个 demo 用户,客户端和服务端都用这个账号,绕开 guest 的 loopback 限制。这是跑 demo 阶段最不折腾的姿势。

装好之后验证一下容器状态:

docker ps | grep rabbitmq-demo # 期望看到 STATUS 为 Up,端口映射显示 0.0.0.0:5672->5672/tcp

如果端口被占用,docker run会直接报 Address already in use。换端口时记得生产者消费者代码里的 port 要和这里保持一致,管理台端口也同步换,这是一个容易忽略的联动点。

2.2 最小拓扑设计:一个 direct 交换机、两条队列、一组绑定关系

环境起来后,下一步不是写代码,而是先把拓扑画清楚。RabbitMQ 的核心模型是「生产者 → 交换机 → 队列 → 消费者」,消息不会直接进队列,而是带着 routing key 发给交换机,由交换机按绑定关系投递。所以拓扑设计的本质是回答三个问题:用什么交换机、建几条队列、每条队列绑定什么 key。

demo 里我推荐用 direct 类型交换机。它按 routing key 做精确匹配,语义最直白,排查问题最方便。相比 fanout(广播给所有绑定队列)和 topic(按点分通配符匹配),direct 在演示阶段能让你一眼看出「这条消息为什么进了这个队列、没进那个队列」。很多人拿到 demo 需求第一件事就是写 producer 发消息,结果发着发着发现消息不见了,或者进了错误队列。拓扑先行的习惯能省掉一半以上的排查时间。

组件名称说明
交换机demo.exchange类型 direct,设置 durable
队列一demo.queue.order订单业务,绑定 keyorder.create
队列二demo.queue.log操作日志,绑定 keyorder.log
绑定key 精确匹配同一交换机挂两条队列,各收各的消息

为什么是两条队列而不是一条?单单跑通一条队列之后,你会自然地想验证「路由」这件事:生产者发一条order.create消息,只有订单队列收到,日志队列纹丝不动,这就理解了 direct 的精髓。而订单和日志分开消费也是生产环境的真实姿态——不同业务关注不同的数据,消费速率、失败处理和监控目标都不一样。

2.3 用管理台验证拓扑是否生效

启动后打开管理台,用 demo/demo123456 登录,按以下顺序做三轮检查:

  1. Overview 页面:确认节点状态是 running,没有 memory alarm 或 disk alarm。出现 alarm 时 broker 会阻塞所有写操作,这在生产环境是最常见的问题之一,demo 阶段先养成看一眼的习惯。
  2. Queues 页面:脚本声明过的队列会出现在列表里,点进队列能看到 Bindings 标签页,确认绑定关系不是空的。
  3. Exchanges 页面:确认 demo.exchange 存在,点进去能看到它绑定了哪些队列、每个绑定的 routing key 是什么。

管理台还有一个实用功能:在 Queues 页面点进某个队列,可以直接手动 Publish 一条测试消息,也可以手动 Get 取出队列里的消息。这个功能在排查「生产端代码没写好」还是「消费端代码有问题」时特别好用——绕过代码,用控制台手动投递,立刻能分清责任。demo 过程中我至少会用一次它来验证链路。

3. 生产者端动手:连接复用、声明拓扑与消息发布

3.1 RabbitMQ 连接与信道的关系:为什么每条消息一条连接是错的

在写发布代码之前,先要把「连接」和「信道」这两个概念分清楚。Connection 是一条底层 TCP 连接,建立时要经过 AMQP 握手、认证(SASL),可能还有 TLS 握手,代价高;Channel 是建立在 Connection 之上的轻量复用通道,本质是 TCP 连接上的多路复用。RabbitMQ 官方推荐的做法是:一个进程维护一个 Connection,进程内的每个线程、每路业务各开一个 Channel。Channel 的创建销毁非常廉价,Connection 不是。

初学者最容易犯的错是每条消息 new 一个 Connection,发完立刻 close。短连接场景下,一次完整的连接握手开销可能比业务本身还大,而且高频建连会被 broker 判定为异常甚至触发连接数限制。之前有 A 同学写的 demo 就是这种写法,压测时每秒只能发两百条,改成单连接多 Channel 后直接上到每秒几千条。这个参数不体现在任何一行业务逻辑里,但对吞吐的影响是数量级的。

我把连接参数单独做成一个函数,方便在不同脚本里复用:

import pika def build_connection(): params = pika.ConnectionParameters( host='localhost', port=5672, credentials=pika.PlainCredentials('demo', 'demo123456'), heartbeat=60, blocked_connection_timeout=30 ) return pika.BlockingConnection(params)

参数说明:heartbeat=60 表示空闲时每 60 秒发送一次心跳保活包,防止连接被服务端或中间网络设备回收;blocked_connection_timeout 是 broker 因内存或磁盘告警导致连接被阻塞时,客户端等待的最长时间,避免一直挂死。注意ConnectionParameters里没有「连接池」这个参数,池化是业务层自己的事——进程内持有一个长连接就够了,多线程场景给每个工作线程分配一个 Channel。

3.2 声明交换机、队列和绑定:一份能反复执行的幂等代码

拓扑声明在 pika 里是显式调用。声明本身是幂等的,重复执行不会报错,只要参数一致就会复用已存在的对象。先写完整代码再逐行拆参数:

import json import time import pika connection = build_connection() channel = connection.channel() # 交换机:持久化 direct 交换机 channel.exchange_declare( exchange='demo.exchange', exchange_type='direct', durable=True ) # 队列:持久化,非排他,不自动删除 channel.queue_declare( queue='demo.queue.order', durable=True, exclusive=False, auto_delete=False ) # 队列绑定到交换机,routing key 精确匹配 channel.queue_bind( queue='demo.queue.order', exchange='demo.exchange', routing_key='order.create' ) # 构造一条订单消息 payload = { 'order_id': f'ORD-{int(time.time())}', 'amount': 99.50, 'user_id': 'u_10001' } # 发布消息 channel.basic_publish( exchange='demo.exchange', routing_key='order.create', body=json.dumps(payload, ensure_ascii=False).encode('utf-8'), properties=pika.BasicProperties( delivery_mode=2, content_type='application/json' ) ) print('producer: 发布成功') connection.close()

关键参数逐个说:

  • exchange_type='direct':交换机类型。demo 用精确匹配,生产环境用 topic 的情况更多,但原理一致。
  • durable=True:交换机和队列都持久化。broker 重启后,durable 的交换机和队列定义会恢复;非持久化的会被清除。注意 durable=True 只保证队列和交换机定义不丢,消息本身能否幸存取决于 BasicProperties 的 delivery_mode。
  • exclusive=False:队列允许其他连接访问。exclusive=True 的队列只能被声明它的连接使用,连接关闭队列即删除,适合临时任务队列,不适合业务队列。
  • auto_delete=False:最后一个消费者断开后队列不自动删除。auto_delete=True 的队列在消费者全部断开后会消失,demo 阶段容易误伤。
  • delivery_mode=2:消息标记为持久化。broker 会把这类消息落盘,重启后恢复。不设这个值,即使队列是 durable 的,消息依然只活在内存里。

还有一个细节:basic_publish默认不检查交换机是否存在,如果 exchange 名字拼错,消息不会报错而是被静默丢弃。这是 RabbitMQ 最坑的默认行为之一:发布失败不抛异常,你看到的只是队列里永远没有消息。排查手段是看管理台 Exchanges 页面的 publish 计数,或者按下文的 publisher confirm。

3.3 用 publisher confirm 确认消息真的被 broker 收下

默认发布是发完即忘:basic_publish 返回只代表本地客户端缓存成功,不代表 broker 收到了消息。对订单这类不能丢的消息,需要开启 publisher confirm 机制:

# 打开发布确认模式 channel.confirm_delivery() try: channel.basic_publish( exchange='demo.exchange', routing_key='order.create', body=json.dumps(payload, ensure_ascii=False).encode('utf-8'), properties=pika.BasicProperties(delivery_mode=2) ) print('producer: broker 已确认接收') except pika.exceptions.NackError: print('producer: broker 返回 nack,需要重发') except pika.exceptions.UnroutableError: print('producer: 交换机或路由不存在,消息未被接受')

confirm_delivery() 打开确认模式后,每次 basic_publish 会同步等待 broker 的 Basic.Ack 才能继续。它的价值在于:只要发布成功返回,消息就保证进了 broker 的内存或磁盘;失败时会抛异常拿到明确信号,便于重试或告警。同步确认的代价是每条消息多一个网络往返 RTT。生产环境要吞吐的话,常见做法是改为异步批量确认:发送多批消息后再统一处理确认事件,或者用事件回调监听 Basic.Ack。demo 阶段不用做这么复杂,但你要知道这个取舍存在。

4. 消费者端动手:手动 ACK、prefetch 与幂等消费实践

4.1 自动 ACK 的陷阱:消息一投递就被标记为成功

消费者端第一个要决策的参数是 auto_ack。pika 的 basic_consume 里默认 auto_ack=True,含义是 broker 把消息投递给消费者后立刻从队列里删除,无论消费者是否真的处理成功。

自动 ACK 的问题在于:消息从队列发出到你处理完这中间,如果消费者进程崩溃、网络断开或回调抛出异常,这条消息已经不在 broker 里了,彻底丢失。没有重试机会,没有恢复入口。对日志类消息也许可以接受,对订单、支付这类消息等于自毁。所以 demo 里直接演示手动 ACK 的标准姿势:

  • 回调里把消息处理完,调用basic_ack(delivery_tag=method.delivery_tag)通知 broker 这条可以删除。
  • 处理失败时调用basic_nack(delivery_tag=method.delivery_tag, requeue=False),把消息标记为失败,按后续死信配置转移。
  • 如果希望失败消息重新投递,用basic_nack(requeue=True),broker 会把它放回队列重新派发。但注意:重投递意味着可能反复失败形成循环,生产环境要有重试次数上限,demo 里建议直接用 requeue=False 配合死信。

看到服务端返回 channel 级异常、消息又确实没被 ack 时,一个常见误判是把错误归因于「broker 不消费」,其实只是忘了在回调末尾写 ack。消费者回调跑完而不 ack,消息会一直压在 unacked 状态,不超时不移交,看起来像队列卡死,实际上卡的是自己那行代码。

4.2 预取数量 prefetch:控制消费者的「在途消息」窗口

prefetch_count 是消费者端另一个决定命运的旋钮。它表示消费者在收到 ack 确认之前,broker 最多一次性给该消费者投递多少条未确认消息。直白说,就是你手上最多同时「在途」几条消息。

prefetch 设得太大,比如默认 0 表示无限制,broker 会把消息哗哗地全部推给消费者,消费者内存被积压占满,一旦崩溃所有未确认消息全部回到队列重新投递——好消息是消息不丢,坏消息是恢复时要重新处理一遍,重复消费风险上升。prefetch 设得太小,比如永远是 1,则每处理一条消息都要等一轮网络往返,吞吐上不去,单消费者尤其明显。

prefetch_count常见场景典型风险
1多个慢消费者轮询分摊、顺序敏感吞吐上限明显偏低
10 ~ 50CPU 型轻处理、同步回调内存占用和吞吐较均衡
100 ~ 300IO 密集型、回调快速转线程池未确认堆积多,崩溃恢复成本高

一个粗略的估算方式:消费者能接受的最大在途数量 = 单条消息预估内存占用 × prefetch_count × 消费者数,建议别超过消费者进程可用内存的三分之一。demo 里订单消息很小,prefetch 设在 10 到 50 之间不会出错。另外要提醒:prefetch 只在手动 ACK 模式下生效,如果用 auto_ack=True,broker 根本不管你的 prefetch 设置,这个参数等于白设。

4.3 消费者完整代码:回调、ACK、幂等检查一次写完

直接给一份可跑的消费者实现:

import json import time import pika def already_processed(order_id: str) -> bool: """模拟幂等检查:用 Redis SETNX 或数据库唯一索引查询""" return False # demo 阶段默认未处理过 def save_order(order: dict) -> None: """模拟落库""" print(f'save order: {order["order_id"]}') def on_order_message(channel, method, properties, body): order_id = None try: msg = json.loads(body.decode('utf-8')) order_id = msg.get('order_id') print(f'consumer: 收到 {order_id}') # 幂等检查:重复投递直接 ack,不重复处理 if already_processed(order_id): print(f'consumer: {order_id} 重复消息,跳过') channel.basic_ack(delivery_tag=method.delivery_tag) return # 模拟业务耗时 time.sleep(0.2) # 落库 save_order(msg) # 处理成功才 ack channel.basic_ack(delivery_tag=method.delivery_tag) except Exception as exc: print(f'consumer: 处理异常 {exc}') # requeue=False 让消息进死信队列,避免无限重投 channel.basic_nack( delivery_tag=method.delivery_tag, requeue=False ) params = pika.ConnectionParameters( host='localhost', port=5672, credentials=pika.PlainCredentials('demo', 'demo123456'), heartbeat=60 ) connection = pika.BlockingConnection(params) channel = connection.channel() # 预取窗口设为 10 channel.basic_qos(prefetch_count=10) channel.basic_consume( queue='demo.queue.order', on_message_callback=on_order_message, auto_ack=False ) print('consumer: 开始监听,Ctrl+C 退出') try: channel.start_consuming() except KeyboardInterrupt: connection.close()

这段代码有几个地方值得停下来解释:

  • already_processed是幂等兜底。RabbitMQ 的投递语义是 at-least-once,重复投递不是异常而是承诺的一部分——ACK 丢失、连接断开、broker 重投都会导致同一条消息被消费多次。接收方必须幂等,这是消息中间件架构里绕不开的约束。进数据库前先查、或者给 order_id 建唯一索引,重复消息直接 ACK 跳过即可。
  • 异常分支里用了requeue=False,意思是这条消息我不认领了,broker 按死信配置把消息转移到死信队列。如果这里写成requeue=True,失败消息会立刻重新投递给消费者,大概率再次失败,反复循环。demo 阶段看不出问题,生产环境里这就是典型的「毒消息死循环」现场。
  • start_consuming()是阻塞调用,主线程会一直挂着。需要和其他逻辑共存时,把它放进独立线程,或者改用channel.consume()生成器分批拉取。

5. rabbitmqdemo 避坑指南:五个最容易翻车的现场

demo 能跑通只是第一步,下面这五个坑基本覆盖了我在不同 MQ 项目里反复遇到、也亲手踩过的现场,每条按现象、原因、解决来拆。

5.1 翻车点一:broker 重启之后,消息和队列一起消失

现象:昨天发的消息还查得到,早晨一看队列空了,管理台里队列列表也少了。RabbitMQ 日志没有异常,纯粹是重启后一切归零。

原因:队列声明时用的默认非持久化,消息发布时 delivery_mode 也没设 2。RabbitMQ 正常运行时把这些消息放在内存里,内存够用就表现正常,一旦进程重启,内存消息全部丢失,非持久化队列连定义都不会恢复。很多人只记着把队列设为 durable=True,忘了消息本身也要持久化标记,这两件事是独立维度:durable 管队列定义,delivery_mode=2 管消息存储。

解决:三处都要到位——交换机 durable=True、队列 durable=True、消息 properties 的 delivery_mode=2。要验证是否真的落盘,可以在管理台 Queues 页面看某个队列的消息总数旁边是否有持久化计数,重启 broker 后再看队列和消息是否恢复。我用这种方式做过「断电恢复」演练,效果比任何代码测试都直观。

5.2 翻车点二:消费者日志突然刷 connection reset

现象:消费者进程跑了两三个小时,日志开始刷 ConnectionResetError 或 ReadTimeout,程序没崩但消息不再被处理;重启消费者又能撑一阵。生产环境里经常是凌晨出这个错,白天没事。

原因:最常见的元凶是 heartbeat 超时。消费者的回调函数里如果做了同步阻塞调用——比如调外部 HTTP 接口没设超时、查数据库被慢查询卡住——pika 主线程在等回调返回时无法发送心跳包,broker 超过连续两个心跳间隔没收到心跳,就从服务端断开连接。另一个低频原因是容器或宿主机回收空闲 TCP 连接,部署在虚拟化环境时要一并排查。

解决:给 ConnectionParameters 显式传入 heartbeat=60,别用默认值。回调函数里不要做重耗时操作,把外部调用丢到线程池、设置明确的 HTTP/DB 超时时间,保证回调能在几百毫秒内返回。再加一层断线重连兜底:捕获连接关闭异常后重建 connection 和 channel,重新执行 basic_consume。连接断不可怕,可怕的是断了自己不知道,或知道了不重建。

5.3 翻车点三:消息被重复消费,账单多算了一次

现象:订单消息被消费了两次,订单表里出现了两条相同 order_id 的记录,或者下游收到两条短信。查日志发现第一次已经显示处理成功,过了几秒又来了一条一模一样的。

原因:这就是 at-least-once 语义。消费者处理完消息、ACK 还没发出去时连接断了,broker 认为投递失败,重新投递同一条消息。消息本身没有丢,但接收方看到了两次。这不是配置错误,而是 RabbitMQ 保证「不丢」的代价——投递可能会重复。

解决:接收端必须幂等。最简单可靠的做法是用 order_id 建唯一索引,插入时数据库会拒绝重复;也可以先用 Redis 做SETNX consumed:{order_id} 1去重,命中说明已处理。关键在顺序:先查重、再处理、成功后 ACK。如果先处理再查重,重复消息依然可能执行两次副作用操作。这个顺序问题就是重复消费事故里最隐蔽的细节,不少人栽在把查重写在了业务逻辑之后。

5.4 翻车点四:队列堆积疯涨,消费者看着却没闲着

现象:管理台 Queues 页面 Ready 消息数持续上涨,但消费者的日志每隔几秒就有一条处理记录,看起来一直在干活,总量就是下不去。

原因:典型的是 prefetch 设成 1 加上消费者只有一两个;或者回调里同步等待外部依赖(数据库连接池被打满、下游接口响应变慢),处理速度被外部系统拖住。还有一种隐蔽原因:消费者连接建立成功,但 basic_consume 指定的队列名和生产者写的不是同一个,消息全部堆在另一个队列,而那个队列根本没有消费者。

解决:先看管理台,确认目标队列的 Consumers 列不是 0,确认消息确实在被消费。然后把 prefetch 从 1 提到 10~50,或者给消费者增加并发。排查外部依赖时,用监控看回调平均耗时——如果耗时本身是秒级,加多少消费者都没用,瓶颈在下游系统而不是 MQ。

5.5 翻车点五:多个消费者实例同时开,消息顺序乱了

现象:订单有「创建→支付→完成」三个状态消息,发的是同一个 order_id,逻辑上要求按顺序处理。开了两个消费者实例后,完成消息先被处理、创建消息后处理,业务状态机直接错乱。

原因:RabbitMQ 只保证单个队列内部消息的投递顺序,而且只在「该队列只有一个消费者」时才严格成立。多个消费者并行消费同一队列时,队列投递按序,但每个消费者处理速度和并发不同,谁先处理完没有保证。顺序和并发放一起,天然冲突。

解决:对严格有序的消息,收敛并发度。一个队列只挂一个消费者,或者用 routing key 把同一个业务 ID 的消息哈希到同一个队列,保证同一个 order_id 的消息只进一个消费者;接收侧也可以给消息带 sequence 序号,消费后合流重排,但成本较高,不建议在 demo 阶段引入。最务实的做法是把「顺序需求」和「高吞吐需求」分开建模,别让一条队列同时扛这两件事。

6. 从 demo 到能落地:死信队列与延迟消息这两个必配项

6.1 死信队列:给失败消息一条退路

到这里 demo 已经能跑通生产和消费,但离能上线还差两个配置。第一个是死信队列。消息在三种情况下会变成死信:消费者basic_nack(requeue=False)或basic_reject明确拒绝、消息在队列里超过 TTL 过期、队列长度达到上限触发溢出删除。默认情况下这些消息直接被丢弃,生产环境绝对不行。

配置方法:给业务队列声明时挂两个参数,指定死信交换机和死信路由键:

channel.exchange_declare(exchange='demo.dlx', exchange_type='direct', durable=True) channel.queue_declare(queue='demo.queue.dead', durable=True) channel.queue_bind(queue='demo.queue.dead', exchange='demo.dlx', routing_key='order.dead') channel.queue_declare( queue='demo.queue.order', durable=True, arguments={ 'x-dead-letter-exchange': 'demo.dlx', 'x-dead-letter-routing-key': 'order.dead' } )

这样消费者basic_nack(requeue=False)的消息会自动进入 demo.queue.dead 死信队列,保留完整消息体,方便事后捞出来人工排查或重新投递。我在生产上的习惯是给死信队列配独立消费者,失败消息统一抄送告警,每周统计一次死信量——量突然变大,说明有消费者正在处理高错误率的业务。

6.2 延迟消息:用 TTL + 死信实现 30 秒延期

第二个必配项是延迟消息。RabbitMQ 没有内置延迟队列,标准做法是利用「TTL 过期消息转死信」这个机制:消息进一个设置了 TTL 且没有消费者的保留队列,超时后变成死信转到目标交换机,再路由到真正的业务队列。

# 声明一个 30 秒后过期的保留队列 channel.queue_declare( queue='demo.queue.delay_30s', durable=True, arguments={ 'x-message-ttl': 30000, # 单位毫秒 'x-dead-letter-exchange': 'demo.exchange', 'x-dead-letter-routing-key': 'order.create' } )

生产者把消息发到 demo.queue.delay_30s,30 秒后消息过期,broker 把它作为死信转投到 demo.exchange,按order.create路由进订单队列,消费者感知到的就是一条延迟了 30 秒的业务消息。支付超时关闭、秒杀订单限时锁定、定时补偿任务都能用这个模式做,这也是死信和 TTL 组合最经典的用法。

6.3 验证方法

验证延迟链路的方法很朴素:往延迟队列发一条消息,打开管理台盯 demo.queue.delay_30s 的消息数,前 30 秒它稳定有 1 条,超时后瞬间归零,同时在 demo.queue.order 里看到这条消息出现。管理台上这个时间差就是最好的验收证据。TTL 的单位是毫秒,别顺手填成 30,那会得到一条 30 毫秒就过期的消息,排查半天还以为是链路问题。

我这几年搭 RabbitMQ demo 有一个固定习惯:不管业务多简单,先把死信链路和 TTL 延迟链路搭好,再写任何业务代码。因为线上真正让你半夜起来的,从来不是「消息发不出来」,而是「消息去哪了」和「消息为什么晚到了这么久」——死信队列和 TTL 就是这两类问题最初的双保险。希望帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询