幂等消费者与恰好一次语义
对消息处理进行去重和检查点记录,实现下游有效的恰好一次传递。
幂等消费者与恰好一次语义 是 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 反馈 — 无需本地设置。
此课程中的所有课时
- 异步生产与消费 Kafka 事件
- 模式注册表与 Avro 契约演进
- 事务型发件箱模式
- 幂等消费者与恰好一次语义