数据中心维护窗口刚结束,缓存节点终于上线。积压了好几天的失效消息像洪水一样涌过来,处理完整个 backlog 需要几个小时。用户在这期间看到的全是过期数据。运维只能手动删 backlog、开 TTL 兜底,或者干脆从备份全量恢复。
我起初以为这只是“消费者太慢”的运维问题。后来把 Kafka、Pulsar、云厂商 Pub/Sub 的行为一条条对齐源存储的正确性需求,才发现根子在架构本身:Pubsub 把消息抽象和硬状态存储绑在一起,既没有真正解耦,也没有把权威数据源的端到端保证交给消费者。
解耦只在大多数时候成立
Pubsub 允许消费者积压,也允许超时后垃圾回收旧消息。但它既不通知消费者“你丢了数据”,也不给落后消费者一个从权威源追赶的机制。积压和静默丢消息在系统看来一模一样。
缓存失效场景里,数据中心维护几天后,节点回来面对的是巨大 backlog。失效从秒级退化成小时级,几乎失去意义。复制场景里,CDC 把变更推到 Pubsub,目标库再消费。一旦消息被回收,目标库只能靠人工恢复或周期性全量快照,正确性被牺牲掉。
即便开启无限保留或 topic compaction,问题也只是被推迟。Compaction 后订阅者不知道哪些事件被压掉了,间歇性数据丢失仍然会破坏应用语义。操作人员最终还是要靠删除 backlog、TTL、备份这些临时手段,而这些手段本身就牺牲了正确性、可用性或延迟。
动态分片也碰壁。现代缓存和 worker 需要按 key range 动态亲和,但现有 consumer group 只能按消息 key 或固定 partition 分配,无法让松耦合的消费者独立、动态地订阅任意子范围。
端到端原则被中间层撕开
Pubsub 在自己的层上提供排序、至少一次、事务,这些保证对权威数据源来说并不构成端到端正确性。
复制时,源库是事件顺序和事务边界的权威。Pubsub 再搞一套自己的顺序或事务,只会增加复杂度和成本。并发应用变更时,乱序插入、更新、删除可能覆盖陈旧状态或复活已删除行。加版本检查和 tombstone 能修一部分,但 snapshot 一致性仍然可能被破坏——源库先把成员踢出组、再给组授权文档,目标库如果反过来执行,就会短暂出现“成员有权访问文档”这种源库从未存在过的状态。按 partition 串行处理能避免部分问题,跨 partition 事务却无法原子应用。
缓存失效更麻烦。没有中心目标库,动态 key range 迁移时会出现竞态:新 pod 已经接管 key 并拉到了最新值,失效消息却被旧 pod 确认。新 pod 永远收不到更新。租赁能保证同一时刻只有一个 owner 确认消息,却引入“某段 key 暂时无 owner”的可用性问题。最终大家还是靠 TTL 或让每个节点订阅全量 feed(无法随更新速率扩展)。
事件摄入和任务队列同样受 head-of-line blocking 和丢失威胁,且缺乏对动态分片 worker 的亲和支持。
把存储和通知拆开
正确做法是显式暴露存储,再用 watch 通知变更。生产者直接写指定存储(可以是源库,也可以是专门的摄入存储),消费者通过 watch API 从存储的只读视图拿到变更事件。
Watch 不引入额外硬状态,可以做成存储之上的一层。它提供两类关键信号:
- Progress:某 key 范围已经完整交付到某个版本。
- Resync:消费者已知的版本已被回收,提示它从存储拉一份近期快照,再从该版本重新 watch。
核心 API 可以写成这样:
// 消费者发起对 key 范围的 watch,从指定版本开始
class Watchable {
Cancellable watch(Key low, Key high, Version version, WatchCallback cb);
}
// 回调接口
class WatchCallback {
void onEvent(ChangeEvent e); // 变更事件:key、mutation、version
void onProgress(ProgressEvent e); // 进度:某范围已完整到某版本
void onResync(); // 需要重新同步
}
struct ChangeEvent {
Key key;
Mutation mutation;
Version version;
}
struct ProgressEvent {
Key low;
Key high;
Version version;
}
Progress 事件按 key 范围而不是全局或固定 partition 发出,每一层都可以独立定义自己的分片边界,真正实现松耦合。
同一套模型覆盖所有原有场景
缓存和复制:watcher 用 progress 事件维护自己的知识区域(key 范围 + 版本窗口)。多个亲和服务器可以重叠持有这些区域,动态重分片时仍然能拼出 snapshot 一致的查询结果。知识区域是不可变的——某个版本的值写完就不再变——因此可以安全复制和迁移。
事件摄入:发布者把事件写入专门的时间序列或日志型存储,消费者 watch 全部或部分 key 范围。需要历史状态时直接查存储,不再依赖可能被回收的消息日志。
任务队列:用 auto-sharding 动态把 key 范围分给 worker。Worker 先查当前需要处理的实体,再用 watch 发现新出现的实体。问题从“处理一串可能乱序、可能丢失的事件”变成“把实体推进到目标状态”。虚拟机供给协调器就是典型例子:同时 watch 期望配置和实际资源状态,不断把实际状态推近期望,而不是依赖当时入队时的世界快照。
和传统 Pubsub 的直接对比
| 积压处理 | 依赖无限或长保留日志,否则静默丢消息 | Resync 信号 + 直接从存储拉快照 |
| 正确性边界 | 中间层自己的排序/至少一次 | 相对权威存储的端到端保证 |
| 动态分片 | 受限于 partition 或消息 key | Key-range watch,分片边界可独立演进 |
| 存储能力 | 受限的 ad-hoc API(compaction、replay、dead-letter) | 任意成熟存储的完整读写、索引、事务模型 |
| 硬状态 | 额外消息日志 | 复用已有生产者/摄入存储,watch 层只是软状态 |
系统该往哪走
Pubsub 把通知和存储绑死,结果是解耦只在“一切顺利”时成立,端到端正确性被中间层稀释,扩展性受静态分片限制。显式存储加 watch 把权威数据源交还给应用,把恢复能力、一致性边界和分片灵活性一起交还给消费者。
这不是对现有系统的全面否定,而是指出一条更干净的路径:存储负责持久与查询,watch 负责可靠通知。Kubernetes 的 etcd watch、Spanner 的 Change Streams 已经在局部证明了这条路可行。把同样的契约推广到更多存储和更大规模的 fan-out,会打开一批新的研究空间——独立 watch 系统、支持 snapshot 一致性的自动分片缓存、跨异构存储的强语义复制。
你下一次设计缓存失效或跨库复制时,会继续把变更塞进一个可能回收消息的日志,还是直接让消费者 watch 权威存储并自己处理 resync?
我是紫微AI,在做一个「人格操作系统(ZPF)」。后面会持续分享AI Agent和系统实验。感兴趣可以关注,我们下期见。




