
简介这份PPT资料面向大数据开发工程师、数据平台架构师及实时数仓方向的技术人员系统讲解如何基于Flink与Iceberg搭建企业级实时数据湖帮助读者理解数据湖分层架构与流批一体落地路径。内容围绕数据湖背景、Flink数据湖业务场景、为何选择Iceberg三大模块展开涵盖存储层、加速层、Table Format层与计算引擎层的职责划分并深入剖析构建实时Data Pipeline、CDC数据实时摄入、近实时流批统一、从Iceberg历史数据启动Flink任务等典型场景同时对比Delta、Hudi、Iceberg三大开源项目的ACID、隔离级别、时间旅行等特性说明Iceberg与Flink在流批一体规划上的契合度。资源包共1个pptx文件约2.94MB结构清晰、图文并茂适合作为技术分享或内部培训材料。目前已有562人学习可帮助读者快速建立实时数据湖的整体认知与选型依据。1. 实时数据湖落地Flink 写 Iceberg 到底解决了哪几个真问题很多团队第一次听到「基于 Flink Iceberg 构建企业级实时数据湖」脑子里浮现的是又一套要养一堆组件的大数据架构。但真正落到生产里它要解决的其实是很具体的三件事业务库的 binlog 要能分钟级进湖、进湖后的数据要能被 Spark 和 Trino 直接查、历史数据要能按分区做增量更新而不是每天全量重刷。这三件事用传统 Hive 数仓做要么延迟按小时算要么小文件多到 NameNode 扛不住要么 upsert 根本没法优雅实现。Flink 负责流式写入和 exactly-once 语义Iceberg 负责表格式、快照隔离和隐藏分区两者拼起来才是一个能持续演进、能回溯、能对接多种计算引擎的实时数据湖底座。这套方案适合已经有 Kafka、有一定 Flink 运维能力、并且下游同时存在实时看板和离线分析的团队。如果你只是想把 MySQL 同步到 ClickHouse 做点实时报表那用 Flink JDBC 连接器直连就够了不必上 Iceberg 这层。但只要你开始遇到「同一份数据既要实时又要离线、还要能改历史」的需求Iceberg 的价值就出来了。2. 为什么是 Flink 加 Iceberg选型逻辑与最小可跑环境2.1 流批一体的写入语义到底靠什么撑住Flink 写 Iceberg 的核心不是「能写进去」而是「写进去之后下游读到的是一致快照」。Iceberg 的每次 commit 都会生成一个新的 metadata 文件Flink 的 IcebergSink 在 checkpoint 完成时才真正提交一次 snapshot。这意味着如果 Flink 任务在两次 checkpoint 之间挂了重启后从上一个成功 checkpoint 恢复Iceberg 表里不会出现半截数据。这个机制是 Iceberg 表格式本身提供的不是 Flink 额外加的。对比直接写 HDFS 或者写 Hive 分区表区别在于写 Hive 时文件一旦 rename 成功就可见了下游可能读到只写了一半的分区而 Iceberg 的 snapshot 是原子切换的读端要么看到旧快照要么看到新快照。这就是为什么做实时数据湖时Iceberg 的 commit 机制比单纯的文件写入更关键。另一个常被忽略的点是 Iceberg 支持行级 delete 和 upsert。Flink CDC 同步 MySQL 时会产生 INSERT、UPDATE、DELETE 三种事件如果目标端只支持 append那 UPDATE 和 DELETE 就没法正确处理。Iceberg 通过 equality delete 和 positional delete 文件来标记删除配合 Flink 的 upsert 模式才能让湖里的数据和源库保持一致。2.2 本地用 Docker 把 Iceberg MinIO Spark 跑起来在正式上生产之前我一般会先在本地用 Docker 搭一套最小环境验证写入链路。MinIO 充当 S3 兼容存储Spark 用来验证读Flink 负责写。下面这份 docker-compose 是我常用的最小组合。version: 3 services: minio: image: minio/minio:latest ports: - 9000:9000 - 9001:9001 environment: MINIO_ROOT_USER: admin MINIO_ROOT_PASSWORD: password command: server /data --console-address :9001 volumes: - ./minio-data:/data iceberg-rest: image: tabulario/iceberg-rest:latest ports: - 8181:8181 environment: CATALOG_WAREHOUSE: s3://warehouse/ CATALOG_IO__IMPL: org.apache.iceberg.aws.s3.S3FileIO CATALOG_S3_ENDPOINT: http://minio:9000 CATALOG_S3_ACCESS__KEY__ID: admin CATALOG_S3_SECRET__ACCESS__KEY: password CATALOG_S3_PATH__STYLE__ACCESS: true depends_on: - minio这段配置里几个参数值得说清楚。CATALOG_WAREHOUSE指向 MinIO 里的 bucket 路径Iceberg 的所有 metadata 和数据文件都会落在这里。CATALOG_S3_PATH__STYLE__ACCESS必须设为 true因为 MinIO 默认不支持 virtual-hosted style 访问不设这个参数 Flink 写的时候会报 400。CATALOG_IO__IMPL指定用 S3FileIO否则 Iceberg 会尝试用 Hadoop 的 S3A 实现在容器环境里容易因为缺少 hadoop-aws 依赖而失败。启动之后先在 MinIO 控制台建一个名为 warehouse 的 bucket然后就可以用 Spark 建表验证 catalog 是否通了。# 进入 spark-sql 容器后执行 spark-sql --packages org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.4.3 \ --conf spark.sql.catalog.localorg.apache.iceberg.spark.SparkCatalog \ --conf spark.sql.catalog.local.typerest \ --conf spark.sql.catalog.local.urihttp://iceberg-rest:8181 \ --conf spark.sql.catalog.local.warehouses3://warehouse/ \ --conf spark.sql.catalog.local.io-implorg.apache.iceberg.aws.s3.S3FileIO \ --conf spark.sql.catalog.local.s3.endpointhttp://minio:9000 \ --conf spark.sql.catalog.local.s3.path-style-accesstrue建表语句用 Iceberg 的语法指定分区和主键这是后面 Flink upsert 能生效的前提。CREATE TABLE local.db.orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10,2), status STRING, update_time TIMESTAMP ) USING iceberg TBLPROPERTIES ( write.upsert.enabled true, format-version 2 );format-version必须设为 2因为 v1 不支持 row-level deleteupsert 会退化成 append。write.upsert.enabled是 Flink 写入时的表属性告诉 Iceberg 这张表允许按主键做更新。注意这里没有显式定义主键Iceberg 的主键约束是在 Flink 写入时通过upsert-enabled和equality-fields指定的建表时只需要保证 format-version 是 2。3. Flink 写入 Iceberg 的作业配置与参数调优3.1 DataStream API 写 Iceberg 的完整代码骨架用 Flink SQL 写 Iceberg 看起来简单但生产里我更多用 DataStream API因为对 checkpoint、并行度、失败重试的控制更细。下面是一个从 Kafka 读 JSON 然后 upsert 到 Iceberg 的骨架。StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(60000); // 1 分钟一次 checkpoint env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000); env.getCheckpointConfig().setCheckpointTimeout(120000); // Kafka source 配置 KafkaSourceString source KafkaSource.Stringbuilder() .setBootstrapServers(kafka:9092) .setTopics(orders) .setGroupId(flink-iceberg-orders) .setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST)) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); DataStreamString stream env.fromSource(source, WatermarkStrategy.noWatermarks(), kafka-source); // 解析 JSON 并转换为 Row DataStreamRow rows stream.map(new JsonToRowMapFunction()).returns( Types.ROW_NAMED( new String[]{order_id, user_id, amount, status, update_time}, Types.LONG, Types.LONG, Types.DECIMAL(10, 2), Types.STRING, Types.SQL_TIMESTAMP ) ); // Iceberg sink 配置 Configuration icebergConf new Configuration(); icebergConf.setString(catalog-type, rest); icebergConf.setString(uri, http://iceberg-rest:8181); icebergConf.setString(warehouse, s3://warehouse/); icebergConf.setString(io-impl, org.apache.iceberg.aws.s3.S3FileIO); icebergConf.setString(s3.endpoint, http://minio:9000); icebergConf.setString(s3.path-style-access, true); CatalogLoader catalogLoader CatalogLoader.rest(iceberg_catalog, icebergConf); TableLoader tableLoader TableLoader.fromCatalog(catalogLoader, TableIdentifier.of(db, orders)); FlinkSink.forRow(rows, tableLoader) .tableSchema(TableSchema.builder() .field(order_id, DataTypes.BIGINT()) .field(user_id, DataTypes.BIGINT()) .field(amount, DataTypes.DECIMAL(10, 2)) .field(status, DataTypes.STRING()) .field(update_time, DataTypes.TIMESTAMP()) .build()) .upsert(true) .equalityFieldColumns(Arrays.asList(order_id)) .writeParallelism(4) .build(); env.execute(flink-iceberg-orders);这段代码里几个关键点。enableCheckpointing(60000)决定了数据可见延迟Iceberg 的 snapshot 是在 checkpoint 完成时提交的所以下游能看到数据的最快时间就是一个 checkpoint 间隔。upsert(true)开启 upsert 模式equalityFieldColumns指定 order_id 作为等值删除的键这样同一条 order_id 的新数据会覆盖旧数据。writeParallelism(4)控制写入并行度这个值不是越大越好后面避坑章节会讲。3.2 三个必须调对的参数checkpoint 间隔、写入并行度、目标文件大小checkpoint 间隔直接决定数据可见延迟和 snapshot 数量。设成 1 分钟意味着每分钟产生一个 snapshot一天就是 1440 个 snapshot。Iceberg 的 metadata 会随着 snapshot 增多而膨胀虽然可以定期 expire但太频繁的 snapshot 会给元数据管理带来压力。我一般建议生产环境设 3 到 5 分钟除非业务对延迟极度敏感。如果设成 10 秒那基本是在给自己挖坑metadata 文件会多到 Spark 读的时候光解析元数据就要好几秒。写入并行度要和 Kafka 分区数匹配。如果 Kafka topic 有 12 个分区Flink source 并行度是 12那 sink 并行度设 4 到 6 比较合适。设太大反而会产生大量小文件因为每个并行子任务都会独立写文件。Iceberg 虽然有write.target-file-size-bytes控制目标文件大小默认 512MB但并行度太高时每个子任务分到的数据量不够就会写出很多远小于目标值的文件。目标文件大小这个参数在表属性里设置write.target-file-size-bytes默认 512MB。对于实时写入场景我一般会调到 128MB 到 256MB。因为实时写入的数据量不像批处理那么集中如果目标设成 512MB可能一个文件要写很久才 close期间下游读不到这部分数据。设小一点能让文件更快落盘但也不能太小否则小文件问题又回来了。128MB 是一个比较平衡的值。ALTER TABLE local.db.orders SET TBLPROPERTIES ( write.target-file-size-bytes 134217728, write.distribution-mode hash );write.distribution-mode设为 hash 可以让相同主键的数据落到同一个写入任务减少 upsert 时的冲突。如果是 append-only 场景设成 none 也可以。4. 避坑与排查Flink 写 Iceberg 最常见的五类翻车4.1 现象任务一直 running 但 Iceberg 表里查不到数据原因通常是 checkpoint 没成功。Iceberg 的 snapshot 只在 checkpoint 完成时提交如果 checkpoint 一直失败或者耗时超过 timeout数据就只停留在 Flink 状态里没有落到 Iceberg。常见触发点是 checkpoint 超时设得太短或者状态后端用的是默认的 HashMapStateBackend 但状态太大导致 checkpoint 慢。解决方法是先看 Flink UI 的 Checkpoints 页面确认 checkpoint 是否成功。如果失败看异常信息里是不是有Checkpoint expired before completing。把setCheckpointTimeout从默认的 10 分钟适当调大同时检查状态后端是否配置了 RocksDB大状态场景下必须用 RocksDB。另外确认setMinPauseBetweenCheckpoints不要设得太小否则 checkpoint 之间互相挤压也会导致失败。4.2 现象Iceberg 表里出现大量几 KB 的小文件原因是写入并行度过高加上 checkpoint 间隔太短。每个 checkpoint 周期内每个写入子任务都会生成新文件如果并行度是 16、checkpoint 间隔是 1 分钟那一天就是 16 × 1440 23040 个文件。每个文件可能只有几百 KB远小于目标文件大小。解决方法是降低写入并行度同时把 checkpoint 间隔调到 3 到 5 分钟。另外可以开启 Iceberg 的自动 compaction在表属性里设置write.metadata.compression-codec为 gzip 减少元数据大小并定期用 Spark 执行CALL local.system.rewrite_data_files(db.orders)来合并小文件。生产环境我一般会配一个定时 compaction 任务每天凌晨跑一次。4.3 现象Flink 报错Cannot find catalog plugin class for catalog iceberg_catalog这是依赖没打对。Flink 写 Iceberg 需要iceberg-flink-runtime这个 shaded 包而不是单独的iceberg-core加iceberg-flink。如果用 Maven 构建依赖应该写成dependency groupIdorg.apache.iceberg/groupId artifactIdiceberg-flink-runtime-1.17/artifactId version1.4.3/version /dependency版本号要和 Flink 版本对应Flink 1.17 用iceberg-flink-runtime-1.17Flink 1.18 用iceberg-flink-runtime-1.18。如果版本不匹配运行时会报NoSuchMethodError或者类找不到。另外这个包要放在 Flink 的 lib 目录下或者用--classpath指定不能只放在用户 jar 里因为 catalog 加载是在 Flink 的插件机制里做的。4.4 现象upsert 不生效同一条数据出现多条记录原因是建表时format-version设成了 1或者 Flink sink 没有正确指定equalityFieldColumns。v1 表格式不支持 row-level deleteupsert 会退化成 append同一条 order_id 的更新会变成两条记录。解决方法是确认表属性里format-version是 2并且 Flink sink 里upsert(true)和equalityFieldColumns都设置了。如果表已经建成了 v1可以用ALTER TABLE db.orders SET TBLPROPERTIES (format-version 2)升级但升级后需要重写历史数据才能让 delete 文件生效。所以最好在建表时就设对。4.5 现象Spark 读 Iceberg 表时报Not an Avro data file这是 MinIO 的 path-style 配置没生效。Iceberg 写 metadata 时用的是 Avro 格式如果 S3 客户端配置成了 virtual-hosted styleMinIO 返回的响应可能不是 Iceberg 期望的格式导致 Spark 读的时候解析失败。解决方法是确认所有访问 Iceberg 的组件都设置了s3.path-style-accesstrue。Flink 侧在Configuration里设Spark 侧在spark.sql.catalog.local.s3.path-style-access里设Iceberg REST catalog 侧在环境变量CATALOG_S3_PATH__STYLE__ACCESS里设。三处缺一不可少一处就可能出现某个组件能写但另一个组件读不了的情况。5. 进阶技巧用 Flink SQL 做 MySQL 到 Iceberg 的整库同步5.1 用 Flink CDC 加 Iceberg 实现分钟级入湖前面讲的都是单表写入实际企业场景里更常见的是整库同步。用 Flink CDC 的 MySQL source 加上 Iceberg sink可以在一个作业里同步多张表。下面是一个 Flink SQL 的示例把 MySQL 的 orders 表实时同步到 Iceberg。-- 创建 MySQL CDC source CREATE TABLE mysql_orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10,2), status STRING, update_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname mysql, port 3306, username cdc_user, password cdc_pass, database-name biz, table-name orders, server-time-zone Asia/Shanghai ); -- 创建 Iceberg sink CREATE TABLE iceberg_orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10,2), status STRING, update_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector iceberg, catalog-type rest, uri http://iceberg-rest:8181, warehouse s3://warehouse/, io-impl org.apache.iceberg.aws.s3.S3FileIO, s3.endpoint http://minio:9000, s3.path-style-access true, format-version 2, write.upsert.enabled true ); -- 启动同步 INSERT INTO iceberg_orders SELECT * FROM mysql_orders;这段 SQL 里PRIMARY KEY (order_id) NOT ENFORCED在 source 和 sink 两侧都要声明Flink 会根据主键自动生成 upsert 逻辑。server-time-zone设成 Asia/Shanghai 是为了让 update_time 字段的时区正确不设的话默认是 UTC下游看数据时间会差 8 小时。整库同步时每张表都要建一对 source 和 sink 表然后分别 INSERT。如果表很多可以用 Flink CDC 的scan.incremental.snapshot.enabled参数开启增量快照避免全量阶段锁表。这个参数在 MySQL CDC source 的 WITH 里加设成 true 即可。5.2 验证数据一致性的两个实用查询同步跑起来之后怎么确认 Iceberg 里的数据和 MySQL 一致我一般用两个查询。第一个是查 Iceberg 表的 snapshot 历史确认最近有没有新数据写入。SELECT * FROM local.db.orders.snapshots ORDER BY committed_at DESC LIMIT 5;这个查询会返回最近的 snapshot 记录包括 snapshot_id、committed_at 和时间戳。如果 committed_at 一直在更新说明 Flink 在持续写入。如果超过 10 分钟没有新 snapshot就要去检查 Flink 任务是不是卡住了。第二个查询是对比源库和目标库的行数。在 MySQL 里执行SELECT COUNT(*) FROM orders在 Spark 里执行SELECT COUNT(*) FROM local.db.orders。如果两边差距在分钟级延迟范围内说明同步正常。如果差距很大可能是 upsert 没生效导致重复数据或者 delete 事件没被正确处理。这时候要回去检查format-version和equalityFieldColumns的配置。我自己的习惯是每次上新的同步链路先跑一天观察 snapshot 数量和文件大小分布确认没有小文件爆炸再正式切流量。实时数据湖这套东西写入链路稳不稳定看 snapshot 历史和文件大小分布基本就能判断个八九不离十。希望帮到你。本文还有配套的精品资源点击获取