将事务数据库的数据实时同步到分析平台,是现代数据架构中最基础也最头疼的问题之一。传统方案依赖外部 CDC 工具,架构脆弱、运维复杂、成本高企。Snowflake 工程师 Marco Slot 日前撰文详解了该公司如何从根本上重新设计了 Postgres 复制机制——通过一个数据库扩展,把变更数据捕获直接推入对象存储。

从「拉」到「推」:CDC 的范式转变
Postgres 原生的逻辑解码(logical decoding)只负责把 WAL 日志翻译成行级变更流,剩下的重担全部甩给客户端:回填、Schema 变更处理、表创建与删除、快照对齐、故障重启……这些步骤中的任何一个出问题,整条管道就会断裂。
Snowflake 的方案名为 Data Mirroring(数据镜像),核心理念极其简单:不再让外部系统来「拉」数据,而是从 Postgres 内部主动「推」。具体实现是一个名为 snowflake_cdc 的 Postgres 扩展,它在后台持续将变更批量写入每张表的变更日志(change log)和一个元日志(meta log),底层存储格式为 Apache Iceberg(压缩 Parquet 文件)。
为什么选对象存储?因为 Amazon S3 这类服务本身具有高可扩展性和高可靠性,Postgres 备份早已在使用。作为 CDC 数据的目的地,它比任何自定义中间件都更稳。
四阶段时间线解耦
Data Mirroring 将一次写入拆解为四个连续但解耦的阶段:
1. 写入(Write):数据修改表并产生 WAL 记录
2. 解码(Decode):历史 WAL 被翻译成行级变更,利用 Postgres 的「历史快照」功能读取写入时刻的目录表状态
3. 捕获(Capture):行级变更被分批收集,定期追加到 Iceberg 变更日志
4. 应用(Apply):Snowflake 端以有限状态机方式执行元日志指令,将变更批次合并到目标表
这种设计的精妙之处在于:即使某张表在解码时已被修改甚至删除,历史快照仍能保证 WAL 记录被正确理解。Schema 变更走的是同一条 Write→Decode→Capture 路径,自然地嵌入变更流中,不会出现时序错乱。
交易边界:用分布式系统的最简原语解决最难的问题
数据库系统用一个简单原语隐藏了海量复杂性:事务(transaction)。事务失败就重试,成功就保证恰好执行一次——这正是 ETL/CDC 最缺乏的保证。
Snowflake 此前发布的 pg_lake(已正式可用)让 Postgres 能够跨 Postgres 表和 Iceberg 表执行事务。Data Mirroring 在此基础上更进一步:Postgres 端在一个事务内把数据和 Schema 变更批量推入多个 Iceberg 变更日志,Snowflake 端在一个事务内合并多个批次。这意味着所有 Snowflake 表精确推进到某个 Postgres 事务边界,外键和联接的正确性得到保证。
不做 Upsert:插入密集型场景的性能飞跃
传统复制方案普遍采用 upsert(插入或更新)策略,原因是很难解决快照与变更的一致性问题。但 upsert 在列式存储上有严重缺陷:每次插入都要匹配目标表已有行,代价极高;且无法高效地将近期变更与已有数据合并。
由于 Data Mirroring 是受控的事务性过程,它可以生成完美的删除流和插入流,精确执行一次,不存在重复应用的风险。结果是:插入密集型工作负载(通常是最大的表)的复制速度极快且成本低廉,因为插入只是追加,永远不做 upsert。
Live Views:低频应用也能获得低延迟
基于上述能力,Snowflake 推出了 Live Views(实时视图)功能:将未应用的变更日志与目标表数据结合,查询时直接将过滤条件下推到存储层,同时扫描 Parquet 变更日志和基表。即使变更应用频率不高,Live Views 的延迟仍能控制在 1 分钟以内,且查询性能几乎不受影响。
Slot 总结道:「你设置一次,它永远运行。」没有外部连接器会落后,没有快照冲突,没有随表增长而减速的 upsert——只有一个 Postgres 扩展往对象存储推批次,一个 Snowflake 端事务性地应用它们,两边独立运行,永不停歇。
来源:Snowflake 工程博客 / Marco Slot



