Postgres CDC 变天:Snowflake 把数据同步做成了瑞士钟表
Postgres CDC 变天:Snowflake 把数据同步做成了瑞士钟表
你凌晨三点被 PagerDuty 叫醒过吗?
Postgres → 数据仓库的同步挂了,WAL 堆积把磁盘撑爆,Debezium 消费滞后 30 分钟还没恢复。你一边手动 truncate WAL 一边祈祷别丢数据。
如果你经历过这个凌晨三点,Snowflake 这篇工程博文你应该读。
他们没再修修补补传统 CDC 方案,而是把整个复制逻辑推倒重来——把 CDC 塞进 Postgres 自己体内,让它直接把变更推到数据湖里。
传统 CDC 的通病:拉取模式注定脆弱
Postgres 的 CDC 不是不能用,只是那个"能用"要拿人力去换。
标准做法是走 logical decoding——Postgres 把 WAL 解码成行级 insert/update/delete 操作,通过 replication slot 暴露给外部消费者。剩下的 pipeline 你自己搭:Debezium → Kafka → Sink Connector → 目标库。
问题出在最后这一跳:消费者完全不知道 Postgres 内部发生了什么。
Schema 变了?消费者不知道,处理到一半炸了。表被 drop 了?消费者还在那等变更,等到 WAL 溢出。Postgres 挂了还是网络断了?消费者分不清,要么重复消费要么丢数据。
再加上 snapshot 和增量变更的对齐问题、故障恢复时的 exactly-once 语义、大表 upsert 的性能退化……任何一个环节出问题,凌晨三点的电话就是你的。
Snowflake 的答案是:别拉了,让 Postgres 自己推。
推翻重来:把 CDC 写进 Postgres 扩展
Snowflake 写了一个叫 snowflake_cdc 的 Postgres 扩展,跑在数据库进程内部。
它干的事不复杂:持续把变更批量化,推到 S3 上的 Iceberg 表里(Parquet 格式),Snowflake 那边自动 apply。
推送到对象存储、让生产者和消费者解耦,这套思路本身不新。新的是位置——扩展跑在 Postgres 进程里,schema 变更、事务边界、表的生老病死它都第一手知道,不用外面的消费者去猜。
流水线分四段:
- Write:数据写入表,WAL 追加记录
- Decode:后台 worker 解码"过去"的 WAL 为行级变更
- Capture:变更打包成 batch,写入 Iceberg changelog
- Apply:Snowflake 端按事务边界批量合并
四段各自向后错开一个时间窗口,互不阻塞。解码走的是 Postgres 的 historic snapshot 机制——表后来被改过、被删了都没关系,解码器看到的仍是写入那一刻的 catalog。
把事务带到跨系统复制里
整篇博文里,这一节最戳我。
做过 ETL 或 CDC 的都懂:事务在单库里好用得像作弊,一跨系统就没了。剩下的幂等、去重、乱序、部分失败全得自己手写,而且每一样都挑你最不想加班的那天出问题。
Snowflake 的做法是在两端都用事务包住:
- Postgres 端:一次事务把多个 batch 写入多个 Iceberg changelog,要么全写,要么一条都不写
- Snowflake 端:一次事务合并多个 batch 到目标表,精确停在 Postgres 事务边界上
换来的是:外键完整、join 结果对得上、不需要 upsert。
说到 upsert——不少 CDC 方案图"简单",把所有操作一律转成 upsert。省事在前,账在后面:
- 表中间状态不一致(新表刚创建时的 snapshot 对齐问题)
- Insert 变 upsert 意味着每次插入都要和列存的目标表做 match,越跑越慢
- 合并历史数据和增量变更几乎是不可完成的任务
Snowflake 不用 upsert,因为它的复制是事务驱动的——insert 就是 append,delete 就是 delete,不会重复执行。结果:insert-heavy 的大表(通常也是最大的表)复制极快,几乎零开销。
Live Views:不急着 apply 也能低延迟
这一段我看了两遍才反应过来。
传统 CDC 想让延迟低,就得拼命 apply;apply 越频繁,合并开销越大,目标表查得越慢。两头堵。
Snowflake 的解法是 live views:查询时把 changelog 里还没 apply 的变更和目标表数据现场合并,过滤和投影直接下推到 Parquet 文件层。
于是 apply 就不必频繁了。间隔拉长,查询延迟照样在一分钟以内,代价几乎看不见——"apply 频率"和"查询延迟"这两个旋钮被拆开了。
对你的意义
大部分团队没有 Snowflake 的工程资源,也犯不上从零造 CDC。但里面有三条取向可以直接搬:
- Push > Pull。让源库主动推送变更,比外部消费 WAL 稳定一个量级。如果你的同步任务不稳定,先不要加机器,看看能不能把逻辑往源端挪。
- 事务边界是 CDC 的基石。用事务包住跨系统操作,能消灭一大类故障。即使是简单的"写个 checkpoint 文件 + 合并变更"也比无状态的 upsert 靠谱。
- Apply 频率和查询延迟可以解耦。live views 的思路用物化视图也能近似实现——不定时强制合并,查询时 UNION 增量数据。
今天你可以做的一件事:
看看你现在的 CDC pipeline,同步延迟超过 5 分钟的话,先查 WAL retention 策略,不要急着加机器。如果延迟波动大,考虑在源端加一个推送层,哪怕是一个简单的 pgoutput → S3 的脚本,也比 Kafka 中间那一层少一个故障点。
✨ 本文由 DeepSeek 生成初稿,Claude 审核润色。
参考来源: