
做数据平台这些年我观察到一个很有意思的现象一聊数据湖大家满脑子都是 HDFS、S3一聊实时第一反应就是 Kafka、Flink。但真正把“实时”和“湖”这两个字接起来的往往是被忽略的那一层表格式。Apache Flink Apache Iceberg 这对组合表面看是计算引擎和存储格式搭伙干活实际上它们靠着一套非常深的协作机制才支撑起今天大家常说的流批一体、数据湖实时入湖、湖仓一体这些架构。这篇文章我想基于自己实际用下来的经验把这层协作关系一次讲透包括底层机制、写入提交的原理、流式读取的处理方式、生产环境里的调优参数以及几个很容易踩的坑。适合正在做实时数仓、数据湖平台或者准备把离线批处理和实时链路统一到一套架构上的同学参考。1. 为什么偏偏是 Flink 和 Iceberg 走到了一起1.1 先搞清楚这俩各管什么Flink 是计算引擎负责的是“怎么算”它接收数据流做清洗、聚合、关联然后把结果写出去。Iceberg 是表格式负责的是“怎么存”它定义了一张表在分布式存储上的布局、元数据组织方式、并发读写规则、事务语义。一句话Flink 干活Iceberg 记账。很多人把这个关系理解成“Flink 往 Iceberg 表里写文件”这么想也不算错但会把二者协作的深度低估了。Hive 时代也有一套表和目录的约定但那张“表”基本就是“Metastore 里存 schema HDFS 目录里放文件”没有任何事务能力。两个任务同时写一个 Hive 分区数据是能互相污染的上游任务写到一半挂了下游已经能看到半截文件表结构改个字段物理路径和文件布局全得跟着大动。Iceberg 的出现就是把这些痛点全部用一个中间层接住了。1.2 Iceberg 到底解决了什么Iceberg 最大的三个能力我总结成三个词快照隔离、模式演进、分区管理。快照隔离意味着每张表的每次写入都会生成一个不可变的快照读任务可以选择任意一个快照读写任务也永远不去修改已经存在的数据文件。这直接让 Flink 的实时写入和下游的批量读取不打架了。以前写 Hive 表时最怕的就是跑批任务刚好碰上数据文件被覆盖现在 Iceberg 的快照机制天然规避了这个冲突。模式演进更好理解。你要给一张几 TB 的表加一列在传统 Hive 里基本意味着重写一遍所有数据文件或者用各种取巧的补列方式。Iceberg 的 schema 是独立于文件的元数据加列、减列、改列顺序改的就是元数据里的一个结构定义数据文件不需要动。Flink 上游加了字段下游表结构直接兼容。分区管理这块是最容易被低估的。传统分区表的分区规则写死之后基本不能改想从按天改成按小时得把历史数据全部重刷。Iceberg 的隐藏分区和分区演化机制解决了这个分区规则变了新数据按新规则落盘老文件还在老地方scan 的时候元数据层会把新旧文件统一组织起来。后面我会详细说这个机制。1.3 Flink 在这套组合里的不可替代性Iceberg 不是只能配 Flink它也能配 Spark、Trino。但 Flink 有一个别人比不了的优势它天然是流批一体的引擎批处理和流处理用的是同一套 API、同一套 SQL 语义更重要的是同一套容错机制。Kafka 里的实时数据Flink 消费之后可以直接以很低的延迟写进 Iceberg这个链路不需要任何中间环节。如果换成 Spark要么走 Structured Streaming 做个微批要么先落 Kafka 再离线导一次。所以你会发现在实时写入数据湖这个场景里Flink 和 Iceberg 几乎是一对没有替代选择的组合。另外一个隐藏优势是 Flink 的 checkpoint 机制和 Iceberg 的快照提交能形成天然的配合。Flink 的 checkpoint 保证状态一致Iceberg 的 commit 保证存储原子可见二者协作的结果是端到端的 exactly-once 写入语义。这个配合关系不是简单地调 API而是两个系统在容错模型上的深度对齐。2. 先把 Iceberg 的底层机制看清楚后面协作才不难理解2.1 一张表到底由哪些文件组成理解 Flink 和 Iceberg 怎么协作必须先知道 Iceberg 表在存储上是四层结构。我见过不少同学调参数全靠猜就是因为不知道每个参数影响的是哪一层文件。层次文件类型作用元数据文件metadata.json记录当前表元数据版本、schema、快照列表等快照清单列表snap-*.avro一个快照下所有 manifest 文件的集合清单文件manifest 文件记录一批数据文件的统计信息、分区信息、列统计数据文件Parquet/ORC/Avro真正的业务数据每次对表做一次写操作Iceberg 不会改动任何已有文件而是写出一批新的数据文件生成新的 manifest再通过一次原子操作把表的当前元数据指针切换到新的快照上。这个“切换指针”的动作就是 commit。理解这层结构后你就能明白为什么 Iceberg 能做时间旅行为什么并发读写互不阻塞也就能明白 Flink 的流式读机制到底在做什么——它其实就是在不断发现“指针又切到新快照了”。2.2 Manifest 文件的元数据过滤能力Manifest 文件里存了每个数据文件的最小值、最大值、空值数量等统计信息。查询的时候Iceberg 的 scan 引擎会先在 manifest 层做过滤把一个几万文件的表直接裁剪到只需要读几个文件。这不是靠分区目录路径去猜的而是靠列统计做的精确过滤。这对 Flink 的读取也很有意义。Flink 作为计算引擎并不直接感知 Iceberg 的文件布局它拿到的是一组“需要读取的文件列表”。Iceberg 的扫描规划器在生成这批文件列表时已经帮你做了大量裁剪所以 Flink 侧任务实际上只扫描必要文件整个任务的 I/O 压力大幅降低。2.3 隐藏分区和分区演化这俩经常被搞混传统 Hive 分区表你要在表里单独维护一个分区字段。Iceberg 的隐藏分区不需要你手动维护分区列而是通过一个 transform 自动从某个普通字段推导出来。比如你有 ts 字段可以在建表时声明PARTITIONED BY (days(ts))Iceberg 会为每一整天的数据生成一个隐式分区值分区目录的命名和物理布局完全由 Iceberg 自己管理。你不需要像以前那样自己拼一个dt2024-01-01的字段在数据里。分区演化就更关键了。假设你一开始是按天分区后来业务上想改成按小时分区传统方案需要全量重刷数据。Iceberg 不需要它会在元数据里记录两套分区方案老数据按旧方案组织新数据按新方案组织。查询的时候Iceberg 负责把两种布局的数据文件合并成一个统一的结果集。在 Flink 的写入场景里这意味着你可以随时调整表的分区规则完全不影响正在运行的写入作业。我在生产里调整过一次分区粒度流任务零重启只需要重新提交一下 DDL这个体验是 Hive 给不了的。3. Flink 写 Iceberg提交机制里藏着的协作智慧3.1 一条数据从 Flink 到 Iceberg 快照要走几步很多人以为 Flink 写 Iceberg 就是“边读边写文件”其实远没这么简单。Flink 写 Iceberg 的完整过程是分阶段的。Flink 的每个写入算子先把数据攒在本地临时文件里这个阶段的数据外面是看不到的。等触发一次 checkpointcheckpoint barrier 对齐后Flink 完成自己的状态快照然后 Iceberg 的提交器会拿到这一批数据文件的清单调用 Iceberg 的 commit API把文件提交为新的快照。这里有个非常关键的点commit 是在 checkpoint 完成之后才发生的。所以 Iceberg 表的快照生成频率就等于 Flink checkpoint 的频率。你有多少个 checkpoint这张表就有多少个快照。这就带出两个直接的调优结论第一checkpoint 间隔不要太短。我看到有人为了“实时性”把 checkpoint 设成 10 秒一次结果一天 8640 个快照元数据文件刷得飞快NameNode 都快报警了。对数据湖这种场景checkpoint 间隔在 1 到 5 分钟是比较合理的范围入库延迟控制在分钟级完全够用。第二小文件和快照频率强相关。每次 checkpoint 提交一批文件提交频率越高单个文件越小。控制小文件不能只靠调文件大小参数关键还是先控住 checkpoint 频率。3.2 为什么 Upsert 模式要额外生成 Delete 文件普通实时入湖就是 append只追加不修改。但真实业务里比如用 Flink CDC 同步业务库 MySQL 的数据同一主键的数据可能会更新你需要的是 UPSERT 语义而不是简单地追加。Iceberg 对 Upsert 的实现是 Merge-on-Read。你在 Flink 表上开了write.upsert.enabledtrue后写入时除了常规数据文件还会为这一批数据里涉及到的主键生成一个 equality delete 文件。这个 delete 文件记录的是主键值意思是“读的时候下面这些主键的数据要用新数据替换旧数据不要了”。读取时Iceberg 会先把数据文件里的内容读出来再根据 delete 文件过滤掉老版本数据最后把新数据合并进去。物理上老数据文件还在硬盘上只是逻辑上被删了。这套设计对 Flink 写入非常友好因为 Flink 不需要去随机更新已有的数据文件只需要顺序写文件这是流式引擎最擅长做的事。代价是如果长期不整理文件delete 文件会越积越多读性能会下降。所以 Upsert 表必须配套定期跑 compaction把旧数据文件里被标记删除的数据真正清理掉。3.3 写入参数该怎么配我给一份实际在用的从实际经验看下面这份配置组合比较稳适合大多数分钟级延迟的场景。CREATE TABLE ods.user_event ( userId BIGINT, eventTime TIMESTAMP(3), eventType STRING, eventData STRING, PRIMARY KEY (userId, eventTime) NOT ENFORCED ) PARTITIONED BY (days(eventTime)) WITH ( connector iceberg, catalog-name iceberg_catalog, catalog-type hive, write.format.default parquet, write.target-file-size-bytes 134217728, write.upsert.enabled true, write.parquet.compression-codec zstd, commit.retry.num-retries 5 );write.target-file-size-bytes我一般设 128MB。Iceberg 默认值偏大Partition 数量多的时候容易产生超大文件太小也不行文件数量会失控。write.parquet.compression-codec建议统一。如果下游主要用 Trino、Sparkzstd 的压缩率和解压速度综合表现最好如果集群是老版本 Hive 组件用 snappy 更兼容。commit.retry.num-retries必须给够。Flink 多个并发写同一张表时commit 冲突是常态默认重试次数太少会直接任务失败设到 5 次能扛住大部分抖动。写完之后你在 Flink SQL 里正常执行 INSERT INTO 就行不用关心文件提交这些底层动作。如果想要更精细的控制可以使用 DataStream API 来指定快照 ID 或自定义序列化 Schema。大多数生产场景下Flink SQL 已经够用。4. Flink 读 Iceberg批读、流读和增量读三种姿势要分清4.1 批式读取本质是给快照加一个光标Flink 批读 Iceberg执行的是一个普通的批查询默认读最新快照的数据。如果你要查历史某个时间点的数据可以在查询时指定scan.snapshot-id或者scan.as-of-timestamp。这个能力叫时间旅行。时间旅行在实际排障里太好用了。数据出问题的时候不用重跑整个链路直接读任务出问题之前的那个快照对比一下就知道数据从哪个环节开始坏的。以前用 Hive 定位这种问题只能靠日志猜现在几行 SQL 就能查。批读场景下 Flink 就是标准的批引擎把 Iceberg scan 出来的文件列表全部拉过来做分布式处理。这里没有流的概念状态、checkpoint、watermark 这些都不涉及就是一个纯粹的分布式 SQL 查询。4.2 流式读取其实是在持续发现新快照Flink 流读 Iceberg 是我觉得最能体现二者协作深度的一个能力。流式读的底层逻辑是冰伯格表每提交一个新快照Flink 就把它当成一个微批次的数据源投喂给下游计算。Flink 的任务长期运行一遍遍地发现新快照、读取新数据、交给下游算子。用 SQL 开流读的话DDL 里需要加 startup 参数。例如从最早快照开始读CREATE TABLE user_event_rt ( userId BIGINT, eventTime TIMESTAMP(3), eventType STRING ) WITH ( connector iceberg, catalog-name iceberg_catalog, scan.startup.mode earliest );这里有几个实际经验要分享第一流读默认模式下如果没仔细设置 startup很可能只会拿到任务启动之后新产生的快照历史数据一概不读。具体行为跟 connector 的默认参数有关但我不建议依赖默认值显式声明scan.startup.mode是最稳的做法。第二流读的单位是快照不是单条记录。上游连续提交了 3 个快照Flink 可能一次性把这 3 个快照的数据都读出来。所以 Iceberg 流读天然带一点小批量性质和 Kafka 的逐条消费不太一样。做窗口聚合的时候watermark 可能会有小幅跳动需要做一下测试验证结果稳定性。第三任务重启后Flink 会通过 checkpoint 记录已经处理到哪个快照 ID 了恢复后从那个位置继续不会把历史快照全部重读一遍。这是 Flink 联合 Iceberg 协议里已经做好的位点管理比你自己去维护 Kafka offset 要省心得多。4.3 增量读取比流读更精确地定位到两次快照之间流读是持续追踪新快照增量读则是明确指定一个起始快照 ID然后读取从这个快照之后新增的所有快照数据。对于需要做“重跑某段时间增量数据”的场景这个能力非常实用。比如你有一套每日指标计算任务某一天结果算错了你只需要重跑那一天的增量数据找到当天零点的快照 ID 作为起点读到最新快照即可不用把全表数据再拉一遍。计算成本节省得不止一点点。5. Catalog 选型与实操从建库到写读直接照抄的配置5.1 三种 Catalog 怎么选我直接给结论Iceberg 支持在 Flink 里配置三种 CatalogHiveCatalog、HadoopCatalog、RESTCatalog。Catalog 类型适用场景优点缺点HiveCatalog生产环境、已存在 Hive Metastore和 Hive 全家桶元数据互通下游 Hive/Spark 能直接看到表依赖 HMS 高可用HadoopCatalog独立测试环境、纯 Flink 场景不需要 HMS配置最简单直接指向 warehouse 路径其他引擎要自己接一套元数据RESTCatalog多引擎共享、云原生环境统一元数据服务跨 Flink/Spark/Trino 都走同一套 API需要额外部署 REST 服务我的建议是只要你的集群里已经跑着 Hive 或者 Spark直接用 HiveCatalog。它能最大程度利用现有基础设施下游引擎查表也方便。如果只是自己搭一套纯 Flink Iceberg 的实验环境HadoopCatalog 三分钟就能跑起来。RESTCatalog 适合公司级的数据湖平台建设但初期没必要把复杂度拉这么高。5.2 Flink SQL 完整建库建表读写示例下面这段 SQL 可以直接抄着用。-- 创建 Catalog CREATE CATALOG iceberg_catalog WITH ( type iceberg, catalog-type hive, uri thrift://hive-metastore:9083, clients 5, warehouse hdfs://nameservice/warehouse ); USE CATALOG iceberg_catalog; -- 创建数据库 CREATE DATABASE IF NOT EXISTS ods; -- 创建 Iceberg 表 CREATE TABLE IF NOT EXISTS ods.user_log ( user_id BIGINT, event_ts TIMESTAMP(3), event_action STRING, event_detail STRING ) PARTITIONED BY (days(event_ts)) WITH ( write.format.default parquet, write.target-file-size-bytes 134217728 ); -- 实时写入 INSERT INTO ods.user_log SELECT user_id, event_ts, event_action, event_detail FROM kafka_user_log; -- 批读最新快照 SELECT * FROM ods.user_log WHERE event_action click; -- 批读指定时间点的历史快照 SET scan.as-of-timestamp 1700000000000; SELECT count(*) FROM ods.user_log;Flink SQL 里跑批查询时如果要用时间旅行可以在环境级别 SET 一个参数也可以直接在建视图时把参数放到表属性里。具体写法各家版本略有差异以你实际使用的 Flink 版本对应的文档为准。跑一下就能通不算复杂。5.3 DataStream API 写入怎么做纯 SQL 写不够用的时候比如要从一个自定义 Source 拿数据写 Iceberg可以直接用 Flink 的 Iceberg Sink。参考写法如下DataStreamRowData stream ...; FlinkSink.forRowData(stream) .table(icebergTable) .overwrite(false) .build();底层它还是走我们前面说的 checkpoint 提交机制所以不管 SQL 还是 DataStream提交语义是一致的。在 DataStream 里用的时候建议手动检查一下待写数据的 Schema 和 Iceberg 表结构的字段名是否完全一致这个位置很容易出字段顺序对不上、数据写错列的问题。6. 生产环境里我踩过的几个坑按排查过程写给你6.1 多个 Flink 任务写同一张表的 Commit 冲突现象两个流任务同时往一张 Iceberg 表里写数据跑了一段时间后其中一个任务突然报错日志里出现CommitFailedException提示当前元数据落后于表的最新元数据重试几次后失败退出。排查链路先看是不是存储压力导致排除再看是不是两个任务写同一批分区发现分区不同也会冲突。最后定位到根因Iceberg 的提交是乐观并发控制两个事务同时基于同一个旧元数据版本提交后提交的那个会失败因为表的最新元数据已经变了。解决办法是两层。第一层调大重试参数commit.retry.num-retries设 5 到 10commit.retry.min-wait-ms设 100commit.retry.max-wait-ms设 60000。第二层从架构上规避高并发写同一张表能合并的写入任务尽量合并或者按业务域拆表。我后来就是把三个写入任务合并成一个冲突彻底消失。6.2 快照文件暴涨HDFS NameNode 差点扛不住现象某天收到 NameNode 服务告警RPC 处理延迟升高。查了一下 warehouse 目录发现一张表的 metadata 目录下全是快照文件和 manifest 文件数量体量巨大。排查链路检查这张表的写入任务的 checkpoint 间隔发现被设置成了 15 秒一次。每次 checkpoint 提交一个快照一天 5760 个快照。再查表属性里有没有配置快照过期发现完全没有。也就是说所有历史快照都保存着元数据越积越多。解决办法是改表属性加上快照过期策略ALTER TABLE ods.user_log SET TBLPROPERTIES ( history.expire.max-snapshot-age-ms 43200000, history.expire.min-snapshots-to-keep 20 );同时把 checkpoint 间隔改成 2 分钟。改完之后过一天再查元数据文件数量稳定在一个低位。这个坑非常典型尤其是在实时接入高峰流量的时候很多同学很容易被“写入延迟越低越好”带偏最后付出的是元数据层面的性能代价。6.3 Flink 增量加列后下游新字段全是 Null现象给 Iceberg 表加了一列country上游 Flink 任务也正常输出了这个字段。但下游查数时这一列全是 null没有任何报错。排查链路先查上游任务是否真的输出新字段没问题。再查 Flink 写入端的表 DDL发现写入任务里用的建表语句还是老的没有加country。问题就出在这写入任务读取源数据时由于其自身的 Flink 表结构里没有country该字段根本没有被写入到 Iceberg 表。这个坑的核心是Iceberg 表的 schema 演进了但写入任务的“视图”没更新。Flink 写入端要确保自己读到的字段和写出的字段都对齐最新的 Iceberg schema不能只改 Iceberg 那边。解决办法是重启写入任务前把 Flink SQL 里的 DDL 同步改成新结构再提交新的写入任务。从此之后我所有 Iceberg 相关任务上线前都会加一条检查比对源端字段列表、Flink 表字段列表、Iceberg 实际表字段列表三者是否一致。这个检查 10 分钟能省下后面一整天排查的时间。6.4 基于隐藏分区的过滤条件在 Flink 里失效了现象一张按days(event_ts)做隐藏分区的表在 Flink SQL 里查WHERE event_ts ...执行计划显示还是全表扫描数据量特别大的时候查询慢到没法接受。排查链路先怀疑是不是 Iceberg 和 Flink 的谓词下推不生效。查官方文档发现Iceberg 的扫描器是支持 Flink 谓词下推到 manifest 层做过滤的但对分区 transform 函数的推导支持没有 Spark 那么完整。在某些 Flink 版本上days(ts)这种虚拟分区列的过滤没法直接推到物理数据文件裁剪。解决办法是两条路都走第一在查询条件里显式写过滤表达式第二更稳妥的做法是在红格表里冗余一个物理分区字段比如直接存一个dt字符串字段查询时直接用dt过滤。这样做虽然牺牲了一些“隐藏分区”的优雅性但换来了 Flink 查询的可预测性。数据入湖时多写一个字段的成本微乎其微查询性能的收益却是数量级的。7. 把这套组合落到架构里常见模式与配套动作7.1 最常用的实时入湖 流读架构我这边落地最多的架构长这样业务库的 Binlog → Flink CDC 采集 → Kafka 做缓冲 → Flink SQL 做清洗、维表关联 → 写入 Iceberg ODS 层 → 下游用 Flink 流读 ODS 层做实时汇总或者用 Trino、Spark 做批量分析。这套架构里 Flink 和 Iceberg 的协作体现在两个位置一个是在写入端Flink 的 checkpoint 把实时数据变成一个个 Iceberg 快照另一个是在读取端Flink 流读把 Iceberg 的快照变成持续的微批数据流。整条链路从业务库到最终分析延迟在分钟级别同时拥有完整的快照隔离和回溯能力。这里有个设计上的小心得不要在写入 Iceberg 之前把数据全量攒在 Kafka 里太久。Kafka 的存储成本和数据回溯能力远不如 Iceberg既然要入湖就让数据以更快速度落到 Iceberg。Kafka 只做削峰和缓冲保持半天到一天的保留期足够。7.2 CDC 场景下很值得关注的几点Flink CDC 写 Iceberg 是我最推荐的组合之一。CDC 天然是主键更新流Iceberg 的 Upsert 模式正好匹配Flink 的 checkpoint 正好可以和 Binlog 位点管理统一起来实现从源头到存储的端到端一致性。但要注意CDC 场景写入频率通常很高如果不做限制每秒一个事务就要产生一个新快照。我是这样处理的要求写入任务开启 mini-batch 攒批同时 checkpoint 间隔不低于 1 分钟表上开启定期 compaction。Compaction 这块Flink 没有内置自动执行我是写一个独立的 Java 任务或者用 Spark 任务每天凌晨跑一次RewriteDataFilesActionResult result IcebergActions.rewriteDataFiles(table);具体写法不细展开了。核心思路是读取 snapshot 里那些 size 明显偏小的数据文件把它们重写为接近目标大小的大文件同时清理掉已经无效的 delete 文件。7.3 版本选型和日常巡检建议Flink 和 Iceberg 的版本兼容非常敏感不是随便拿两个最新版本拼一起就能跑的。我整理了一份目前比较稳的组合表供参考Flink 版本Iceberg 版本Flink 1.14Iceberg 0.13 / 0.14Flink 1.15Iceberg 1.1 / 1.2Flink 1.16Iceberg 1.2Flink 1.17Iceberg 1.3Flink 1.18Iceberg 1.4 / 1.5Flink 1.19Iceberg 1.5 / 1.6选型原则就一条以官方兼容矩阵为准不要跨越太大跨度升级。曾经升级到新版本后connector 类路径变了任务根本起不来回滚又花了一个小时这种教训经历一次就够了。日常巡检我主要看五件事快照数量增长率、孤儿文件数量、数据文件大小分布、commit 失败率、任务重启恢复时间。这五件事能提前发现大部分潜在问题。快照看增长曲线小文件看分布commit 失败看日志统计。我自己写了一个巡检脚本每天跑一次超过阈值就报警。做数据平台到了后期拼的不是技术上限是这种底线监控的细致程度。最后再啰嗦一句Flink 和 Iceberg 的最佳协作方式不是把参数堆到最高也不是追求极致的低延迟而是找到你的业务能接受的延迟和系统能承受的元数据开销之间的平衡点。批读、流读、快照回溯、增量计算这些能力都是在这个平衡点上长出来的。先把监控和巡检做起来再考虑其他优化顺序别搞反了。