THS_Allen头像
关注
Apache Iceberg + Spark Streaming 构建 Lakehouse 实时数仓封面图

Apache Iceberg + Spark Streaming 构建 Lakehouse 实时数仓

Apache Iceberg + Spark Streaming 构建 Lakehouse 实时数仓:CDC 增量入湖与查询加速实践

摘要:传统 Lambda 架构在实时数仓建设中面临数据一致性与维护成本的双重挑战。本文以 Apache Iceberg 表格式为核心,结合 Spark Structured Streaming 与 CDC(变更数据捕获)技术,详细阐述如何构建一套支撑分钟级延迟的 Lakehouse 实时数仓。内容涵盖 Iceberg 元数据管理机制、Spark-Iceberg 增量读写原理、Flink CDC 整库同步方案以及查询层的分区裁剪与文件编排优化,提供可直接落地的代码与配置。


在这里插入图片描述

一、Lambda 架构的困境与 Lakehouse 的破局

在过去五年的数据平台建设中,我们团队一直沿用经典的 Lambda 架构:

  • 批处理层:Spark 每日 T+1 处理 Hive 数据,生成历史全量报表;
  • 速度层:Flink 实时消费 Kafka,写入 HBase / ClickHouse 供实时查询;
  • 服务层:合并批与流的结果,暴露给 BI 工具。

这套架构的问题在数据规模突破 PB 级后集中爆发:同一业务逻辑需要维护批流两套代码,数据口径不一致引发的"数字对不上"成为分析师的日常噩梦。更严重的是,HBase 的 Schema 变更几乎等同于停服重建,无法满足业务快速迭代的需要。

Lakehouse 架构的提出,核心在于用开放的表格式(Apache Iceberg / Hudi / Delta Lake)在对象存储之上实现数仓的 ACID 语义。其中 Apache Iceberg 凭借其:

  • 纯开源、无 vendor lock-in:不绑定特定计算引擎;
  • 优秀的生态系统:Spark、Flink、Trino、StarRocks 均有成熟连接器;
  • 先进的元数据设计:隐式分区、Time-Travel、Partition Evolution;

成为我们在 2024-2025 年技术升级的首选。


二、Iceberg 元数据架构:理解隐式分区的钥匙

很多工程师初次接触 Iceberg 时,会困惑于"为什么查询时不需要指定分区字段"。要回答这个问题,必须深入理解 Iceberg 的三层元数据模型。

2.1 Catalog → Table → Snapshot → Manifest → DataFile

Iceberg 的元数据分为以下层级:在这里插入图片描述

Catalog(Hive / Hadoop / JDBC / REST)
    └── Table Metadata JSON
            └── Snapshot List(多版本快照)
                    └── Manifest List
                            └── Manifest File(分区统计信息 + DataFile 列表)
                                    └── DataFile(Parquet / ORC / Avro)

关键设计:每个 Snapshot 是一个不可变的表状态。当执行 INSERTUPDATEDELETE 时,Iceberg 不会修改任何已有 DataFile,而是写入新的 DataFile,并在 Manifest 中记录文件的 min/max 统计信息,最后生成一个新的 Snapshot 并切换 current-snapshot-id

这意味着:

  • Time-Travel 天然支持:SELECT * FROM table TIMESTAMP AS OF '2025-06-01 10:00:00' 只需回溯到对应 Snapshot;
  • 并发写入安全:乐观锁机制下,两个 Spark Job 同时提交时,后提交的 Job 会检测到元数据版本变化并自动重试;
  • 分区演进无痛:修改分区策略不会影响历史数据,新数据按新分区写入,查询时引擎自动选择最优裁剪策略。

2.2 隐式分区与查询裁剪

传统 Hive 表需要用户显式匹配分区字段(如 WHERE dt='2025-08-11'),而 Iceberg 在 Manifest 文件中记录了每个 DataFile 的列级统计信息(min/max、null count、distinct values)。当查询带有过滤条件时,Iceberg 通过底层 API 的 planFiles() 方法,在读取任何 Parquet 文件之前,先根据统计信息过滤掉不相关的 DataFile。

在我们的生产环境中,一张 500 亿行的用户行为表,通过 Iceberg 的隐式分区 + Z-Order 排序,将全表扫描查询的 IO 量降低了 97%。


三、CDC 增量入湖:Flink CDC + Iceberg 整库同步

实时数仓的核心是数据新鲜度。我们需要将 MySQL、Oracle 等业务库的变更实时捕获并写入 Iceberg。
在这里插入图片描述

3.1 技术选型:Flink CDC 3.0

Flink CDC 3.0 引入了整库同步能力,支持 Schema Evolution(加列、改类型)自动同步到下游 Iceberg 表。相比早期版本需要为每张表单独写 Flink Job,3.0 版本只需一个 YAML 配置文件即可同步整库。

3.2 YAML 配置实战

以下是我们同步核心订单库的配置:

source:
  type: mysql
  hostname: mysql-primary.db.svc
  port: 3306
  username: cdc_user
  password: ${CDC_PASSWORD}
  tables: order_db.order_info,order_db.order_detail
  server-id: 5400-5404

sink:
  type: iceberg
  catalog:
    type: hadoop
    warehouse: hdfs://namenode:8020/warehouse/iceberg
  table-prefix: cdc_
  table-defaults:
    format-version: 2
    write.metadata.metrics.default: counts   # 默认收集列统计
    write.distribution-mode: hash            # 按主键哈希分布,避免小文件

pipeline:
  parallelism: 8
  schema-change-mode: evolve               # 自动同步 Schema 变更

关键参数解析

  • format-version: 2:Iceberg V2 表格式支持行级 UPDATE/DELETE(基于 position/equality delete files),是 CDC 场景的必要条件;
  • write.distribution-mode: hash:确保同一主键的数据落在同一文件,减少后续 Merge-on-Read 时的文件扫描范围;
  • schema-change-mode: evolve:当上游 MySQL 执行 ALTER TABLE ADD COLUMN 时,Flink CDC 自动在 Iceberg 表上执行对应变更,无需人工介入。

3.3 写入优化:WAL 与 Checkpoint 调优

Flink CDC 写入 Iceberg 时,每个 Checkpoint 会触发一次 Iceberg Commit。如果 Checkpoint 间隔过短(如 1 秒),会产生大量小文件和元数据膨胀;如果过长(如 10 分钟),数据延迟又会超标。

我们的调优策略是分层 Checkpoint

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(30000);  // Checkpoint 30 秒,平衡延迟与小文件
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(20000);

配合 Iceberg 的 write.target-file-size-bytes=134217728(128MB),使每个 DataFile 大小稳定在 100-150MB 之间,兼顾查询效率与写入吞吐。

3.4 Merge-on-Read vs Copy-on-Write

Iceberg V2 提供两种更新模式:

模式写入路径读取路径适用场景
Copy-on-Write(COW)重写整个 DataFile直接读 DataFile写少读多
Merge-on-Read(MOR)写 Delete File + 新 DataFile合并读取写多读少

CDC 场景下,由于变更频率高,我们选择 MOR 模式。读取时通过 read.delete.mode=merge-on-read 配置,Iceberg 会自动将 Delete File 中的记录排除。对于准实时报表,我们额外配置了 Spark 定时 Compaction Job(每小时一次),将 Delete File 合并到 DataFile 中,避免读放大。


在这里插入图片描述

四、Spark Structured Streaming:流式增量计算

数据入湖后,需要在 Iceberg 之上进行增量 ETL,生成 DWD(明细层)和 DWS(汇总层)。Spark Structured Streaming 与 Iceberg 的集成提供了微批(Micro-batch)和连续处理(Continuous Processing)两种模式。

4.1 增量读取原理

Spark 通过 Iceberg 的 IncrementalChangelogScan API 实现增量读取。每次微批启动时,Spark 查询自上次 Checkpoint 以来新增的所有 Snapshot,只读取新增的 DataFile。

val df = spark.readStream
  .format("iceberg")
  .option("stream-from-timestamp", startTimestamp)
  .load("warehouse.iceberg.cdc_order_info")

val query = df.writeStream
  .format("iceberg")
  .outputMode("append")
  .option("checkpointLocation", "/checkpoints/order_dwd")
  .toTable("warehouse.iceberg.dwd_order_event")

4.2 DWD 层构建:事件清洗与打宽

在 DWD 层,我们需要将订单主表与详情表关联,并补充用户维度信息。由于 Iceberg 支持 ACID,我们可以使用 MERGE INTO 实现幂等的 Upsert:

MERGE INTO warehouse.iceberg.dwd_order_event t
USING (SELECT * FROM streaming_batch) s
ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *

关键优化点:

  • Broadcast Hint:维度表(如用户表)仅 200MB,通过 /*+ BROADCAST(dim_user) */ 强制广播,避免 Shuffle;
  • Z-Order 排序:对 DWD 表执行 OPTIMIZE table ZORDER BY (user_id, event_time),将同一用户的数据聚类到相邻文件,极大加速后续用户级聚合查询。

4.3 窗口聚合与 Watermark

DWS 层需要按 5 分钟滚动窗口统计订单金额。Structured Streaming 的 Watermark 机制用于处理乱序数据:

val windowedCounts = df
  .withWatermark("event_time", "10 minutes")
  .groupBy(
    window($"event_time", "5 minutes"),
    $"region"
  )
  .agg(sum($"amount").as("total_amount"))

WaterMark 延迟设为 10 分钟,意味着 10 分钟前的窗口会被触发并写入 Iceberg。由于 Iceberg 的 Snapshot 隔离性,下游查询不会读到未闭合的窗口数据,保证了"读到即完整"的语义。


五、查询加速:分区演进、隐藏分区与文件编排

5.1 分区演进(Partition Evolution)

业务初期,订单表按 days(order_time) 分区即可满足需求。随着数据量增长,我们发现同一分区内的文件过多(每日 10 万+ 文件),查询启动时的文件列表耗时成为瓶颈。

Iceberg 支持分区演进:在不重建表的情况下,修改分区策略使新数据按更细的粒度分区。

-- 原始分区策略:按天
ALTER TABLE order_info ADD PARTITION FIELD hours(order_time);

执行后,历史数据仍按天组织,新写入数据按小时组织。查询引擎根据时间范围自动选择最优的分区粒度进行裁剪。

5.2 隐藏分区(Hidden Partitioning)

传统 Hive 表中,分区字段必须是表中的显式列(如 dt STRING),导致业务 SQL 中充斥 WHERE dt='2025-08-11' 这类与业务无关的过滤条件。

Iceberg 的隐藏分区允许从现有列派生分区,而无需添加冗余列:

CREATE TABLE warehouse.iceberg.order_info (
  order_id BIGINT,
  user_id BIGINT,
  order_time TIMESTAMP,
  amount DECIMAL(16,2)
) USING iceberg
PARTITIONED BY (days(order_time), bucket(16, user_id));

这里 days(order_time) 是隐藏分区,业务查询只需写 WHERE order_time >= '2025-08-01',Iceberg 自动将条件转换为分区过滤。

5.3 文件编排:OPTIMIZE 与 REWRITE DATA

CDC 持续写入会产生大量小文件,严重影响查询性能。我们通过 Spark 定时作业进行文件编排:

-- 合并小文件,目标 128MB
OPTIMIZE warehouse.iceberg.cdc_order_info;

-- Z-Order 重排,加速多维过滤
REWRITE DATA TABLE warehouse.iceberg.dwd_order_event
  USING ZORDER (user_id, product_id);

-- 清理过期 Snapshot,释放存储
VACUUM warehouse.iceberg.dwd_order_event;

生产环境中,我们将上述 SQL 封装为 Airflow DAG,每日凌晨 2:00 执行,将前一天的小文件合并后,查询 P95 耗时从 45 秒降至 3 秒。


在这里插入图片描述

六、查询层集成:Trino / StarRocks 统一查询入口

湖仓的价值最终体现在查询层。我们在 Iceberg 之上搭建了统一的查询网关:

  • Ad-hoc 查询:Trino 连接 Iceberg Catalog,分析师通过 SQL 直接探查原始数据;
  • 高并发报表:StarRocks 3.x 支持 Iceberg 外表查询,通过 Data Cache 将热数据缓存到本地 SSD,QPS 可达 5000+;
  • 湖仓一体加速:对于查询频率极高的 DWS 汇总表,通过 StarRocks 的 CREATE MATERIALIZED VIEW 将 Iceberg 数据异步导入内表,实现亚秒级响应。

StarRocks 查询 Iceberg 的关键配置:

CREATE EXTERNAL RESOURCE iceberg_resource
PROPERTIES (
  "type" = "iceberg",
  "iceberg.catalog.type" = "HIVE",
  "hive.metastore.uris" = "thrift://hive-metastore:9083"
);

CREATE EXTERNAL TABLE ext_order_info (
  order_id BIGINT,
  amount DECIMAL(16,2)
) ENGINE=ICEBERG
PROPERTIES (
  "resource" = "iceberg_resource",
  "database" = "warehouse",
  "table" = "dwd_order_event"
);

StarRocks 的 CBO(Cost-Based Optimizer)会自动将过滤条件下推到 Iceberg,利用 Manifest 层的统计信息跳过不满足条件的文件,实现与原生数仓表接近的查询性能。


七、总结

本文围绕 Apache Iceberg + Spark Structured Streaming 的技术组合,系统阐述了 Lakehouse 实时数仓的构建路径:

  1. 数据入湖:利用 Flink CDC 3.0 实现 MySQL 整库分钟级同步,借助 Iceberg V2 的 MOR 模式支撑高频更新;
  2. 分层计算:Spark Structured Streaming 读取 Iceberg 增量 Snapshot,通过 MERGE INTO 构建 DWD/DWS,Watermark 机制保证窗口完整性;
  3. 查询加速:分区演进、隐藏分区、Z-Order 排序与定时文件编排,层层削减查询 IO;
  4. 统一查询:Trino 负责灵活探查,StarRocks 负责高并发加速,实现"一份数据、多种负载"。

相比传统 Lambda 架构,该方案将数据链路维护成本降低了约 60%,数据一致性达到 Snapshot 隔离级别。随着 Iceberg Spec V4 的推进(列式元数据、更快 Commit),以及 Paimon 在实时更新场景的持续演进,Lakehouse 正在成为实时数仓的事实标准。

参考资料

  • Apache Iceberg 官方文档:Table Spec & Partitioning
  • Flink CDC 3.0 官方文档:Pipeline YAML 配置
  • StarRocks 官方文档:Iceberg 外表查询
  • Dremio Blog: Looking back the last year in Lakehouse OSS (2025)

转载自 CSDN-专业IT技术社区

原文链接:https://blog.csdn.net/DK_Allen/article/details/163656854

文章来源crawl

评论

赞0

评论列表

微信小程序
QQ小程序

关于作者

点赞数:0
关注数:0
粉丝:0
文章:0
关注标签:0
加入于:--