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

