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

启用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 12:05:56