ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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:mysqlhostname:mysql-primary.db.svcport:3306username:cdc_userpassword:${CDC_PASSWORD}tables:order_db.order_info,order_db.order_detailserver-id:5400-5404sink:type:icebergcatalog:type:hadoopwarehouse:hdfs://namenode:8020/warehouse/icebergtable-prefix:cdc_table-defaults:format-version:2write.metadata.metrics.default:counts# 默认收集列统计write.distribution-mode:hash# 按主键哈希分布,避免小文件pipeline:parallelism:8schema-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

StreamExecutionEnvironmentenv=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 的IncrementalChangelogScanAPI 实现增量读取。每次微批启动时,Spark 查询自上次 Checkpoint 以来新增的所有 Snapshot,只读取新增的 DataFile。

valdf=spark.readStream.format("iceberg").option("stream-from-timestamp",startTimestamp).load("warehouse.iceberg.cdc_order_info")valquery=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:

MERGEINTOwarehouse.iceberg.dwd_order_event tUSING(SELECT*FROMstreaming_batch)sONt.order_id=s.order_idWHENMATCHEDTHENUPDATESET*WHENNOTMATCHEDTHENINSERT*

关键优化点:

  • 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 机制用于处理乱序数据:

valwindowedCounts=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 支持分区演进:在不重建表的情况下,修改分区策略使新数据按更细的粒度分区。

-- 原始分区策略:按天ALTERTABLEorder_infoADDPARTITIONFIELD hours(order_time);

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

5.2 隐藏分区(Hidden Partitioning)

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

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

CREATETABLEwarehouse.iceberg.order_info(order_idBIGINT,user_idBIGINT,order_timeTIMESTAMP,amountDECIMAL(16,2))USINGiceberg PARTITIONEDBY(days(order_time),bucket(16,user_id));

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

5.3 文件编排:OPTIMIZE 与 REWRITE DATA

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

-- 合并小文件,目标 128MBOPTIMIZEwarehouse.iceberg.cdc_order_info;-- Z-Order 重排,加速多维过滤REWRITEDATATABLEwarehouse.iceberg.dwd_order_eventUSINGZORDER(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 的关键配置:

CREATEEXTERNAL RESOURCE iceberg_resource PROPERTIES("type"="iceberg","iceberg.catalog.type"="HIVE","hive.metastore.uris"="thrift://hive-metastore:9083");CREATEEXTERNALTABLEext_order_info(order_idBIGINT,amountDECIMAL(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)
返回列表