当前位置:首页>python>Python 玩转 RabbitMQ:路由、主题一次讲清

Python 玩转 RabbitMQ:路由、主题一次讲清

  • 2026-10-11 07:04:45
Python 玩转 RabbitMQ:路由、主题一次讲清

📚 消息队列系列 · 第 3 篇(共 6 篇)📌 上一篇:RabbitMQ 从入门到不丢消息📌 本篇主角:交换器(Exchange)—— RabbitMQ 最核心的路由机制📌 你将学到:fanout 广播、direct 路由、topic 主题、实战应用📌 难度指数:⭐⭐(需要 Python 基础)


上一篇我们学了工作队列:一个队列,多个消费者,任务轮流分发。

但真实业务里,需求往往更复杂:

场景1:订单产生了 → 邮件服务、短信服务、积分服务【都要】收到场景2:订单产生了 → 只有支付服务需要处理【支付成功】的消息场景3:订单产生了 → 物流服务只想收【华东地区】的订单

如果只用上篇的简单队列,你得写死很多逻辑。而 RabbitMQ 的交换器(Exchange)机制,就是专门解决这类"消息往哪送"的问题。


一、先搞懂:生产者到底把消息发给谁?

你可能有个疑问——上篇代码里,生产者发消息时明明写的是:

channel.basic_publish(    exchange='',          # ← 这个参数是什么?    routing_key='task_queue',    body=message)

生产者从来不直接把消息发给队列,而是发给"交换器"(Exchange)。

生产者 → 交换器(Exchange) → 队列(Queue) → 消费者               ↑         消息的"快递分拣中心"

交换器就是快递分拣中心,它按照"路由规则"决定:

  • • 这条消息送给哪个队列?
  • • 送一份?还是复制多份?
  • • 还是直接丢弃?

RabbitMQ 有 4 种交换器类型,本篇讲最常用的 3 种:

类型
中文名
路由规则
一句话
fanout
广播
发给所有绑定的队列
群发消息
direct
直连
按路由键精确匹配
精准投递
topic
主题
按路由键模糊匹配
灵活订阅
headers
头匹配
按消息头匹配
很少用,略过

二、fanout 广播:一条消息,全员通知

场景

订单系统:下单成功啦!   ↓ 广播邮件服务:📧 发通知邮件    ✅短信服务:💬 发验证短信    ✅积分服务:⭐ 加积分        ✅

fanout 广播:只要绑定了这个交换器的队列,全部都能收到消息。

生产者代码

# broadcast_producer.py —— 广播模式生产者import pikaimport jsonconnection = pika.BlockingConnection(    pika.ConnectionParameters(host='localhost'))channel = connection.channel()# 1. 声明 fanout 交换器channel.exchange_declare(    exchange='order_events',   # 交换器名字    exchange_type='fanout'     # 广播类型)# 2. 发消息(fanout 模式下 routing_key 会被忽略)order = {    "order_id": "A1001",    "user": "张三",    "amount": 299.00,    "status": "paid"}channel.basic_publish(    exchange='order_events',    routing_key='',            # fanout 不关心路由键    body=json.dumps(order, ensure_ascii=False))print(f"📤 订单消息已广播: {order['order_id']}")connection.close()

消费者代码(邮件服务)

# email_worker.py —— 邮件服务消费者import pikaimport jsonconnection = pika.BlockingConnection(    pika.ConnectionParameters(host='localhost'))channel = connection.channel()# 1. 声明同一个交换器channel.exchange_declare(    exchange='order_events',    exchange_type='fanout')# 2. 创建临时队列(名字随机,断开后自动删除)#    每个消费者有自己的临时队列,避免互相抢消息result = channel.queue_declare(queue='', exclusive=True)queue_name = result.method.queueprint(f"📬 邮件服务队列: {queue_name}")# 3. 绑定队列到交换器channel.queue_bind(exchange='order_events', queue=queue_name)def callback(ch, method, properties, body):    order = json.loads(body)    print(f"📧 邮件服务收到订单 {order['order_id']},给 {order['user']} 发确认邮件")channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True)print("👂 邮件服务等待订单事件...")channel.start_consuming()

把消费者代码复制成两份(email_worker.py、sms_worker.py),都跑起来,再运行生产者:

# 终端 1python email_worker.py# 终端 2python sms_worker.py# 终端 3python broadcast_producer.py

输出:

终端1(邮件): 📬 邮件服务队列: amq.gen-xxxx1终端2(短信): 📬 短信服务队列: amq.gen-xxxx2终端3: 📤 订单消息已广播: A1001终端1: 📧 邮件服务收到订单 A1001,给 张三 发确认邮件终端2: 💬 短信服务收到订单 A1001,给 张三 发通知短信

💡 关键点:

  • • 每个消费者用 exclusive=True 的临时队列,谁也不会抢谁的消息
  • • fanout 就像广播电台——不管你想不想听,发了就是全频道覆盖

三、direct 直连:按路由键精准投递

场景

订单消息产生   ↓ 路由键 = "payment.success"日志服务:只收 routing_key="payment.success" 的 ✅风控服务:只收 routing_key="payment.failed" 的 ✅

direct 直连:路由键(routing_key)完全匹配的队列才能收到消息。

        生产者的路由键            ↓   ┌────────────────┐   │  direct交换器    │   └────────────────┘      ↓ 精确匹配        ↓ 精确匹配 队列A(绑定 payment.success)    队列B(绑定 payment.failed)

生产者代码

# direct_producer.py —— 直连模式生产者import pikaimport jsonconnection = pika.BlockingConnection(    pika.ConnectionParameters(host='localhost'))channel = connection.channel()# 声明 direct 交换器channel.exchange_declare(exchange='payment_events', exchange_type='direct')# 发两条不同级别的消息events = [    ("payment.success", {"order_id": "B2001", "status": "成功"}),    ("payment.failed",  {"order_id": "B2002", "status": "失败"}),]for routing_key, event in events:    channel.basic_publish(        exchange='payment_events',        routing_key=routing_key,          # ← 路由键决定投递目标        body=json.dumps(event, ensure_ascii=False)    )    print(f"📤 发送 [{routing_key}]: {event}")connection.close()

消费者代码(日志服务:只收成功消息)

# log_worker.py —— 只处理支付成功import pikaimport jsonconnection = pika.BlockingConnection(    pika.ConnectionParameters(host='localhost'))channel = connection.channel()channel.exchange_declare(exchange='payment_events', exchange_type='direct')queue_name = channel.queue_declare(queue='', exclusive=True).method.queue# 只绑定 routing_key = 'payment.success'channel.queue_bind(    exchange='payment_events',    queue=queue_name,    routing_key='payment.success')def callback(ch, method, properties, body):    event = json.loads(body)    print(f"📝 日志服务记录: 订单 {event['order_id']} 支付{event['status']}")channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True)print("👂 日志服务只关注 [payment.success]...")channel.start_consuming()

测试:

# 终端 1python log_worker.py# 终端 2python direct_producer.py

输出:

终端1: 👂 日志服务只关注 [payment.success]...终端2: 📤 发送 [payment.success]: {'order_id': 'B2001', ...}终端2: 📤 发送 [payment.failed]:  {'order_id': 'B2002', ...}终端1: 📝 日志服务记录: 订单 B2001 支付成功   ← 只收到这一条!       (B2002 的失败消息没人要,被丢弃)

💡 direct 与 fanout 的区别:

  • • fanout:不看路由键,全部广播
  • • direct:只看路由键,精确匹配

四、topic 主题:路由键模糊匹配(重点⭐)

场景

订单消息的路由键 = "order.east.beijing"   ↓ topic 模糊匹配物流服务:绑定 "order.*.beijing"   → 只收北京的 ✅统计服务:绑定 "order.east.*"      → 只收华东的 ✅所有服务:绑定 "order.#"          → 所有订单都收 ✅

topic 主题:支持通配符匹配的路由键,比 direct 灵活得多。

两个通配符

通配符
含义
示例
*
匹配一个词
order.*
 → order.east ✅ order.east.beijing ❌
#
匹配零个或多个词
order.#
 → order ✅ order.east.beijing ✅

路由键由点号分隔的词组成:

order.east.beijing │     │     │词1   词2   词3order.*.beijing   → 匹配 order.east.beijing / order.north.beijingorder.east.*      → 匹配 order.east.beijing / order.east.shanghaiorder.#           → 匹配所有 order 开头的

生产者代码

# topic_producer.py —— 主题模式生产者import pikaimport jsonconnection = pika.BlockingConnection(    pika.ConnectionParameters(host='localhost'))channel = connection.channel()channel.exchange_declare(exchange='log_topic', exchange_type='topic')# 模拟不同地区的订单日志logs = [    ("order.east.beijing",  {"order_id": "C001", "region": "北京"}),    ("order.east.shanghai", {"order_id": "C002", "region": "上海"}),    ("order.north.harbin",  {"order_id": "C003", "region": "哈尔滨"}),]for routing_key, log in logs:    channel.basic_publish(        exchange='log_topic',        routing_key=routing_key,        body=json.dumps(log, ensure_ascii=False)    )    print(f"📤 发送 [{routing_key}]: {log['order_id']}")connection.close()

消费者代码(华东地区监控)

# east_worker.py —— 只关注华东订单import pikaimport jsonconnection = pika.BlockingConnection(    pika.ConnectionParameters(host='localhost'))channel = connection.channel()channel.exchange_declare(exchange='log_topic', exchange_type='topic')queue_name = channel.queue_declare(queue='', exclusive=True).method.queue# 绑定模式:order.east.*  → 华东所有地区channel.queue_bind(    exchange='log_topic',    queue=queue_name,    routing_key='order.east.*')def callback(ch, method, properties, body):    log = json.loads(body)    print(f"📊 华东监控: 订单 {log['order_id']} 来自 {log['region']}")channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True)print("👂 华东监控绑定 [order.east.*]...")channel.start_consuming()

另一个消费者绑定 order.*.beijing(只收北京):

# beijing_worker.py —— 只关注北京订单# 和 east_worker.py 几乎一样,只改一行:channel.queue_bind(    exchange='log_topic',    queue=queue_name,    routing_key='order.*.beijing'   # ← 只匹配北京的)

测试:

# 终端 1python east_worker.py      # 绑定 order.east.*# 终端 2python beijing_worker.py   # 绑定 order.*.beijing# 终端 3python topic_producer.py

输出:

终端1(华东): 📊 华东监控: 订单 C001 来自 北京终端1(华东): 📊 华东监控: 订单 C002 来自 上海             (C003 哈尔滨不是华东,不收)终端2(北京): 📊 北京监控: 订单 C001 来自 北京             (C002 上海、C003 哈尔滨都不是北京,不收)

💡 同样的消息,不同的绑定,各取所需——这就是 topic 的威力。


五、三种交换器对比

fanout 广播
direct 直连
topic 主题
匹配规则
不看路由键
完全相等
通配符匹配
灵活度
⭐
⭐⭐
⭐⭐⭐
场景
群发通知
精准投递
按业务维度订阅
典型用例
系统广播、全员通知
日志分级、支付状态
地区/类型筛选
代码难度
简单
简单
中等

一句话选择口诀:

全都要      → fanout只要某一种  → direct按规则挑    → topic

六、实战:订单系统的消息分发

把三种模式用在一个完整场景里:

"""订单系统消息分发实战- fanout  :下单广播给所有服务- direct  :支付结果给日志/风控- topic   :按地区分发给物流"""import pikaimport jsonconnection = pika.BlockingConnection(    pika.ConnectionParameters(host='localhost'))channel = connection.channel()# ===== 声明 3 个交换器 =====channel.exchange_declare(exchange='order_broadcast', exchange_type='fanout')channel.exchange_declare(exchange='payment_result', exchange_type='direct')channel.exchange_declare(exchange='delivery_region', exchange_type='topic')# ===== 声明并绑定队列 =====def bind_queue(exchange, routing_key=''):    """创建临时队列并绑定"""    q = channel.queue_declare(queue='', exclusive=True).method.queue    channel.queue_bind(exchange=exchange, queue=q, routing_key=routing_key)    return q# 订阅订单广播的服务q_email = bind_queue('order_broadcast')      # 邮件服务q_sms = bind_queue('order_broadcast')        # 短信服务# 订阅支付结果的服务q_log = bind_queue('payment_result', 'payment.success')   # 日志服务q_risk = bind_queue('payment_result', 'payment.failed')   # 风控服务# 订阅地区订单的服务q_east = bind_queue('delivery_region', 'order.east.*')    # 华东物流q_beijing = bind_queue('delivery_region', 'order.*.beijing')  # 北京物流# ===== 模拟业务流程 =====def on_order_created(order):    """用户下单:广播给所有服务"""    channel.basic_publish(        exchange='order_broadcast',        routing_key='',        body=json.dumps(order, ensure_ascii=False)    )    print(f"📤 广播下单事件: {order['order_id']}")def on_payment(order, success=True):    """支付完成:按结果精准分发"""    routing_key = 'payment.success' if success else 'payment.failed'    channel.basic_publish(        exchange='payment_result',        routing_key=routing_key,        body=json.dumps({**order, "paid": success}, ensure_ascii=False)    )    print(f"📤 支付结果[{routing_key}]: {order['order_id']}")def on_delivery(order):    """进入物流:按地区分发"""    routing_key = f"order.{order['region']}.{order['city']}"    channel.basic_publish(        exchange='delivery_region',        routing_key=routing_key,        body=json.dumps(order, ensure_ascii=False)    )    print(f"📤 物流事件[{routing_key}]: {order['order_id']}")# ===== 执行:模拟一单完整流程 =====order = {"order_id": "D10086", "user": "李四", "region": "east", "city": "beijing", "amount": 99}on_order_created(order)   # 1. 下单 → 所有服务收到on_payment(order, True)   # 2. 支付成功 → 日志服务收到on_delivery(order)        # 3. 物流 → 华东+北京物流都收到print("\n🎯 各队列收到的消息分布:")print(f"  邮件服务队列: {q_email}   → 收到广播 ✅")print(f"  短信服务队列: {q_sms}     → 收到广播 ✅")print(f"  日志服务队列: {q_log}     → 收到支付成功 ✅")print(f"  风控服务队列: {q_risk}    → 无支付失败,空闲")print(f"  华东物流队列: {q_east}    → 收到华东订单 ✅")print(f"  北京物流队列: {q_beijing} → 收到北京订单 ✅")connection.close()

七、常见问题

❓ 路由键的命名规范

推荐用 点分结构,方便 topic 匹配:user.created          (不要 user_created)order.east.beijing    (不要 orderEastBeijing)payment.success

❓ 消息发出去没人收会怎样?

direct/topic:没有队列匹配 → 消息被丢弃(RabbitMQ 默认行为)

想避免?用 mandatory=True 参数,路由不到时回调通知你。

❓ 一个队列可以绑定多个交换器吗?

# 可以!一个队列绑定多个交换器,任何一边来消息都能收到channel.queue_bind(exchange='order_broadcast', queue=q)channel.queue_bind(exchange='payment_result',  queue=q, routing_key='payment.success')

❓ fanout 和 topic 怎么选?

完全没规则、全员通知     → fanout有明确分类、要筛选       → topic精确一对一              → direct

总结

生产者 → 交换器 → 队列 → 消费者          │          ├─ fanout  广播给所有队列(群发)          ├─ direct  路由键精确匹配(精准)          └─ topic   通配符模糊匹配(灵活)
知识点
一句话
交换器
消息的分拣中心,生产者真正发送的目标
fanout
广播模式,绑定就收
direct
直连模式,路由键完全匹配
topic
主题模式,* 匹配一个词,# 匹配多个词
临时队列exclusive=True
,各消费者独立收消息

生产环境里三种交换器经常混用,就像订单系统的实战演示——广播通知用 fanout,支付结果用 direct,地区分发用 topic。


📝 下篇预告:入门 Kafka,看懂分区和消费者组就够了

关注我,软件测试实战干货持续更新 🚀

最新文章

随机文章