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 已推进”拆成不同状态。

resumeAfterstartAfterstartAtOperationTime 如何选择?

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 文档列出了 pipelinefull_documentresume_afterstart_afterstart_at_operation_timemax_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 和对账状态。

上线前如何验收?

功能测试

  1. 插入、更新、替换、删除各触发预期处理。
  2. update delta、fullDocument 与 delete 的字段缺失按设计处理。
  3. pipeline 不修改事件 _id
  4. 两个不同 collection 的相同文档 _id 不会冲突。

故障测试

  1. 在副作用完成、checkpoint 前强制终止进程,验证重复处理无害。
  2. 在副作用前终止,验证事件会重放。
  3. 断开网络并恢复,确认驱动与应用的重试边界。
  4. 使用无效或过期 token,确认系统停止并进入 backfill/reconcile,而不是跳到现在。
  5. 在隔离环境触发 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、幂等、背压、恢复和对账的复制管道。只要把这些状态显式化,网络抖动和进程重启就不会被误报成数据一致性保证。