You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Snowflake Iceberg表疑似存在数据不可变性与CDC异常问题

Snowflake Iceberg表CDC与数据不可变性问题分析

问题1:UPDATE操作在CDC中显示为INSERT

以下是复现步骤及结果:

CREATE ICEBERG TABLE T1 (i INT ) EXTERNAL_VOLUME = 'exvol' CATALOG = 'SNOWFLAKE' BASE_LOCATION = 'S/T1';

CREATE STREAM T1_STREAM ON TABLE T1;

BEGIN TRANSACTION;
INSERT INTO S.T1 VALUES (1);
COMMIT;
SELECT * FROM T1_STREAM;

执行结果:

+---+-----------------+-------------------+------------------------------------------+
| I | METADATA$ACTION | METADATA$ISUPDATE | METADATA$ROW_ID                          |
|---+-----------------+-------------------+------------------------------------------|
| 1 | INSERT          | False             | 43825f04d34552c4a8524f87c5e328fa70f2cd0d |
+---+-----------------+-------------------+------------------------------------------+
BEGIN TRANSACTION;
UPDATE T1 SET I=2 WHERE I=1;
COMMIT;
SELECT * FROM T1_STREAM;

执行结果:

+---+-----------------+-------------------+------------------------------------------+
| I | METADATA$ACTION | METADATA$ISUPDATE | METADATA$ROW_ID                          |
|---+-----------------+-------------------+------------------------------------------|
| 2 | INSERT          | False             | 43825f04d34552c4a8524f87c5e328fa70f2cd0d |
+---+-----------------+-------------------+------------------------------------------+
BEGIN TRANSACTION;
UPDATE T1 SET I=3 WHERE I=2;
COMMIT;
SELECT * FROM T1_STREAM;

执行结果:

+---+-----------------+-------------------+------------------------------------------+
| I | METADATA$ACTION | METADATA$ISUPDATE | METADATA$ROW_ID                          |
|---+-----------------+-------------------+------------------------------------------|
| 3 | INSERT          | False             | 43825f04d34552c4a8524f87c5e328fa70f2cd0d |
+---+-----------------+-------------------+------------------------------------------+

问题2:Parquet文件被原地重写,违反数据不可变性

SELECT file_name, file_size, row_count, row_group_count, etag, md5 FROM
TABLE(INFORMATION_SCHEMA.ICEBERG_TABLE_FILES(TABLE_NAME => 'T1'));

执行结果:

+-----------------------------------------------------------------+-----------+-----------+-----------------+------+----------------------------------+
| FILE_NAME                                                       | FILE_SIZE | ROW_COUNT | ROW_GROUP_COUNT | ETAG | MD5                              |
|-----------------------------------------------------------------+-----------+-----------+-----------------+------+----------------------------------|
| S/T1.eB1BwjT4/data/snow_QLd-Nw_mD1Q_gAD3iV1Qbxg_0_1_002.parquet |      1536 |         1 |               1 |      | 114ce57400f3251920df4121c3f0f5f3 |
+-----------------------------------------------------------------+-----------+-----------+-----------------+------+----------------------------------+
BEGIN TRANSACTION;
UPDATE T1 SET I=4 WHERE I=3;
COMMIT;
SELECT * FROM T1_STREAM;

UPDATE T1 SET I=5 WHERE I=4;

执行结果:

+------------------------+-------------------------------------+
| number of rows updated | number of multi-joined rows updated |
|------------------------+-------------------------------------|
|                      1 |                                   0 |
+------------------------+-------------------------------------+
SELECT * FROM T1_STREAM;

执行结果:

+---+-----------------+-------------------+------------------------------------------+
| I | METADATA$ACTION | METADATA$ISUPDATE | METADATA$ROW_ID                          |
|---+-----------------+-------------------+------------------------------------------|
| 6 | INSERT          | False             | 43825f04d34552c4a8524f87c5e328fa70f2cd0d |
+---+-----------------+-------------------+------------------------------------------+
SELECT file_name, file_size, row_count, row_group_count, etag, md5 FROM
TABLE(INFORMATION_SCHEMA.ICEBERG_TABLE_FILES(TABLE_NAME => 'T1'));

执行结果:

+-----------------------------------------------------------------+-----------+-----------+-----------------+------+----------------------------------+
| FILE_NAME                                                       | FILE_SIZE | ROW_COUNT | ROW_GROUP_COUNT | ETAG | MD5                              |
|-----------------------------------------------------------------+-----------+-----------+-----------------+------+----------------------------------|
| S/T1.eB1BwjT4/data/snow_QLd-Nw_mD1Q_wN4R6odQbxg_0_1_002.parquet |      1536 |         1 |               1 |      | 9ce3fd2cba20a0493dbe252b5bfb0818 |
+-----------------------------------------------------------------+-----------+-----------+-----------------+------+----------------------------------+

问题原因分析

  1. CDC事件失真原因:
    Snowflake的Iceberg表流捕获依赖于快照的文件变更日志。对于仅包含单条记录的极小Parquet文件,执行UPDATE时,Snowflake没有遵循Apache Iceberg标准的"标记旧文件为删除+添加新文件"的快照变更模式,而是直接原地重写原文件。这种操作不会生成CDC所需的"DELETE+INSERT"成对事件,导致流只能捕获到最终的INSERT动作,METADATA$ISUPDATE始终为False。

  2. 文件重写原因:
    这是Snowflake针对Iceberg表的小文件优化逻辑。当文件大小远低于默认阈值(通常为16MB)时,Snowflake会触发原地重写来合并小文件、降低元数据管理开销。但该行为违反了Apache Iceberg核心的"数据文件不可变"原则——Iceberg要求所有数据变更都通过新增快照、调整文件引用实现,不允许修改已存在的数据文件。

可能忽略的关键点

  • Snowflake的Iceberg实现并非完全对齐Apache Iceberg原生规范,存在特有的性能优化逻辑,这类细节在官方文档中可能未被重点强调。
  • 单条记录的测试场景会触发极端的小文件优化,生产环境中当文件达到一定大小后,UPDATE将触发标准的Iceberg变更流程,不会出现此类问题。
  • Iceberg表的流CDC捕获依赖于快照的文件级变更,原地重写会绕过这一机制,导致CDC事件无法准确反映实际操作类型。

验证建议

  • 插入大量数据生成超过16MB的Parquet文件,再执行UPDATE操作,观察CDC是否能正确识别UPDATE事件,以及文件是否被重写。
  • 检查Snowflake Iceberg表的相关配置参数(如小文件合并阈值),确认是否有参数可以禁用或调整该优化行为。

内容的提问来源于stack exchange,提问作者Sumeet Keswani

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 04:44:51