直接用@retry装饰器会失败,因为其复用已断开的channel导致ChannelClosedByBroker等错误;可靠重试必须每次重建连接、重声明exchange/queue、重获取channel,并启用mandatory、confirm_delivery和合理heartbeat。
为什么直接用 retry 装饰器会失败
在 rabbitmq 生产端(publish 端)加重试,不能直接套用通用的 @retry(比如 tenacity 或 backoff),因为大多数消息发送失败是瞬时网络抖动或连接断开,而装饰器默认重试的是同一段函数调用——如果底层 pika.blockingconnection 已断开,后续重试仍会复用已失效的 channel,抛出 channelclosedbybroker 或 connectionclosed 错误,反而掩盖真实问题。
可靠重试必须和连接生命周期对齐:每次重试都应尝试重建连接 + 重声明 exchange/queue(若需要)+ 重获取 channel。
- 不要在装饰器里缓存
connection或channel实例 - 避免在重试循环中反复
channel.basic_publish而不检查channel.is_open - 不要假设 exchange 已存在;生产端重试时,exchange 可能被运维临时删掉
publish_with_retry 必须封装连接重建逻辑
核心是把“建立可用 channel”抽成可重试的子过程,再在其之上做 publish 重试。推荐用 tenacity,因为它支持嵌套重试策略和自定义 stop/wait 条件:
from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type import pika <p>def get_ready_channel(): connection = pika.BlockingConnection(pika.ConnectionParameters(host='localhost')) channel = connection.channel() channel.exchange_declare(exchange='my_exchange', exchange_type='direct', durable=True) return connection, channel</p><p>@retry( stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=1, max=10), retry=retry_if_exception_type((pika.exceptions.AMQPConnectionError, pika.exceptions.ChannelClosedByBroker)) ) def publish_with_retry(routing_key, body): connection, channel = get_ready_channel() try: channel.basic_publish( exchange='my_exchange', routing_key=routing_key, body=body, mandatory=True, # 触发 ReturnListener 若路由失败 properties=pika.BasicProperties(delivery_mode=2) # 持久化 ) finally: channel.close() connection.close()
-
mandatory=True很关键:若消息无法路由到 queue,会触发ReturnListener,此时应记录并告警,而不是静默重试 - 每次重试都调用全新
get_ready_channel(),确保连接和 channel 都是 fresh 的 - 不用
add_callback_threadsafe或异步回调——BlockingConnection 不支持
如何处理 Unroutable 和 Undeliverable 消息
RabbitMQ 的 basic.publish 默认不反馈路由结果。要捕获“发出去但没进任何 queue”的情况,必须启用 mandatory + 注册 return_listener,但这和重试装饰器有冲突:装饰器只捕获异常,不捕获正常返回下的业务失败。
解决方案是把 publish 拆成两步:先注册 return 回调,再发消息,最后同步等待结果(通过 condition 或 flag):
立即学习“Python免费学习笔记(深入)”;
def publish_with_return_check(routing_key, body):
connection, channel = get_ready_channel()
result_flag = {'returned': False, 'exception': None}
<pre class="brush:php;toolbar:false;">def on_return(channel, method, properties, body):
result_flag['returned'] = True
channel.add_on_return_callback(on_return)
try:
channel.basic_publish(
exchange='my_exchange',
routing_key=routing_key,
body=body,
mandatory=True,
properties=pika.BasicProperties(delivery_mode=2)
)
if result_flag['returned']:
raise ValueError('Message was returned (unroutable)')
finally:
channel.close()
connection.close()
- 这个函数本身不适合直接套
@retry,因为on_return是异步回调,result_flag的读取时机不可靠 - 更稳妥的做法是:重试装饰器只覆盖连接层错误;对
Unroutable这类业务错误,应在上层捕获后走降级逻辑(如写本地日志、发告警),而非盲目重试 - 如果必须重试 unroutable 场景,需确认是 routing_key 临时错配(比如下游服务未启动),这时重试前应 sleep 并检查依赖服务健康状态
生产环境必须设置的三个参数
很多重试失效,是因为没关掉自动确认或忽略心跳超时。以下三项不设好,重试可能永远卡死或伪造成功:
-
ConnectionParameters(heartbeat=30):避免中间设备(如 ELB、NAT)断连;值太小会导致频繁重连,太大则故障发现慢 -
channel.confirm_delivery():开启 publisher confirms,让basic_publish在消息真正入队后才返回(否则只是写入 socket 缓冲区) -
connection.parameters.blocked_connection_timeout = 30:防止 broker 主动阻塞连接时,客户端无限等待
confirm mode 下,basic_publish 会阻塞直到 broker 返回 ack/nack,这正是重试需要的真实成功信号。没有它,装饰器看到的“成功”可能只是消息还在 client 内存里。


















