📚 消息队列系列 · 第 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 广播:一条消息,全员通知
场景
订单系统:下单成功啦! ↓ 广播邮件服务:📧 发通知邮件 ✅短信服务:💬 发验证短信 ✅积分服务:⭐ 加积分 ✅
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 的区别:
四、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 :按地区分发给物流"""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,看懂分区和消费者组就够了
关注我,软件测试实战干货持续更新 🚀