我第一次在生产环境碰到 RabbitMQ,是因为一个支付回调把数据库打挂了。
事情是这样的:用户在第三方支付平台付完钱,支付平台会回调我们的接口通知「钱到账了」。问题是每天中午有个促销时段,几千笔回调几乎同时到达。MySQL 连接池直接爆了,订单状态改不过来,用户那边显示「已付款」但我们这边还是「待支付」。客服电话被打爆。
后来我们加了一层 RabbitMQ,回调接口只做一件事:把消息扔进队列,立刻返回 200。后台的 Worker 从队列里慢慢消费,一个一个改订单状态。数据库稳了,客服安静了。
这大概是消息队列最经典的场景:削峰填谷。流量突然涌进来的时候,不让它直接冲击下游,而是先存起来,让下游按自己的节奏处理。
AMQP 和 RabbitMQ 的模型
RabbitMQ 实现的是 AMQP 0-9-1 协议。理解它的模型,核心就四个东西:
Producer(生产者) 发消息。消息不是直接发到队列的,而是先发给 Exchange(交换机)。Exchange 根据 Binding(绑定) 规则,把消息路由到一个或多个 Queue(队列)。Consumer(消费者) 从队列里取消息。
Producer → Exchange → [Binding规则] → Queue → Consumer
这个设计比「生产者直接往队列里写」灵活得多。同一个消息可以被路由到多个队列,不同类型的 Exchange 有不同的路由逻辑。
四种 Exchange,四种路由策略
RabbitMQ 内置了四种 Exchange 类型。你大概率只会用到前三种。
Direct Exchange
消息的 routing key 和队列的 binding key 精确匹配时才路由。
假设有一个日志系统。你创建两个队列:error_queue 绑定了 error,info_queue 绑定了 info。发一条 routing key 是 error 的消息,它只去 error_queue。简单直接。
Fanout Exchange
忽略 routing key,把消息广播到所有绑定的队列。最像「大喇叭」。
之前那个支付回调的场景,除了改订单状态,你还需要发短信通知用户、更新营销积分。三个操作对应三个队列,全部绑定到同一个 fanout exchange。支付回调消息进来,三个队列各拿一份,各干各的。
Topic Exchange
routing key 和 binding key 做模式匹配。binding key 用 . 分段,* 匹配一段,# 匹配零段或多段。
比如 binding key 是 order.payment.*,匹配 order.payment.success 和 order.payment.failed,但不会匹配 order.payment 或 order.shipping.success。
真实的微服务场景里 topic exchange 用得最多。一个 order.# 绑订单服务,一个 order.payment.# 绑支付服务,路由粒度足够精细。
Headers Exchange
不看 routing key,看消息头里的键值对。几乎没人用这玩意儿。我也见过一两个老系统用了 headers exchange 做条件路由,后来都迁移到 topic 了。维护成本远大于收益。
消息丢了怎么办:确认和持久化
消息队列最大的恐惧不是性能,是丢消息。
RabbitMQ 的确认机制分两层。Producer 层叫 Publisher Confirm,Consumer 层叫 Consumer Ack。
Publisher Confirm
Producer 发消息给 RabbitMQ 后,RabbitMQ 确认收到。如果 RabbitMQ 没确认(网络抖动、磁盘满了、队列不存在),Producer 重发。
代码长这样:
channel.confirm_delivery()
try:
channel.basic_publish(exchange='orders', routing_key='new', body=msg)
except pika.exceptions.UnroutableError:
# 消息没发出去,重试或者记日志
logger.error("Message unroutable, will retry")
注意一个细节:confirm_delivery 模式下,RabbitMQ 确认的是「消息到达了 Exchange 并且被路由到了至少一个队列」。如果消息路由不到任何队列,会触发 UnroutableError。所以你要么用 mandatory 参数 + return listener 处理不可路由消息,要么提前确保 binding 存在。
Consumer Ack
Consumer 拿到消息,处理完,手动告诉 RabbitMQ「我搞定了,你可以删了」。如果 Consumer 处理到一半挂了,没发 ack,RabbitMQ 会把消息重新投递给其他 Consumer。
def callback(ch, method, properties, body):
try:
process_order(body) # 处理订单
ch.basic_ack(delivery_tag=method.delivery_tag)
except Exception:
# 不 ack,让 RabbitMQ 重新投递
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True)
一个常见的坑:忘了 basic_ack。消息会一直留在队列里,Consumer 那边看上去处理成功了,但 RabbitMQ 不知道。消息越积越多,最后内存爆了。
持久化
确认只管「消息有没有被正确处理」。消息本身会不会因为 RabbitMQ 重启而消失,那是持久化的事。
三个地方需要持久化:Exchange 声明时 durable=True,Queue 声明时 durable=True,消息发送时 delivery_mode=2。
这里有个细节我踩过坑:发消息的时候忘了设 delivery_mode=2,消息确实进队列了,Consumer 也正常消费了,一切正常。结果服务器重启了一次——重启前积压的几万条消息全部消失。Exchange 和 Queue 还在,消息没了。因为消息本身没有被持久化到磁盘。
什么时候不该用 RabbitMQ
RabbitMQ 不是银弹。有几个场景它真的不合适:
日志收集和流式处理。 RabbitMQ 的吞吐量在几万条/秒这个量级。如果你要处理百万级日志,直接用 Kafka。Kafka 的磁盘顺序读写模型天然适合高吞吐场景,RabbitMQ 的索引式存储做不到。
超大消息。 RabbitMQ 对单个消息的大小没有硬限制,但默认最大 128MB 可以改。问题是消息存在内存里,大消息会让内存暴涨,触发 flow control,整个集群变慢。遇到过有人往 RabbitMQ 里塞图片 Base64,一首歌的功夫整个集群跪了。这种事应该把文件存 OSS,消息里只放 URL。
强顺序要求。 多个 Consumer 消费同一个队列时,消息处理的顺序是不确定的。如果需要严格按顺序处理,要么单 Consumer,要么用 Kafka partition 的机制。RabbitMQ 本身不保证全局顺序。
最后
RabbitMQ 属于那种「平时想不起来,出了问题才意识到它多关键」的基础设施。它不性感,没有 AI 光环,不会出现在技术大会的 Keynote 里。但当你凌晨三点被报警叫起来发现订单状态又没更新的时候,你会感谢那个往架构里加了消息队列的人。
说回技术选型。如果你是中小团队,消息量在万级/秒以下,需要一个可靠、好运维、文档丰富的消息中间件——RabbitMQ 几乎是最稳妥的选择。Protocol 标准、管理界面好用、社区庞大、踩坑经验足够多。这些比「最高吞吐量」之类的 benchmarks 重要得多。