1.3 生产者发布消息第一课 本节摘要:生产者的职责远不止调用一次发送函数:建立连接、打开通道、声明资源、构造消息、发布、确认送达,六步缺一不可。本节是追踪单上第一处真正的"丢失风险高发区"——basicpublish 是一个火后不理的调用,本节先把完整发布流程跑通,风险点的堵法留给第 3 章。 前两节解决了"该不该发",这一节解决"怎么发"。别小看这六步,生产环节每一步都有它存在的理由,跳过任何一步的理解,到了排查事故时都要加倍偿还。 发布一条消息要过几道门 第一道门是连接(Connection)。生产者与 Broker 之间是一条 TCP 长连接,建立它需要三次握手加 AMQP 握手,开销不小,所以连接应当长驻复用,而不是每发一条消息建一次。第二道门是通道(Channel)。
本节摘要:生产者的职责远不止调用一次发送函数:建立连接、打开通道、声明资源、构造消息、发布、确认送达,六步缺一不可。本节是追踪单上第一处真正的"丢失风险高发区"——basic_publish 是一个火后不理的调用,本节先把完整发布流程跑通,风险点的堵法留给第 3 章。
前两节解决了"该不该发",这一节解决"怎么发"。别小看这六步,生产环节每一步都有它存在的理由,跳过任何一步的理解,到了排查事故时都要加倍偿还。
第一道门是连接(Connection)。生产者与 Broker 之间是一条 TCP 长连接,建立它需要三次握手加 AMQP 握手,开销不小,所以连接应当长驻复用,而不是每发一条消息建一次。第二道门是通道(Channel)。连接内部可以开多条轻量的虚拟连接,所有 AMQP 动作——声明、发布、订阅——都发生在通道上。为什么需要这道门?因为 TCP 连接太贵而操作太频繁,多路复用是必然设计。连接与通道的细节第 2 章有专节,这里先建立"一连接多通道"的心智模型。
第三道门是声明。交换机、队列、绑定这些资源必须先存在,消息才有去处。声明操作是幂等的:资源已存在且参数一致时无事发生,参数不一致时通道报错——这个特性决定了"发布方与消费方都声明一遍"是安全且推荐的做法,谁先启动谁创建,谁都不会因为启动顺序而失败。
第四道门才是发布。basic_publish 把消息交给交换机,带上路由键。第五道门是确认——请注意,这扇门默认是关着的:basic_publish 在协议层是异步 fire-and-forget,调用返回不代表 Broker 收到了。发往一个不存在的交换机会直接抛异常,但路由到"没有匹配队列的交换机"则静默丢弃。追踪单风险点二号就埋在这里,发布确认机制的完整解法见第 3 章,本节先用简单手段暴露它。
背景:为后续章节准备一个可复用的演示环境:交换机 trace.direct,队列 trace.orders,绑定键 order.created。
操作:下面是完整的生产者代码,每一步都带注释说明其必要性:
import pika params = pika.ConnectionParameters( host="localhost", heartbeat=30) # 心跳保活,防半开连接 connection = pika.BlockingConnection(params) channel = connection.channel() # 声明交换机:类型 direct,持久化(durable=True 使其重启后仍在) channel.exchange_declare( exchange="trace.direct", exchange_type="direct", durable=True) # 声明队列并绑定:消费端通常会做同样声明,顺序无关 channel.queue_declare(queue="trace.orders", durable=True) channel.queue_bind( exchange="trace.direct", queue="trace.orders", routing_key="order.created") body = '{"order_id": "A1024", "event": "created"}' channel.basic_publish( exchange="trace.direct", routing_key="order.created", body=body.encode(), properties=pika.BasicProperties( content_type="application/json", delivery_mode=2)) # 消息持久化标记 print("发布动作完成") connection.close()
执行后用管理命令清点库存,验证消息确实入队:
rabbitmqctl list_queues name messages messages_ready messages_unacknowledged # 预期输出: # Listing queues ... # trace.orders 1 1 0
三列数字各有所指:messages 是队列中消息总数,messages_ready 是待投递的,messages_unacknowledged 是已投给消费者但未签收的。现在三个数字分别是 1、1、0——消息已就位,无人签收。
结果:消息顺利入队。解读:把实验推进一步,故意把 routing_key 改成 order.cancelled 再发一条——队列里依然只有 1 条消息。那条消息去哪了?没有任何报错、任何日志提示。它被 direct 交换机评估后找不到匹配绑定,直接丢弃。这就是"静默丢弃"的现场还原:发布方的代码看起来完全正常,消息却人间蒸发。三个堵漏手段——发送方确认、备用交换机、mandatory 标志——分别在 2.2 节和 3.5 节展开,此处先把这个坑刻进记忆。
变式一:批量发布。循环一千次调用 basic_publish 时,性能瓶颈往往不在网络而在每条消息的确认开销,batch 发布与异步确认是第 5 章性能优化的主角。变式二:把 host 换成一个不可达地址,观察 BlockingConnection 抛出的连接超时异常——生产代码必须捕获它并走重试逻辑,否则上游服务会被连累。
把本节内容浓缩成发布前过一遍的检查项,这份清单会贯穿全册:
| 检查项 | 不做会怎样 |
|---|---|
| 连接与通道长驻复用 | 每请求建连,延迟翻倍且可能耗尽连接数 |
| 声明幂等带参数一致 | 启动顺序敏感,或参数漂移引发通道错误 |
| 消息体带 content_type | 消费方只能猜格式,升级时全面崩坏 |
| 持久化标记 delivery_mode=2 | Broker 重启消息蒸发(配合 3.2 节) |
| 送达确认机制 | 静默丢弃无从察觉(第 3 章补全) |
💡 关键直觉:basic_publish 的返回值不能作为"发送成功"的依据。把它理解为"把信塞进了邮筒",而确认机制才是"邮局开的回执"。
消息已出生产者之手。下一章镜头切到 Broker 内部:交换机如何接过这张追踪单,决定它去往哪个队列。