启用upsert的Iceberg表无法触发Flink SQL流快照增量读取
问题:Iceberg启用upsert后Flink流Join停止增量读取新快照
现象
- 基于Iceberg v2表构建Flink流作业,通过Join多源表写入目标Iceberg表,源表未启用
write.upsert.enabled时,流订阅正常,可持续读取新快照; - 只要任意一个源表启用
'write.upsert.enabled'='true',流Join作业仅读取一次初始数据后就停止响应新快照,作业不终止但不再产生新输出。
环境配置
- Flink版本:1.17.2
- Iceberg版本:1.4.2
- 正常工作的流查询SQL:
INSERT INTO iceberg.target_packaging SELECT usr.`user_id` AS `user_id`, usr.`adress` AS `address`, ord.`item_id` AS `item_id`, .... FROM iceberg.source_users /*+ OPTIONS('streaming'='true', 'monitor-interval'='15s') */ usr JOIN iceberg.source_orders /*+ OPTIONS('streaming'='true', 'monitor-interval'='15s') */ ord ON usr.`user_id` = ord.`user_id`;
源表配置对比:
- 正常配置:
CREATE TABLE iceberg.source_users ( `user_id` STRING, `adress` STRING, .... PRIMARY KEY (`user_id`) NOT ENFORCED ) with ('format-version'='2');表属性:
[current-snapshot-id=7980858807056176990,format=iceberg/parquet,format-version=2,identifier-fields=[user_id],write.parquet.compression-codec=zstd]- 启用upsert的异常配置:
CREATE TABLE iceberg.source_users ( `user_id` STRING, `adress` STRING, .... PRIMARY KEY (`user_id`) NOT ENFORCED ) with ('format-version'='2', 'write.upsert.enabled'='true');表属性:
[current-snapshot-id=3566387524956156231,format=iceberg/parquet,format-version=2,identifier-fields=[user_id],write.parquet.compression-codec=zstd,write.upsert.enabled=true]
已确认事项
- Flink作业已启用Checkpoint;
- 流查询已配置
streaming=true和monitor-interval=15s; - 源表和目标表均已预先正确定义。
排查与解决建议
- 检查Iceberg快照生成情况:用Iceberg CLI执行
iceberg snapshot list <table-name>,确认新写入数据后是否生成了新快照,以及快照是否为增量类型。启用upsert后Iceberg会生成带位置删除的快照,需确保Flink能识别这类快照变更。 - 显式配置Flink读取分片参数:在读取选项中添加
read.split.target-size和read.split.metadata-target-size,帮助Flink正确扫描新快照文件,示例:iceberg.source_users /*+ OPTIONS('streaming'='true', 'monitor-interval'='15s', 'read.split.target-size'='134217728', 'read.split.metadata-target-size'='16777216') */ usr - 升级Iceberg版本:Iceberg 1.5.0+对Flink流读取upsert表的逻辑有优化,可尝试升级版本验证是否解决兼容性问题。
- 排查作业日志:查看Flink JobManager和TaskManager日志,搜索
SnapshotFetcher、IncrementalSource相关条目,确认是否存在快照扫描失败、文件读取异常的报错。 - 验证时间属性配置:如果是事件时间Join,确保源表数据携带正确的事件时间戳,Watermark配置合理,避免upsert导致的Watermark推进异常影响Join触发。
内容的提问来源于stack exchange,提问作者Shaflump
相关产品推荐
相关产品推荐

