MongoDB Change Streams 可以让应用订阅 collection、database 或整个 deployment 的数据变化,但“能收到事件”不等于“故障后不会漏或重复”。生产实现的关键是把完整 _id 当作 opaque resume token,只有在业务副作用成功后才原子推进 checkpoint;恢复时复用原 pipeline 与 options,并为 token 过期、invalidate、消费积压和重复投递设计明确路径。
Change Streams 适合解决什么问题?
MongoDB 官方文档将 Change Streams 定义为数据库变化订阅机制,应用无需直接 tail oplog,就可以监听单个 collection、database 或 deployment。典型用途包括搜索索引同步、缓存失效、审计投影、通知触发和跨系统派生数据。
它不是消息队列的自动替代品。change event 反映数据库已经提交的变化,但下游消费仍需要处理进程崩溃、网络中断、重复读取、不可恢复的 token、目标系统限流和业务副作用幂等。若任务需要人工重试、优先级、定时执行或复杂补偿,可以先参考 FastAPI 后台任务选型指南,判断事件监听之后是否还需要 durable queue 或 workflow。
Change Streams 需要 replica set 或 sharded cluster。流可以带 aggregation pipeline,用 $match 等阶段减少客户端收到的事件,但不能修改或移除事件的 _id,因为它就是 resume token。MongoDB 从 4.2 起会拒绝破坏该字段的 change stream pipeline。
一条 change event 中哪些字段必须保留?
至少保存以下语义,而不是只取 fullDocument:
| 字段 | 用途 | 常见误区 |
|---|---|---|
_id |
完整 resume token | 自己拆解、转成时间戳或只保存其中一部分 |
operationType |
insert/update/replace/delete/invalidate 等 | 假设所有事件都有完整文档 |
ns |
database 与 collection | 多 collection 消费时丢失来源 |
documentKey |
文档身份,通常含 _id |
delete 时仍去依赖 fullDocument |
updateDescription |
update 的字段变化 | 把它误当作最终完整状态 |
clusterTime |
服务端事件顺序线索 | 当作全局业务版本或外部幂等键 |
Resume token 是 opaque 值。应用应按 BSON/Extended JSON 可逆地保存完整 _id,不能解析后重新拼装,也不要把 clusterTime 当作等价替代。token 的格式可能随服务端版本和事件类型变化。
默认的 update 事件主要给出 delta。若消费者需要当前完整文档,可选择 fullDocument="updateLookup",但这个 lookup 发生在读取事件时;如果文档随后又变化或被删除,返回的内容未必等于事件提交那一刻的精确快照。对严格审计场景,应评估 MongoDB 的 pre/post-images,而不是把 update lookup 当作历史版本库。
checkpoint 应该在什么时候推进?
最安全的规则是:事件的业务副作用确认成功后,再保存该事件的 resume token。
假设消费者收到事件 E,先写 checkpoint,再更新搜索索引;如果索引请求超时并且进程退出,重启后会从 E 之后恢复,E 的副作用可能永久缺失。反过来,先更新索引再保存 checkpoint;若两步之间崩溃,E 会被重新读取,因此下游必须支持幂等。
这形成常见的 at-least-once 处理模型:允许重复,不允许静默漏掉。示例骨架如下,代码只表达顺序,未在本文中实际运行:
with collection.watch(
pipeline,
full_document="updateLookup",
resume_after=checkpoint.resume_token,
) as stream:
for event in stream:
operation_id = encode_token(event["_id"])
apply_idempotently(operation_id, event)
checkpoint_store.compare_and_set(
expected=checkpoint.version,
token=event["_id"],
)
apply_idempotently 不能只是一个名字。对搜索索引,可用数据库文档 _id 作为目标文档键并覆盖最新版本;对通知,可在 outbox/dedup 表中对 operation ID 建唯一约束;对计费或库存,必须使用业务级版本与事务边界,不能假设重复调用无害。
Checkpoint 也需要 CAS。若两个消费者误用同一个 partition/stream identity,较慢实例可能用旧 token 覆盖较新的进度,造成大段重放。保存记录至少包含 stream identity、pipeline digest、options digest、resume token、单调版本和消费者 owner/lease。
ZoyTown 现有的 FastAPI、MongoDB 与 Redis 原子发布文章展示了 revision CAS 与缓存状态分离的思路;Change Streams checkpoint 也应把“Mongo 事件已提交”“下游副作用完成”“checkpoint 已推进”拆成不同状态。
resumeAfter、startAfter 和 startAtOperationTime 如何选择?
MongoDB 允许用 token 恢复 change stream,但三个入口语义不同:
resumeAfter:从指定事件之后继续,适合普通可恢复中断;不能用来跨过已经关闭流的某些 invalidate 场景。startAfter:同样从指定 token 之后开始,但可用于在 invalidate event 之后启动新的流。startAtOperationTime:从一个 operation time 开始,适合没有 token 的受控起点,不应作为日常 checkpoint 的降级替代。
官方特别警告:用 token 恢复时必须复用生成该 token 时相同的 pipeline 和 options。更改过滤条件、fullDocument、collation 或监听范围后继续复用旧 token,可能导致不可预测行为、数据一致性问题或无法恢复。
因此,部署中应给 stream configuration 计算稳定 digest:
stream_id = orders-search-v3
pipeline_digest = sha256(canonical_pipeline)
options_digest = sha256(canonical_options)
只有 digest 一致时自动恢复;配置变化则创建新的 stream identity,并按明确的 backfill 起点启动。不要悄悄沿用旧 checkpoint。
resume token 为什么会失效?
恢复依赖 oplog 中仍保留足够历史。高写入量、较小 oplog 或消费者长时间停机都可能让所需历史滚出窗口。驱动通常会自动恢复部分网络错误,但无法凭空恢复已经不在历史中的位置。
生产监控至少需要以下指标:
- 最近一次成功处理事件距当前时间的 lag;
- 最近 checkpoint 的保存时间与版本;
- 消费速率、失败率与重试队列长度;
- resume 尝试次数及错误分类;
- 下游幂等命中或重复事件数量;
- 业务对账发现的缺失数量。
如果 token 不可恢复,自动“从现在开始”通常会制造静默数据洞。正确策略取决于投影:可重建的搜索索引或缓存应停止增量流,执行基于权威 Mongo 状态的全量/分片 backfill,再以新的 checkpoint 接续;不可重建的审计或外部副作用则需要人工升级、独立事件存档或事务性 outbox,不能假装数据完整。
对于搜索投影,可与 MongoDB $rankFusion 混合检索指南组合:Change Streams 负责发现权威数据变化,索引更新必须仍使用稳定文档键与可对账版本,检索相关性逻辑不应被塞进消费者的 checkpoint 事务。
遇到 invalidate event 怎么办?
collection drop、rename 等操作可能产生 invalidate 并关闭当前 cursor。普通 resumeAfter 不一定能跨过这个边界;MongoDB 提供 startAfter 处理从 invalidate token 之后开启新流的场景。
但“技术上能继续”不代表“业务上应该继续”。如果监听的 collection 被 rename,namespace、索引、schema 和消费者假设都可能变化。建议把 invalidate 视为控制面事件:停止副作用,记录 token 与 DDL 变更,验证新 namespace 与 schema,决定使用 startAfter、建立新 stream,还是执行 backfill。
不要让一个无限 retry loop 吞掉 invalidate。重试应该区分短暂网络故障、可恢复 cursor 错误、token 历史丢失、认证错误和 DDL 造成的终止。
PyMongo 消费循环应该怎样退出和恢复?
PyMongo Change Streams 文档列出了 pipeline、full_document、resume_after、start_after、start_at_operation_time 与 max_await_time_ms 等选项。阻塞式迭代适合专用 worker;应用关闭时应结束读取、等待正在处理的事件完成,并在 checkpoint 成功后关闭 cursor。
建议将读取与处理解耦,但队列必须有上限。无限内存队列会在下游变慢时把延迟问题变成 OOM。达到 high-water mark 时暂停读取或降低并发,同时持续暴露 lag。每个事件的处理超时、重试次数和 dead-letter 规则都应可配置。
示例状态机可以是:
READ -> VALIDATE -> APPLY -> CHECKPOINT -> ACK_LOCAL
| |
v v
RETRY RECONCILE
ACK_LOCAL 只是消费者内部完成,不代表搜索、邮件或第三方系统已经达到最终一致。报告中应分别描述 Mongo 事件、业务副作用、checkpoint 和对账状态。
上线前如何验收?
功能测试
- 插入、更新、替换、删除各触发预期处理。
- update delta、
fullDocument与 delete 的字段缺失按设计处理。 - pipeline 不修改事件
_id。 - 两个不同 collection 的相同文档
_id不会冲突。
故障测试
- 在副作用完成、checkpoint 前强制终止进程,验证重复处理无害。
- 在副作用前终止,验证事件会重放。
- 断开网络并恢复,确认驱动与应用的重试边界。
- 使用无效或过期 token,确认系统停止并进入 backfill/reconcile,而不是跳到现在。
- 在隔离环境触发 rename/drop,验证 invalidate 路径。
可观测性测试
为每个 batch 或事件记录结构化字段:stream ID、operation type、namespace、处理耗时、结果与匿名化 operation ID。不要打印完整业务文档或凭据。可以使用 FastAPI OpenTelemetry 链路追踪指南中的 trace 传播原则,把入库请求与派生更新关联起来,但高基数 token 不应直接成为 metric label。
常见错误清单
- 处理前推进 checkpoint,导致失败事件被跳过。
- 处理后推进 checkpoint,却没有下游幂等,导致重复副作用。
- 只保存
clusterTime或 token 的一部分。 - 修改 pipeline 后继续复用旧 token。
- 把
updateLookup当作事件时刻的历史快照。 - token 失效后静默从当前时间开始。
- 多实例共享 checkpoint,却没有 lease 与 CAS。
- 使用无限队列吸收下游背压。
- 把 cursor 自动恢复误认为端到端 exactly-once。
FAQ
Change Streams 能保证 exactly-once 吗?
它提供可恢复的数据库变更流,不会替你把 Mongo、checkpoint 与任意外部系统变成一个事务。生产消费者通常按 at-least-once 设计,通过下游幂等、唯一约束、业务版本和对账获得可接受的一致性。
每个事件都写 checkpoint 会不会太慢?
可以按小批次推进,但崩溃时会重放整个未 checkpoint 批次,因此批次大小是写放大、恢复时间和重复处理之间的权衡。对不可重复副作用,先建立可靠幂等,再优化 checkpoint 频率。
可以只依赖驱动自动恢复吗?
不可以。驱动能处理部分 resumable error,但不了解你的副作用是否完成,也无法处理 oplog 历史丢失、配置变化或业务对账。应用仍需持久 checkpoint 和显式恢复状态机。
Change Streams 与 outbox 怎么选?
如果只需把 Mongo 的权威状态投影到可重建索引或缓存,Change Streams 很合适。若每个业务事件必须有不可丢失、可审计的语义,并与业务写入同事务提交,事务性 outbox 往往更清晰;也可以用 Change Streams 消费 outbox collection。
可靠的 Change Streams 系统不是一个 for event in stream 循环,而是一条有 checkpoint、幂等、背压、恢复和对账的复制管道。只要把这些状态显式化,网络抖动和进程重启就不会被误报成数据一致性保证。