FastAPI Backend Development Bootcamp · 课时

幂等消费者与恰好一次语义

对消息处理进行去重和检查点记录,实现下游有效的恰好一次传递。

第 4 / 4 课13 个步骤

幂等消费者与恰好一次语义 是 CoddyKit 上的免费 FastAPI Backend Development Bootcamp 课时。 这是第 4 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 FastAPI Backend Development Bootcamp 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 FastAPI Backend Development Bootcamp 课程共包含 4 节课。

投递问题

Kafka 和 Pulsar 等消息代理默认采用至少一次投递。如果消费者在处理消息后、确认消息前崩溃,消息代理会在消费者重启后重新投递该消息。

对于从订单事件主题中获取消息的 FastAPI 后端来说,这意味着同一个 OrderPlaced 事件可能会两次到达您的处理程序。没有防护措施时,您可能会向客户扣款两次,或发送两封确认邮件。

  • 最多一次:处理前确认——速度快,但崩溃时会丢失消息。
  • 至少一次:处理后确认——不会丢失消息,但可能重复。
  • 有效的恰好一次:至少一次投递加上幂等处理。

幂等性才是真正的目标

在一般情况下,不可能在网络中实现真正的恰好一次投递。您可以构建的是有效的恰好一次:消息代理可能多次投递同一条消息,但消费者只会恰好一次地应用其效果。

当使用相同输入运行两次后得到的最终状态与运行一次相同时,处理程序就是幂等的。具体策略是:为每条消息分配一个稳定且唯一的键,记录已经处理过的键,并跳过之前见过的任何消息。

选择去重键

去重键必须在重新投递期间保持稳定。消息代理偏移量并不稳定——重新平衡分区可能会改变偏移量,而且 Pulsar 消息 ID 与 Kafka 偏移量不同。

优先使用消息中携带的业务级标识符,或由生产者分配的事件 id:

  • 由生产者设置的 event_id(一个 UUID)——通常是最佳选择。
  • 当每个订单只产生一种事件时,可以使用 order_id 这样的自然键。
  • 只有在别无选择时,才使用 (topic, partition, offset)。

下面的纯去重函数只会记住它见过的键。

def make_deduper():
    seen = set()

    def process(event_id, payload):
        if event_id in seen:
            return "skipped (duplicate)"
        seen.add(event_id)
        return f"processed {payload}"

    return process


if __name__ == "__main__":
    handle = make_deduper()
    events = [
        ("evt-1", "order#100"),
        ("evt-2", "order#101"),
        ("evt-1", "order#100"),  # redelivery
    ]
    for eid, payload in events:
        print(eid, "->", handle(eid, payload))

仅靠内存去重还不够

基于集合的去重器会在 FastAPI 进程每次重启时重置——而崩溃恰恰是发生重新投递的时刻。您需要一个能够在重启后保留数据、并能在多个工作进程副本之间共享的持久化去重存储。

常见的两种后端:

  • 数据库表:存储已处理的 id,具有强一致性,可以与业务数据连接,并支持事务。
  • 带 TTL 的 Redis:速度快,适合重新投递只会在有限时间窗口内发生的场景。

正确的选择取决于重复消息在原始消息之后最晚可能延迟多久到达。

Inbox 表模式

Inbox 模式会将每条已处理消息的 id 存入一个带有 UNIQUE 约束的表中。应用效果之前,先尝试插入该 id;如果违反唯一性约束,说明这是重复消息,因此跳过。

关键是,必须将去重行的插入和业务变更放在同一个数据库事务中。要么两者一起提交,要么两者都不提交——不会出现效果已经生效但 id 尚未记录的时间窗口。

-- Inbox table for the consumer (Postgres)
CREATE TABLE processed_messages (
    event_id   TEXT PRIMARY KEY,
    consumer   TEXT NOT NULL,
    processed_at TIMESTAMPTZ NOT NULL DEFAULT now()
);

-- Same transaction: dedup insert + business write
BEGIN;
INSERT INTO processed_messages (event_id, consumer)
VALUES ('evt-1', 'order-service');  -- fails if already present

UPDATE accounts SET balance = balance - 50
WHERE id = 'acct-7';
COMMIT;

FastAPI 中的幂等消费者

下面展示了如何将事务内去重的思路接入使用 SQLAlchemy 的异步消费者。事件 id 的 INSERT 和业务变更共享同一个事务。如果插入引发 IntegrityError,说明我们已经处理过此事件,于是回滚。

由于这依赖正在运行的数据库和 SQLAlchemy 会话,因此它属于框架/基础设施代码,无法由在线评测系统单独执行。

from sqlalchemy.exc import IntegrityError
from sqlalchemy import text

async def handle_event(session, event_id: str, order_id: str, amount: int):
    try:
        async with session.begin():
            await session.execute(
                text("INSERT INTO processed_messages (event_id, consumer) "
                     "VALUES (:eid, 'order-service')"),
                {"eid": event_id},
            )
            await session.execute(
                text("UPDATE orders SET total = total + :amt WHERE id = :oid"),
                {"amt": amount, "oid": order_id},
            )
    except IntegrityError:
        # Duplicate delivery: effect already applied, safe to ignore
        await session.rollback()
        return "duplicate"
    return "processed"

检查点:在效果生效后提交偏移量

幂等性负责处理重复消息;检查点机制控制消息代理何时将某条消息标记为已处理。黄金法则是:只有在效果已持久化后,才提交偏移量。

请禁用自动提交,并在事务成功后手动提交。这样可以保证至少一次投递:提交前发生崩溃会导致消息被重新处理,而幂等性会吸收这次重放。

  • Kafka:enable.auto.commit=false,然后调用 consumer.commit()。
  • Pulsar:关闭自动确认,并在成功后调用 consumer.acknowledge(msg)。
# Kafka consumer loop with manual commit AFTER processing
from confluent_kafka import Consumer

consumer = Consumer({
    "bootstrap.servers": "localhost:9092",
    "group.id": "order-service",
    "enable.auto.commit": False,   # we commit ourselves
    "auto.offset.reset": "earliest",
})
consumer.subscribe(["orders"])

while True:
    msg = consumer.poll(1.0)
    if msg is None or msg.error():
        continue
    process_with_dedup(msg)          # durable, idempotent
    consumer.commit(message=msg)     # advance offset only now

为什么确认与效果的顺序很重要

操作顺序决定了一切。请比较效果已经应用但随后发生崩溃时的两种顺序:

  • 先提交再处理:偏移量先前进。发生崩溃会丢失消息——最多一次。对于资金操作来说不可接受。
  • 先处理再提交:效果先持久化。发生崩溃会重新处理消息——至少一次。去重会使这次重放无害。

请始终选择先处理再提交,并依靠幂等性。下面的模拟展示了如何吸收一次重新投递。

def consume(messages, crash_after=None):
    seen = set()
    balance = 0
    committed = 0
    for i, (eid, amount) in enumerate(messages):
        if eid not in seen:        # idempotent effect
            seen.add(eid)
            balance += amount
        committed = i + 1          # commit AFTER effect
        if crash_after is not None and i == crash_after:
            break
    return balance, committed


if __name__ == "__main__":
    stream = [("e1", 50), ("e2", 30)]
    # Crash before committing e2, then e2 is redelivered
    bal, _ = consume(stream, crash_after=0)
    replay = stream[1:] + [("e2", 30)]  # redelivery duplicate
    seen = {"e1"}
    for eid, amount in replay:
        if eid not in seen:
            seen.add(eid)
            bal += amount
    print("final balance:", bal)  # 80, not 110

使用 TTL 窗口进行 Redis 去重

当重复消息只可能在有限时间窗口内到达(例如消息代理重试加上重新平衡所需的时间)时,Redis 是一种轻量级去重存储。使用 SET key value NX EX ttl:只有在键不存在时才设置该键,并使其自动过期。

原子性的 NX 标志会将检查和占用合并为一个操作,因此两个并发工作进程无法同时成功占用同一个事件 id。下面是使用过期时钟对这一逻辑进行的纯模拟。

class FakeRedisNX:
    def __init__(self):
        self.store = {}

    def set_nx_ex(self, key, now, ttl):
        exp = self.store.get(key)
        if exp is not None and exp > now:
            return False          # still claimed -> duplicate
        self.store[key] = now + ttl
        return True               # claimed -> first time


if __name__ == "__main__":
    r = FakeRedisNX()
    print(r.set_nx_ex("evt:1", now=0, ttl=60))   # True
    print(r.set_nx_ex("evt:1", now=10, ttl=60))  # False (dup)
    print(r.set_nx_ex("evt:1", now=70, ttl=60))  # True (expired)

天然幂等的操作

有时,您可以让操作本身具备幂等性,从而不必使用显式的去重存储;这样,执行两次和执行一次的结果相同。

  • 插入或更新:根据业务 id 使用 INSERT ... ON CONFLICT DO UPDATE。
  • 集合赋值:使用 status = 'PAID',而不是相对切换。
  • 条件写入:仅当版本或状态保护条件匹配时才应用变更。

像 balance = balance - 50 这样的相对变更并不具备幂等性——执行两次会重复扣减。请将这类效果转换为绝对状态,或使用已处理 id 检查来加以保护。

-- Idempotent upsert keyed by business id
INSERT INTO order_totals (order_id, total)
VALUES ('order-100', 250)
ON CONFLICT (order_id)
DO UPDATE SET total = EXCLUDED.total;

-- Idempotent state transition (absolute, not relative)
UPDATE orders SET status = 'PAID'
WHERE id = 'order-100' AND status = 'PENDING';

端到端的恰好一次处理管道

以 FastAPI 消费者为例,完整流程如下:

  • 生产者为每个事件写入稳定的 event_id(UUID)。
  • 消费者在禁用自动提交的情况下读取消息。
  • 在同一个数据库事务中:将 event_id 插入收件箱(去重),并应用业务效果。
  • 如果发生 IntegrityError,则跳过——该事件已经处理过。
  • 只有在事务提交后,才向代理提交或确认偏移量。

这样就能实现实际上的恰好一次处理:代理提供至少一次投递,通过幂等性实现恰好一次效果,并通过检查点排序避免消息丢失。

快速检查

您运行着一个至少一次投递的 Kafka 消费者,用于扣减账户余额。要实现实际上的恰好一次处理,哪种组合才是正确的?

回顾

您已经学习了如何将至少一次投递转化为实际上的恰好一次处理:

  • 目标是幂等性——让重复处理产生与单次处理相同的状态。
  • 选择稳定的去重键,最好使用生产者分配的 event_id,绝不要使用原始代理偏移量。
  • 使用持久化收件箱(数据库唯一约束或 Redis NX+TTL),使去重信息在重启和副本之间仍然保留。
  • 去重插入和业务效果使用同一个事务——要么全部成功,要么全部失败。
  • 最后再创建检查点:只有在效果持久化后,才提交或确认偏移量(先处理,再提交)。
  • 相较于相对变更,优先使用天然幂等的操作(插入或更新、绝对状态)。
免费开始

用 AI 导师学习 FastAPI Backend Development Bootcamp — 免费

在浏览器中编写并运行真实代码,获得全天候 AI 导师的即时帮助,并在网页或应用中继续学习。

课程
21
课程
84

常见问题解答

「幂等消费者与恰好一次语义」课时是免费的吗?

是的 — 「幂等消费者与恰好一次语义」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 FastAPI Backend Development Bootcamp 课程的其余内容,请升级到 CoddyKit PRO。 FastAPI Backend Development Bootcamp 课程共包含 4 节课。

「幂等消费者与恰好一次语义」这节课中我会学到什么?

对消息处理进行去重和检查点记录,实现下游有效的恰好一次传递。 你通过在浏览器中直接运行的动手代码来练习 FastAPI Backend Development Bootcamp,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 FastAPI Backend Development Bootcamp 需要有经验吗?

无需任何先前经验。CoddyKit 上的 FastAPI Backend Development Bootcamp 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 4 节课,共 4 节。

「幂等消费者与恰好一次语义」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 FastAPI Backend Development Bootcamp 课中编写并运行代码吗?

能。每节 FastAPI Backend Development Bootcamp 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. 异步生产与消费 Kafka 事件
  2. 模式注册表与 Avro 契约演进
  3. 事务型发件箱模式
  4. 幂等消费者与恰好一次语义
← 返回 FastAPI Backend Development Bootcamp