Python 基础体系 · 第 88/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。

Python 消息系统:RabbitMQ、Kafka、消费语义、重试和死信

消息系统解决的不是“把函数放到后台执行”这么单一的问题,而是把生产者和消费者之间的时间、故障和吞吐解耦:

生产者 ──发布──> Broker ──投递──> 消费者 ──处理──> 外部系统
   │                 │             │
   │                 │             └── ACK / offset commit
   │                 └──持久化、路由、分区、重试、保留
   └──发布确认

这里至少存在三次需要单独确认的动作:

  1. 生产者是否把消息交给 Broker;
  2. Broker 是否把消息交给消费者;
  3. 消费者是否完成了业务处理。

这三次确认彼此不能互相替代。RabbitMQ 的 publisher confirm 只确认 Broker 接受了发布,consumer acknowledgement 只确认消费者对投递的处理;二者是正交机制。(rabbitmq.com)


一、先建立消息系统的基本模型

1. 消息、生产者、消费者和 Broker

消息是带有数据和元数据的一次传输单元。数据通常是 JSON、Protobuf、Avro 或自定义二进制;元数据可能包含:

  • message_id:消息唯一标识;
  • event_type:事件类型;
  • occurred_at:事件发生时间;
  • trace_id:链路追踪标识;
  • key:分区或路由依据;
  • attempt:重试次数;
  • schema_version:数据结构版本。

生产者负责创建和发布消息。生产者通常不应该认为“写入本地 TCP 缓冲区”就代表消息已经可靠保存,因为连接可能在 Broker 接收前断开。

消费者负责获取消息并执行处理逻辑。消费者处理消息时可能调用数据库、HTTP 服务、文件系统、支付接口或另一个消息系统。

Broker负责暂存、路由、排序、复制、投递和记录消费进度。RabbitMQ 和 Kafka 都是 Broker,但两者的核心数据模型不同:

  • RabbitMQ 更接近“消息投递系统”:消息进入队列,通常被某一个消费者取走并在确认后删除;
  • Kafka 更接近“持久化日志”:消息按分区追加保存,消费者通过 offset 表示读取位置,不同消费组可以独立读取同一份数据。

这一区别决定了 ACK、重试、顺序和扩展方式。


二、RabbitMQ 的消息模型:Exchange、Queue、Binding

RabbitMQ 基于 AMQP 0-9-1 时,生产者通常不是直接向 Queue 发布消息,而是先向 Exchange 发布:

Producer
   │ publish(exchange, routing_key, body)
   ▼
Exchange
   │ 根据类型和 Binding 路由
   ├──> Queue A ──> Consumer A
   └──> Queue B ──> Consumer B

2.1 Exchange

Exchange 是路由器,不是消息最终消费位置。常见类型包括:

  • direct:精确匹配 routing key;
  • topic:使用通配符匹配,例如 order.*
  • fanout:广播到所有绑定队列;
  • headers:根据消息头匹配。

Queue 是保存消息并向消费者投递的缓冲区。多个消费者可以共同消费同一个 Queue,此时每条消息通常只会投递给其中一个消费者。

Binding 是 Exchange 到 Queue 的路由关系。生产者只负责提供 Exchange 和 routing key,是否能进入某个 Queue 取决于 Binding。

2.2 RabbitMQ 示例:发布和消费

下面使用 Python 的 pika 客户端演示 RabbitMQ 的基本模型。需要先启动 RabbitMQ,并安装客户端:

python -m pip install pika

创建 rabbitmq_demo.py

from __future__ import annotations

import json
import os
import uuid

import pika


URL = os.environ.get("RABBITMQ_URL", "amqp://guest:guest@localhost:5672/%2F")
EXCHANGE = "orders"
QUEUE = "order.created.worker"
ROUTING_KEY = "order.created"


def connection() -> pika.BlockingConnection:
    parameters = pika.URLParameters(URL)
    return pika.BlockingConnection(parameters)


def declare_topology(channel: pika.adapters.blocking_connection.BlockingChannel) -> None:
    channel.exchange_declare(
        exchange=EXCHANGE,
        exchange_type="topic",
        durable=True,
    )
    channel.queue_declare(queue=QUEUE, durable=True)
    channel.queue_bind(
        exchange=EXCHANGE,
        queue=QUEUE,
        routing_key=ROUTING_KEY,
    )


def publish() -> None:
    with connection() as conn:
        channel = conn.channel()
        declare_topology(channel)

        message = {
            "message_id": str(uuid.uuid4()),
            "order_id": "order-1001",
            "amount": 199,
        }

        channel.confirm_delivery()

        channel.basic_publish(
            exchange=EXCHANGE,
            routing_key=ROUTING_KEY,
            body=json.dumps(message).encode("utf-8"),
            properties=pika.BasicProperties(
                content_type="application/json",
                delivery_mode=pika.DeliveryMode.Persistent,
                message_id=message["message_id"],
            ),
            mandatory=True,
        )

        print("published:", message)


def consume() -> None:
    with connection() as conn:
        channel = conn.channel()
        declare_topology(channel)

        # 每次最多允许 10 条未确认消息
        channel.basic_qos(prefetch_count=10)

        def handle(
            ch: pika.adapters.blocking_connection.BlockingChannel,
            method: pika.spec.Basic.Deliver,
            properties: pika.BasicProperties,
            body: bytes,
        ) -> None:
            message = json.loads(body)
            print("received:", message)

            try:
                # 这里执行真实业务,例如写数据库
                if message["amount"] < 0:
                    raise ValueError("amount must not be negative")

                # 业务成功后 ACK
                ch.basic_ack(delivery_tag=method.delivery_tag)
                print("acked:", message["message_id"])

            except ValueError as exc:
                print("permanent failure:", exc)

                # 不重新入队,交给死信配置;若没有 DLX,则消息会被丢弃
                ch.basic_nack(
                    delivery_tag=method.delivery_tag,
                    requeue=False,
                )

            except Exception as exc:
                print("temporary failure:", exc)

                # 临时故障可以重新入队,但必须有次数上限,
                # 否则会形成高速重试循环。
                ch.basic_nack(
                    delivery_tag=method.delivery_tag,
                    requeue=True,
                )

        channel.basic_consume(
            queue=QUEUE,
            on_message_callback=handle,
            auto_ack=False,
        )

        print("waiting for messages")
        channel.start_consuming()


if __name__ == "__main__":
    import sys

    if len(sys.argv) != 2 or sys.argv[1] not in {"publish", "consume"}:
        raise SystemExit("usage: python rabbitmq_demo.py publish|consume")

    if sys.argv[1] == "publish":
        publish()
    else:
        consume()

运行:

python rabbitmq_demo.py consume
python rabbitmq_demo.py publish

预期结果类似:

published: {'message_id': '...', 'order_id': 'order-1001', 'amount': 199}
received: {'message_id': '...', 'order_id': 'order-1001', 'amount': 199}
acked: ...

这里有三个关键前置条件:

  1. durable=True 只表示 Exchange 和 Queue 的定义可持久化,不代表所有消息自动可靠;
  2. DeliveryMode.Persistent 表示消息要求持久化,但生产者仍应使用 publisher confirms;
  3. auto_ack=False 让应用明确决定什么时候完成消费。

RabbitMQ 的手动 ACK 应该在业务处理完成后发送。若消费者在 ACK 前断开连接,未确认消息会被重新入队,因此消费者必须能够处理重复投递。RabbitMQ 也会通过 redelivered 标志提示该投递可能是重新投递,但这个标志不能替代业务幂等。(rabbitmq.com)


三、RabbitMQ 的 ACK、NACK、Reject 和 Prefetch

3.1 ACK 的真实含义

RabbitMQ 的 ACK 不是“我看见了消息”,而是:

当前消费者已经完成了对这次投递的责任,Broker 可以认为这条消息不再需要重新投递。

因此以下代码顺序是危险的:

def unsafe_handle(ch, method, body):
    ch.basic_ack(method.delivery_tag)
    save_to_database(body)

如果 ACK 成功后进程在 save_to_database() 前崩溃,RabbitMQ 不会再投递消息,而数据库也没有对应记录,结果是消息丢失。

通常应当反过来:

def safe_handle(ch, method, body):
    save_to_database(body)
    ch.basic_ack(method.delivery_tag)

此时如果数据库已经写入,但 ACK 发送前进程崩溃,消息会再次投递,结果可能是数据库重复写入。因此安全性从“可能丢失”变成了“可能重复”,这正是 至少一次消费 的典型形态。

3.2 basic.nack(requeue=True) 的边界

负确认有两条主要路径:

处理成功 ──> ACK ──> Broker 删除或确认完成

临时失败 ──> NACK(requeue=True) ──> 重新进入队列
永久失败 ──> NACK(requeue=False) ──> DLX 或丢弃

requeue=True 适合临时故障,例如数据库短暂不可用;但它不适合数据格式错误、必填字段缺失、业务规则永久不满足等情况。

错误消息如果每次都立即重新入队,会形成:

取出 ──失败──> 立即入队 ──> 再取出 ──失败──> ...

这会导致:

  • 单条坏消息占据消费者;
  • CPU 和网络被无效重试消耗;
  • 正常消息得不到处理;
  • Broker 指标显示吞吐很高,但有效处理量为零。

3.3 Prefetch 是并发窗口,不是线程数

prefetch_count=N 表示一个 Channel 上最多允许大约 N 条未 ACK 投递处于处理中。它控制的是 未确认消息窗口,不是 Python 线程数,也不是 Queue 长度。

若每条消息平均处理时间为 TT,消费者并发窗口为 NN,理想吞吐上限近似为:

λmaxNT\lambda_{\text{max}} \approx \frac{N}{T}

其中:

  • NN 是同时在途的消息数;
  • TT 是单条消息从投递到 ACK 的平均耗时;
  • λmax\lambda_{\text{max}} 是每秒可完成的消息数。

例如:

  • 单条处理时间 T=100 ms=0.1 sT=100\text{ ms}=0.1\text{ s}
  • prefetch_count=10

则理想吞吐约为:

100.1=100 条/秒\frac{10}{0.1}=100\text{ 条/秒}

但这是上限,不包含网络、Broker、数据库锁、Python 调度和尾延迟。若把 prefetch 调成 10,000,并不会自动得到一万倍并发;它可能只会让更多消息堆积在消费者进程内存中。

RabbitMQ 文档也明确区分了 prefetch 窗口与消费者处理能力,并指出值为 0 表示不限制未确认消息数量,这在处理不稳定或消息体较大的消费者中可能造成内存压力。(rabbitmq.com)


四、RabbitMQ 的发布可靠性:Persistent 不等于 Confirmed

发布端至少有以下几种状态:

应用调用 publish
      │
      ├── 未到达 Broker
      ├── Broker 收到但不可路由
      ├── Broker 路由到 Queue
      ├── 消息持久化
      └── Broker 返回 publisher confirm

Publisher confirm 是 Broker 对生产者的确认。它回答的是:

Broker 是否已经接受并负责这条发布消息?

它不回答:

  • 消费者是否处理成功;
  • 数据库是否写入成功;
  • 下游 HTTP 请求是否成功。

如果消息发布到持久化 Queue,生产者应同时考虑:

channel.confirm_delivery()
channel.basic_publish(
    exchange="orders",
    routing_key="order.created",
    body=b"...",
    properties=pika.BasicProperties(
        delivery_mode=pika.DeliveryMode.Persistent
    ),
    mandatory=True,
)

mandatory=True 用于让不可路由消息通过返回机制通知生产者;否则消息可能没有进入任何 Queue。Publisher confirm 与 mandatory return 解决的是不同问题:

  • confirm:Broker 是否接管了发布;
  • return:消息是否存在无法路由的情况。

RabbitMQ 对持久消息的确认时机与 Queue 类型、持久化状态和复制策略有关;应用不应假设 confirm 一定以发布顺序到达,也不应把 confirm 当作业务处理完成信号。(rabbitmq.com)


五、RabbitMQ 重试:不要把“重新入队”当作延迟重试

5.1 直接重新入队的问题

最简单的重试是:

ch.basic_nack(delivery_tag=tag, requeue=True)

它的语义是“现在再试一次”,不是“过 30 秒再试”。如果目标服务需要 30 秒恢复,这种方式会不断快速失败。

更危险的是,RabbitMQ 的 Queue 通常不提供一个通用的“对当前消息睡眠后再次投递”语义。消费者在处理函数中 sleep(30) 会占用一个消费槽位,且 ACK 超时、连接心跳和整体吞吐都可能受到影响。

5.2 延迟队列拓扑

常见方案是使用 TTL Queue 和 Dead Letter Exchange 构造延迟重试:

主队列 orders.main
   │ 消费失败,发布到 retry queue
   ▼
重试队列 orders.retry.30s
   │ 消息 TTL 30 秒后过期
   │ 过期消息由 DLX 转发
   ▼
主 Exchange
   │
   └──> 主队列 orders.main

超过最大次数
   │
   └──> orders.dlq

示例拓扑声明:

import pika

connection = pika.BlockingConnection(
    pika.ConnectionParameters("localhost")
)
channel = connection.channel()

channel.exchange_declare(
    exchange="orders.main",
    exchange_type="direct",
    durable=True,
)

channel.exchange_declare(
    exchange="orders.retry",
    exchange_type="direct",
    durable=True,
)

channel.exchange_declare(
    exchange="orders.dlx",
    exchange_type="direct",
    durable=True,
)

channel.queue_declare(
    queue="orders.main.queue",
    durable=True,
    arguments={
        "x-dead-letter-exchange": "orders.dlx",
        "x-dead-letter-routing-key": "orders.dead",
    },
)

channel.queue_declare(
    queue="orders.retry.30s",
    durable=True,
    arguments={
        "x-message-ttl": 30_000,
        "x-dead-letter-exchange": "orders.main",
        "x-dead-letter-routing-key": "orders.created",
    },
)

channel.queue_declare(
    queue="orders.dlq",
    durable=True,
)

channel.queue_bind(
    exchange="orders.main",
    queue="orders.main.queue",
    routing_key="orders.created",
)

channel.queue_bind(
    exchange="orders.dlx",
    queue="orders.dlq",
    routing_key="orders.dead",
)

connection.close()

失败后不要简单地 requeue=True,而是:

  1. 给消息增加 attempt
  2. 把消息发布到 orders.retry.30s
  3. 确认重试消息发布成功;
  4. 对原始投递发送 ACK
  5. TTL 到期后由 RabbitMQ 重新路由回主队列。

伪代码如下:

def handle_failure(ch, method, properties, body):
    message = decode(body)
    attempt = int(message.get("attempt", 0)) + 1

    if attempt > 3:
        publish_to_dlq(message, reason="max_attempts_exceeded")
        ch.basic_ack(method.delivery_tag)
        return

    retry_message = {
        **message,
        "attempt": attempt,
    }

    publish_to_retry_queue(retry_message, delay="30s")
    ch.basic_ack(method.delivery_tag)

这里的 ACK 顺序很重要。若先 ACK 原消息,再发布重试消息,而发布过程失败,消息会丢失;因此生产环境需要对重试发布使用 publisher confirm,只有确认成功后才 ACK 原投递。

另一方面,若重试发布成功后应用在 ACK 前崩溃,原消息可能重新投递并再次发布一条重试消息。这又会造成重复,因此重试消息也必须具备幂等标识,或使用消息 ID 去重。

5.3 TTL 和死信的准确语义

RabbitMQ 的死信并不只表示“处理失败”。消息可能因为以下原因进入 Dead Letter Exchange:

  • 消费者以 requeue=False 拒绝;
  • 消息过期;
  • Queue 达到长度限制;
  • 某些队列类型或策略触发死信。

配置了 DLX 时,Broker 会把符合条件的消息重新发布到指定 Dead Letter Exchange;没有配置时,消息可能被丢弃。(rabbitmq.com)

因此,死信队列不是自动修复队列。它只是把无法继续正常处理的消息保存到一个可观察、可审计、可人工或程序化恢复的位置。


六、死信队列的生产设计

死信消息至少应携带这些信息:

{
  "message_id": "msg-123",
  "event_type": "order.created",
  "payload": {
    "order_id": "order-1001"
  },
  "failed_at": "2026-09-01T10:00:00Z",
  "attempt": 3,
  "error_type": "ValidationError",
  "error_message": "amount must be positive",
  "consumer": "billing-worker",
  "trace_id": "trace-456"
}

其中 payload 保留原始业务数据,错误元数据用于诊断。不要只记录字符串 "failed",否则死信无法回答:

  • 哪条消息失败;
  • 失败了几次;
  • 由哪个版本的消费者处理;
  • 是暂时故障还是永久错误;
  • 是否可以安全重放。

6.1 死信处理方式

常见处理流程如下:

DLQ
 │
 ├── 查看原因
 ├── 修复代码或外部依赖
 ├── 确认消息是否已部分生效
 ├── 修改或保留原 message_id
 └── 重新发布到主 Exchange

重放之前必须判断业务副作用。例如订单扣款请求已经成功,但消费者在 ACK 前崩溃,此时重放可能再次扣款。正确做法不是盲目重放,而是让扣款接口以 message_id 或业务幂等键实现去重。

RabbitMQ 的死信重新发布本身也不应被误认为带有完整的端到端发布确认;应用若要求死信不丢,应对自己的重放发布过程启用 confirms,并记录重放结果。(rabbitmq.com)


七、Kafka 的消息模型:Topic、Partition、Offset 和 Consumer Group

Kafka 的核心结构是追加日志:

Topic: orders

Partition 0: offset 0 ──> 1 ──> 2 ──> 3 ──> ...
Partition 1: offset 0 ──> 1 ──> 2 ──> ...
Partition 2: offset 0 ──> 1 ──> 2 ──> ...

7.1 Topic 和 Partition

Topic 是逻辑上的消息流名称。

Partition 是 Topic 的有序日志分片。Kafka 只能保证单个 Partition 内的顺序,不能自动保证整个 Topic 的全局顺序。

如果生产者使用同一个 key:

key = order_id

Kafka 通常会把相同 key 的消息映射到同一个 Partition,从而保证同一订单的事件顺序:

order-1001: CREATED → PAID → SHIPPED

但这并不保证不同订单之间的顺序:

order-1001: CREATED
order-1002: PAID
order-1001: PAID

7.2 Offset

每条消息在 Partition 中有一个 offset。消费者有两个容易混淆的位置:

  • position:消费者已经拉取到哪里,下一条准备读取什么;
  • committed offset:消费者把进度持久化到哪里,重启后从哪里恢复。

Kafka 官方客户端文档明确指出,提交的 offset 应该是“下一条要处理的消息”的 offset,而不是当前消息的 offset。(kafka.apache.org)

假设 Partition 中有:

offset:  10     11     12
message:  A      B      C

消费者处理完 A 后,应该提交:

committed offset = 11

因为 offset 11 是下一条尚未处理的消息。

如果错误地提交 10,重启后会再次读取 A;如果错误地提交 12,重启后会跳过 B。

7.3 Consumer Group

Consumer Group 是一组共同消费 Topic 的消费者。一个 Partition 在同一时刻通常只分配给同一 Group 中的一个消费者:

Topic 有 4 个 Partition

Consumer Group G:
  Consumer 1 -> P0, P1
  Consumer 2 -> P2
  Consumer 3 -> P3

如果消费者数量超过 Partition 数量,多出来的消费者没有 Partition 可消费;因此 Kafka 的并行度上限通常受 Partition 数量约束。

不同 Group 拥有不同消费进度:

orders
 ├── group=billing      从 offset 100 消费
 ├── group=analytics    从 offset 80 消费
 └── group=warehouse   从 offset 120 消费

这与 RabbitMQ 中“同一个 Queue 的消息通常只交给一个消费者”不同。Kafka 的消息可以按保留策略继续存在,并被多个 Group 独立读取。


八、Kafka 消费语义:自动提交为什么危险

Kafka 常见消费循环可以抽象为:

poll()
  │
  ├── position 前进
  ├── 业务处理
  └── commit offset

关键问题是:业务处理和 offset 提交不是天然的一个原子事务

8.1 先提交再处理:至多一次

poll message A
commit offset 11
process A

如果 process A 前进程崩溃,重启后从 offset 11 开始,A 不会再读取。

因此:

  • A 可能丢失;
  • A 不会重复处理。

这接近 至多一次

P(重复)0,P(丢失)>0P(\text{重复}) \approx 0,\quad P(\text{丢失}) > 0

8.2 先处理再提交:至少一次

poll message A
process A
commit offset 11

如果处理完成但提交前崩溃:

poll A
process A 成功
进程崩溃

重启后仍从 offset 10 读取 A,于是 A 被再次处理。

此时:

P(丢失)0,P(重复)>0P(\text{丢失}) \approx 0,\quad P(\text{重复}) > 0

这就是通常更安全的 至少一次

Kafka 默认允许至少一次语义;禁用生产者重试并在处理前提交 offset,可以构造至多一次语义,但要承担丢消息风险。(kafka.apache.org)

8.3 自动提交的失败窗口

enable.auto.commit=true 时,客户端会周期性提交已拉取进度,而不是等待你的业务函数成功。

时间线可能是:

t0: poll() 拉取 A、B、C
t1: 自动提交 offset 13
t2: 处理 A 成功
t3: 处理 B 时进程崩溃

重启后 offset 已经是 13,B 和 C 都不会再次读取。自动提交适合可以容忍丢失、或业务处理本身只是缓存预取的场景;对于订单、账务、库存等业务,通常需要关闭自动提交,改为处理成功后手动提交。


九、Kafka 的重平衡和 poll() 生命周期

Kafka 消费者不仅处理消息,还必须参加 Consumer Group 的协调。消费者长时间不调用 poll(),可能被认为失效并触发重平衡。

两个相关参数分别限制不同问题:

  • max.poll.interval.ms:两次 poll() 之间允许的最长处理间隔;
  • max.poll.records:一次 poll() 最多返回的记录数。

若一次返回 500 条消息,而每条处理 1 秒,那么最坏处理时间约为:

500×1 s=500 s500 \times 1\text{ s}=500\text{ s}

如果 max.poll.interval.ms=300000,也就是 300 秒,消费者可能在处理完这一批之前离开 Group,触发重平衡。

因此不能只看平均吞吐,还要估算:

Tbatch=Nrecords×TrecordT_{\text{batch}} = N_{\text{records}} \times T_{\text{record}}

并满足:

Tbatch<max.poll.interval.msT_{\text{batch}} < \text{max.poll.interval.ms}

实际处理时间有长尾时,应使用更保守的分位数,而不是平均值。

Python 客户端经常通过后台线程或本地缓冲维持 poll,再把业务处理提交给线程池或进程池。但这样会引入新的问题:提交 offset 时必须确保对应消息真的处理成功,不能因为 poll() 已经返回就提前提交。


十、RabbitMQ 与 Kafka 的消费语义对比

维度 RabbitMQ Kafka
主要模型 Queue 中的消息投递 Partition 中的持久化日志
消费进度 ACK/NACK offset commit
消息完成后 通常可被删除 仍可按保留策略存在
并发单位 Queue、Consumer、Channel、prefetch Partition、Consumer Group
顺序 受 Queue、消费者并发和重新入队影响 单 Partition 内有序
重试 NACK、重新入队、TTL+DLX 重读 offset、Retry Topic、DLT
广播 多个 Queue 绑定同一 Exchange 多个 Consumer Group
适合场景 任务分发、路由、工作队列 事件流、日志、回放、流处理

不能简单地说“Kafka 比 RabbitMQ 快”或“RabbitMQ 比 Kafka 可靠”。它们优化的抽象不同:

  • 如果重点是把任务交给一个 Worker,并在成功后确认删除,RabbitMQ 的 Queue 模型自然;
  • 如果重点是保留事件、允许多个系统独立回放和按 Partition 扩展,Kafka 的日志模型自然;
  • 如果需要复杂路由、临时队列、请求响应,RabbitMQ 往往更直接;
  • 如果需要高吞吐追加写、消费者独立追赶历史数据,Kafka 更合适。

十一、重试策略:按错误类型分类,而不是统一重试

消息失败至少分为三类。

11.1 可恢复错误

例如:

  • 数据库连接暂时中断;
  • 下游服务返回 502;
  • 网络超时;
  • 限流响应;
  • 依赖服务正在发布。

这类错误可以重试,但应使用退避:

dn=min(dmax,d0×2n)+Jd_n=\min(d_{\max}, d_0 \times 2^n)+J

其中:

  • d0d_0:初始延迟;
  • nn:已经失败的次数;
  • dmaxd_{\max}:最大延迟;
  • JJ:随机抖动,避免大量消费者同时重试。

例如 d0=1d_0=1 秒、dmax=30d_{\max}=30 秒:

第 1 次:约 1 秒
第 2 次:约 2 秒
第 3 次:约 4 秒
第 4 次:约 8 秒
第 5 次:约 16 秒
第 6 次:约 30 秒

11.2 永久错误

例如:

  • JSON 无法解析;
  • Protobuf schema 版本不兼容;
  • 缺少订单号;
  • 金额为负数;
  • 业务状态不允许执行。

这类消息不应无限重试。重试只能重复得到相同结果,最终应进入 DLQ 或 DLT。

11.3 未知错误

未知异常不能简单归类为“可恢复”。应记录完整堆栈、消息标识、消费者版本和依赖状态,并设置有限重试次数。超过次数后进入死信,否则一个程序 Bug 会让整条消息流不断自旋。


十二、Kafka 的重试和死信:用 Topic 表达状态

Kafka 没有 RabbitMQ 那种由 Queue TTL 和 DLX 直接表达的通用死信拓扑。生产中通常通过多个 Topic 表达重试层级:

orders.main
   │ 失败
   ▼
orders.retry.10s
   │ 到期后由转发程序重新发布
   ▼
orders.main

orders.retry.1m
   │
   ▼
orders.main

orders.dlt
   │ 永久失败或超过次数

一种简单的重试消息结构:

{
  "message_id": "msg-123",
  "original_topic": "orders",
  "original_partition": 2,
  "original_offset": 1050,
  "attempt": 2,
  "available_at": "2026-09-01T10:01:30Z",
  "error_type": "TimeoutError",
  "payload": {
    "order_id": "order-1001"
  }
}

重试消费者检查 available_at

from datetime import datetime, timezone
import time


def ready(message: dict) -> bool:
    available_at = datetime.fromisoformat(
        message["available_at"].replace("Z", "+00:00")
    )
    return datetime.now(timezone.utc) >= available_at


def retry_loop(consumer, producer) -> None:
    for record in consumer:
        message = decode(record.value)

        if not ready(message):
            time.sleep(0.2)
            continue

        producer.produce(
            topic=message["original_topic"],
            key=record.key,
            value=encode(message["payload"]),
        )
        producer.flush()

        # 只有重新发布成功后,才提交 retry topic 的 offset
        consumer.commit()

这个示例用于说明状态转移,生产代码还必须解决两个问题:

  1. produce() 成功调用不等于 Broker 已确认,需要检查 delivery callback 或使用客户端提供的确认机制;
  2. 重新发布成功后、提交 retry offset 前崩溃,会导致原消息重复发布,因此目标消费者必须幂等。

更常见的实现是使用专门的重试转发器,或让消费者通过 seek() 重新读取失败位置。但直接停在一个失败 offset 上会阻塞同一 Partition 后续消息:

P0: A(失败) → B(正常) → C(正常)
             ▲
             └── 如果 A 一直重试,B、C 也无法前进

这体现了 Kafka 顺序和重试之间的冲突:

  • 保留 Partition 顺序:失败消息阻塞后续消息;
  • 追求整体吞吐:把失败消息移到重试 Topic,但可能改变处理顺序。

若业务要求同一订单严格按顺序处理,重试 Topic 设计必须保留同一 key 的路由策略,并避免在旧事件未成功时处理新事件。


十三、消费语义的形式化推导

设:

  • PP:生产者发布;
  • AA:Broker 接收并确认;
  • HH:消费者开始处理;
  • SS:业务副作用成功;
  • CC:消费者确认或提交 offset。

最危险的窗口是:

H ──> S ──崩溃──> C

业务已经成功,但消费进度没有成功保存。恢复后会再次执行 SS

另一种危险窗口是:

H ──> C ──崩溃──> S

消费进度已经保存,但业务副作用尚未发生。恢复后不会再执行 SS

因此:

  • CS 前:至多一次,可能丢失;
  • CS 后:至少一次,可能重复。

要得到端到端的“看起来只生效一次”,需要让以下操作原子化:

业务副作用+消费进度提交\text{业务副作用} + \text{消费进度提交}

但 RabbitMQ ACK 或 Kafka offset 通常由 Broker 管理,业务数据库由另一个系统管理,二者不能自动组成同一个本地事务。

13.1 幂等表解决重复

假设消费者要处理消息并写订单账务表,可以建立消费记录表:

CREATE TABLE consumed_messages (
    message_id VARCHAR(128) PRIMARY KEY,
    consumed_at TIMESTAMP NOT NULL
);

CREATE TABLE account_entries (
    entry_id VARCHAR(128) PRIMARY KEY,
    order_id VARCHAR(128) NOT NULL,
    amount DECIMAL(18, 2) NOT NULL
);

在同一个数据库事务中:

BEGIN;

INSERT INTO consumed_messages(message_id, consumed_at)
VALUES ('msg-123', CURRENT_TIMESTAMP)
ON CONFLICT (message_id) DO NOTHING;

程序检查插入影响行数:

影响 1 行:第一次处理,继续执行副作用
影响 0 行:message_id 已处理,跳过副作用

然后:

INSERT INTO account_entries(entry_id, order_id, amount)
VALUES ('msg-123', 'order-1001', 199.00);

COMMIT;

最后再 ACK RabbitMQ 或提交 Kafka offset。

即使最后一步前崩溃,消息再次到达时,message_id 唯一约束也会阻止重复副作用。

注意:这要求业务副作用也在同一个数据库事务中。若处理逻辑是“调用外部支付接口”,数据库幂等表不能自动回滚外部支付,因此外部接口也必须支持幂等键,或者使用本地 Outbox、状态机和可查询的结果确认。


十四、Outbox:解决数据库提交和消息发布不一致

另一个常见问题是:

写订单数据库成功
发布 Kafka/RabbitMQ 消息失败

或:

发布消息成功
写订单数据库失败

如果数据库和 Broker 不在同一个事务中,直接执行两个动作就存在双写不一致。

Outbox 模式把消息先写入业务数据库:

业务事务
 ├── 写 orders
 └── 写 outbox_events(status='pending')
          │
          ▼
Outbox Relay
          │ publisher confirm
          ▼
RabbitMQ / Kafka
          │
          ▼
更新 outbox_events(status='published')

事务内:

BEGIN;

INSERT INTO orders(order_id, status)
VALUES ('order-1001', 'created');

INSERT INTO outbox_events(
    event_id,
    event_type,
    aggregate_id,
    payload,
    status
)
VALUES (
    'event-1',
    'order.created',
    'order-1001',
    '{"order_id":"order-1001"}',
    'pending'
);

COMMIT;

Relay 程序扫描 pending 记录并发布消息。发布成功后再标记为 published

Relay 在发布成功和更新状态之间崩溃时,会再次发布同一 event_id。因此 Outbox 通常提供的是:

数据库内可靠保存 + 消息至少一次发布

而不是天然 exactly-once。下游仍然需要通过 event_id 幂等。


十五、Exactly-once 到底是什么意思

“Exactly once”必须先说明作用范围。

15.1 投递 exactly-once

如果意思是:

Broker 只把一条消息交给消费者一次,消费者绝不会重复看到它。

那么在存在网络重试、消费者崩溃和 ACK 丢失时,这个目标通常不成立。

15.2 处理 exactly-once

如果意思是:

某个业务副作用在所有故障情况下只执行一次。

那么外部系统必须参与事务或提供幂等协议。

15.3 Kafka 内部 exactly-once

Kafka 支持事务生产和 Kafka Streams 的 exactly-once processing,但它的保证范围主要是 Kafka 输入 offset、Kafka 输出记录以及 Kafka Streams 状态存储之间的原子处理。Kafka 官方设计文档特别提醒,exactly-once 需要阅读适用边界,不能直接推广为所有外部副作用都只发生一次。(kafka.apache.org)

例如:

Kafka input
   │
Kafka Streams 事务
   ├── 提交 input offset
   ├── 写 Kafka output
   └── 更新本地 state store

这三者可以在 Kafka 事务模型中协调。

但如果处理过程还调用:

Kafka input
   │
   ├── 调用支付 API
   ├── 写 MySQL
   └── 写 Kafka output

Kafka 事务并不会自动回滚已经成功的支付 API 或 MySQL 操作。此时仍需要外部系统幂等、事务消息、Outbox 或 Saga 等机制。


十六、Python 中的异步消费者并发

Python 3.14 的 asyncio 提供基于 async/await 的异步 I/O 并发模型,适合网络 I/O 密集型代码;但一个事件循环中的 Python 代码仍然需要避免长时间阻塞,否则其他任务无法运行。(docs.python.org)

一个受限并发的消费者结构可以写成:

import asyncio
from dataclasses import dataclass


@dataclass
class Message:
    message_id: str
    payload: dict


async def handle(message: Message) -> None:
    await asyncio.sleep(0.1)  # 模拟数据库或 HTTP I/O
    print("processed", message.message_id)


async def worker(
    queue: asyncio.Queue[Message],
    semaphore: asyncio.Semaphore,
) -> None:
    while True:
        message = await queue.get()
        try:
            async with semaphore:
                await handle(message)
                # 只有 handle 成功后,才确认消息
                await ack(message)
        except Exception:
            await retry_or_dead_letter(message)
        finally:
            queue.task_done()


async def ack(message: Message) -> None:
    print("ack", message.message_id)


async def retry_or_dead_letter(message: Message) -> None:
    print("retry or dead-letter", message.message_id)


async def main() -> None:
    queue: asyncio.Queue[Message] = asyncio.Queue(maxsize=100)
    semaphore = asyncio.Semaphore(20)

    workers = [
        asyncio.create_task(worker(queue, semaphore))
        for _ in range(4)
    ]

    for i in range(100):
        await queue.put(
            Message(
                message_id=f"msg-{i}",
                payload={"index": i},
            )
        )

    await queue.join()

    for task in workers:
        task.cancel()

    await asyncio.gather(*workers, return_exceptions=True)


asyncio.run(main())

这里有两层并发限制:

  • Queue(maxsize=100):限制本地缓冲,避免无限拉取;
  • Semaphore(20):限制同时执行的业务处理数。

如果 RabbitMQ 的 prefetch 是 1,000,而本地只允许 20 个任务处理,其余消息可能已经从 Broker 投递到进程内,但尚未真正执行。若进程此时崩溃,这些未 ACK 消息会重新投递;这通常没问题,但会增加重复和内存占用。

如果使用阻塞式数据库客户端、阻塞式 HTTP 客户端或 CPU 密集计算,不能直接放在事件循环中:

async def bad_handler(message):
    result = blocking_database_call(message)

应使用异步客户端,或显式放入线程池/进程池。否则 ACK 延迟、心跳响应和消费吞吐都会受到影响。


十七、与 FastAPI、ASGI 和 Celery 的关系

17.1 FastAPI 请求处理不等于消息消费

FastAPI 运行在 ASGI 服务器上,Web 请求生命周期与消息消费者生命周期不同。不要在每一个请求中临时创建一个 RabbitMQ 或 Kafka Consumer:

HTTP 请求
 ├── 创建消费者
 ├── 消费一条消息
 ├── 关闭消费者
 └── 返回响应

这会造成连接创建开销、消费进度不稳定和重复消费。

更合理的关系是:

FastAPI Web 进程
 └── 接收请求并发布消息

独立 Worker 进程
 └── 长期运行并消费消息

ASGI 应用的 lifespan 可以用于管理应用级资源,但长时间消费任务仍应明确处理取消、关闭和部署副本问题。FastAPI 官方文档和 ASGI 规范分别描述了应用生命周期和异步服务器接口;它们不会自动替你定义 RabbitMQ ACK 或 Kafka offset 语义。

17.2 Celery 与 RabbitMQ/Kafka

Celery 是任务队列框架,Broker 可以使用 RabbitMQ 或其他支持的传输后端。它在底层仍然面对相同问题:

任务被取出
  ├── Worker 成功执行
  ├── Worker 失败
  ├── Worker 在执行后、确认前崩溃
  └── Broker 连接中断

因此 Celery 的 ACK、重试、定时和幂等不是独立于消息系统的魔法,而是对上述状态的更高层封装:

  • ACK 过早,可能丢任务;
  • ACK 过晚,可能重复执行;
  • 自动重试需要次数和退避;
  • 定时任务仍需防止重复触发;
  • 业务函数必须幂等。

如果 Celery 任务调用外部支付、发邮件或修改库存,必须把任务 ID、业务 ID 或幂等键写入业务状态,而不能因为“任务队列保证至少一次”就假设业务只执行一次。

17.3 与 gRPC 的组合

gRPC 常用于消费者调用下游服务。此时消息系统的重试和 gRPC 的重试必须分层设计:

消息重试
 └── gRPC 调用
      └── gRPC 内部重试

如果 gRPC 客户端一次调用已经重试 3 次,消息消费者又对整条消息重试 5 次,最坏调用次数可能达到:

3×5=153 \times 5 = 15

如果每一层都设置无限重试,系统可能长时间不 ACK,最终触发重复投递或消费者失活。

应区分:

  • gRPC DEADLINE_EXCEEDED:通常是超时,可有限重试;
  • UNAVAILABLE:服务暂时不可用,可退避重试;
  • INVALID_ARGUMENT:请求永久无效,应进入死信;
  • ALREADY_EXISTS:可能表示幂等请求已经成功,需要查询状态,而不是简单重试。

消息处理的总 Deadline 也应覆盖:

数据库读取 + gRPC 调用 + 本地处理 + ACK/offset 提交

若单条消息的最大允许处理时间为 DD,则应满足:

Tdb+Tgrpc+Tbusiness+Tcommit<DT_{\text{db}} + T_{\text{grpc}} + T_{\text{business}} + T_{\text{commit}} < D

否则消费者可能在业务尚未完成时被 Broker 认为失活,造成重平衡或重复投递。


十八、失败表现和诊断方法

18.1 RabbitMQ 消息不断重复

优先检查:

  1. 是否在业务成功前 ACK;
  2. 消费者连接是否频繁断开;
  3. Channel 是否发生协议异常;
  4. requeue=True 是否处理了永久错误;
  5. 是否存在多个消费者修改同一业务对象;
  6. prefetch 是否远大于实际处理能力。

RabbitMQ 会在连接或 Channel 关闭时重新入队未确认消息,因此重复投递可能是故障恢复的正常结果,而不是 Broker 重复生成了消息。(rabbitmq.com)

18.2 RabbitMQ 消息消失

检查:

  • 是否启用了 auto_ack
  • Queue 是否 durable;
  • 消息是否 persistent;
  • 是否启用 publisher confirms;
  • Exchange 是否存在正确 Binding;
  • 是否使用 mandatory=True 处理不可路由消息;
  • DLX 是否配置错误;
  • Queue 是否因 TTL 或长度限制触发死信。

auto_ack 模式下,消息在发送给消费者后就可能被认为已成功交付;如果消费者随后崩溃,消息可能丢失。(rabbitmq.com)

18.3 Kafka 消费延迟不断增加

检查:

lag = log_end_offset - committed_offset

其中:

  • log_end_offset 是 Partition 当前末尾;
  • committed_offset 是消费组已提交进度;
  • lag 是尚未提交的消息数量。

需要进一步判断:

  • 是生产速度超过消费速度;
  • 还是某个 Partition 的单条慢消息阻塞;
  • 还是消费者频繁重平衡;
  • 还是下游依赖变慢;
  • 还是失败消息被反复重试。

如果只有一个 Partition lag 很高,增加消费者数量通常没有帮助,因为并行度受 Partition 限制。

18.4 Kafka 消费重复

常见时间线是:

处理消息成功
提交 offset 失败或进程崩溃
重启后从旧 offset 读取

这属于至少一次语义的正常后果。诊断时应通过 topic + partition + offset + message_id 判断:

  • 是同一消息重新读取;
  • 还是生产者重复发布了两个不同 offset 的相同业务事件;
  • 还是消费者错误提交了批次边界。

18.5 Kafka 消费跳过消息

重点检查是否启用了自动提交,或代码是否在业务处理前执行了:

consumer.commit()

还要检查批量处理时是否把一批消息全部视为成功。例如:

for record in records:
    process(record)

commit_all(records)

只要中间一条失败而异常被吞掉,后续的整体提交就可能跳过失败消息。


十九、常见误解

误解一:“ACK 表示消息已经被业务系统永久处理”

不一定。RabbitMQ ACK 只表示消费者向 RabbitMQ 确认这次投递;Kafka offset commit 只表示消费进度已保存。业务数据库、支付系统和文件系统是否成功,需要单独保证。

误解二:“持久化消息就不会丢”

RabbitMQ 的 durable Queue、persistent message 和 publisher confirm 分别解决不同层次的问题。缺少发布确认时,生产者无法可靠判断 Broker 是否接管消息。

误解三:“Kafka offset 就是消息 ID”

offset 只在单个 Topic Partition 内有意义。消息迁移到另一个 Topic、Partition 或环境后,offset 不再具有全局身份。业务去重应使用稳定的 message_id 或业务幂等键。

误解四:“重试次数越多越可靠”

重试只能提高暂时性错误的恢复概率。对永久错误,重试次数越多,越会延迟诊断并消耗系统资源。

误解五:“exactly-once 可以让支付只扣一次”

除非支付系统也参与同一事务,或支付接口支持幂等键并能查询历史结果,否则消息 Broker 的 exactly-once 不能保证外部支付只发生一次。

误解六:“提高消费者数量就能提高 Kafka 吞吐”

只有在还有未分配的 Partition 时,增加消费者才可能提高并行度。若 Topic 只有 3 个 Partition,启动 10 个消费者并不能得到 10 路 Partition 并发。


二十、选择和设计的最终落点

可以用下面的判断方式选择模型:

需要任务被一个 Worker 领取并确认完成?
    └── 优先考虑 RabbitMQ Queue / Celery

需要事件长期保留、多个系统独立消费和回放?
    └── 优先考虑 Kafka Topic

需要复杂路由、广播、临时队列、请求响应?
    └── RabbitMQ 的 Exchange 模型更直接

需要按 key 保持顺序并按 Partition 扩展?
    └── Kafka 更自然

需要可靠发布?
    └── RabbitMQ 使用 publisher confirms
    └── Kafka 关注 producer ack、重试、幂等和事务配置

需要可靠消费?
    └── 业务成功后 ACK 或提交 offset
    └── 接受至少一次带来的重复
    └── 通过幂等键消除重复副作用

需要重试?
    └── 区分临时错误和永久错误
    └── 有限次数、指数退避、随机抖动
    └── 超限进入 DLQ/DLT

需要“只生效一次”?
    └── 让业务副作用具备幂等性
    └── 必要时使用数据库事务、Outbox 或外部幂等协议

消息系统最重要的不是某个客户端 API,而是明确每一个状态的所有权:

谁负责保存消息?
谁负责确认发布?
谁负责确认消费?
失败后消息去哪?
重试由谁调度?
重复副作用如何消除?
死信如何恢复?

RabbitMQ 更偏向确认式投递,Kafka 更偏向可回放日志;两者都不能替代业务幂等。只要把发布确认、消费确认、失败重试和外部副作用分别建模,所谓“可靠消息”就不再是模糊承诺,而可以被拆成可验证的状态转换。


系列导航与关联阅读

官方资料

本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。