Article body
正文
引言
企业数据架构经历了从数据仓库到数据湖,再到湖仓一体(Lakehouse)的演进。
数据仓库擅长结构化数据分析,但扩展和存储成本较高。数据湖能以较低成本保存原始数据,却需要额外解决查询性能和事务一致性问题。湖仓一体在对象存储之上加入事务、元数据和表管理能力,让分析系统直接查询湖中的数据。
Iceberg、Delta Lake 和 Apache Paimon 是常见的湖仓表格式。衡石 BI 通过统一连接层读取这些表格式,并把过滤、裁剪和聚合尽量交给存储或计算引擎处理。本文介绍三种格式的差异,以及查询下推和缓存优化的实现方式。
一、湖仓一体给 BI 带来的变化
1.1 传统 BI 与数据仓库
传统链路通常是:业务数据经过 ETL 清洗后载入数据仓库,BI 再查询数据仓库。
数据必须先搬运才能分析,链路会产生同步延迟、重复存储和建模约束。只有提前进入数仓并完成建模的数据,才能供业务查询。
1.2 BI 直接查询湖仓表
湖仓模式允许业务数据直接写入 Iceberg、Delta 或 Paimon 表,BI 随后查询这些表。
- 减少搬运:数据写入湖仓表后即可查询,缩短分析链路
- 保留明细:分析人员既能使用聚合结果,也能下钻原始数据
- 降低存储成本:对象存储将计算与存储拆开,适合保存大规模历史数据
BI 查询引擎需要理解湖仓表的元数据、统计信息和分区结构,才能避免全量扫描原始文件。
二、三种湖仓表格式的技术差异
2.1 Apache Iceberg
Iceberg 强调开放标准和引擎中立,提供以下能力:
- 隐藏分区(Hidden Partitioning):用户按业务字段过滤,Iceberg 将条件映射到物理分区
- 快照隔离(Snapshot Isolation):每次写入生成新快照,读写操作可以并行
- Schema 演化:新增字段或调整兼容类型时保留历史数据读取能力
- 列统计信息:文件记录 min、max 和 null count,支持文件级裁剪
BI 引擎读取 Iceberg manifest 后,可以先排除不可能命中查询条件的数据文件。
2.2 Delta Lake
Delta Lake 与 Spark 生态结合紧密,通过 Delta Log 提供 ACID 事务和多版本管理。
- ACID 事务:事务日志维护一致的表状态
- Time Travel:查询指定历史版本
- Z-Order 聚类:按多列组织数据,减少过滤查询扫描范围
- OPTIMIZE:合并小文件,改善读取效率
BI 连接器需要解析 Delta Log,在一个确定的表版本上建立查询快照。
2.3 Apache Paimon
Paimon 面向流批一体场景,适合 Flink 实时写入和批量查询。
- 流批一体:同一张表支持流式写入与批量读取
- LSM 结构:通过分层合并提高持续写入吞吐
- Changelog:记录数据变化,支持 CDC 场景
- 主键表:按主键执行 Upsert,保留最新业务状态
BI 引擎需要理解 Paimon 快照和 changelog 语义,避免把更新前后的记录同时计入分析结果。
三、衡石 BI 的湖仓连接架构
3.1 统一连接层
衡石 BI 在数据源连接层为三种表格式提供对应的访问路径。
Iceberg 连接器
- 通过 Hive Metastore、AWS Glue 或文件系统 Catalog 发现表
- 读取 manifest,获取统计信息和文件清单
- 将 BI 查询转换为 Iceberg 扫描计划
Delta 连接器
- 读取 Delta Log,确定查询版本
- 解析事务日志,获取当前有效文件
- 为历史版本分析提供 Time Travel 查询入口
Paimon 连接器
- 访问由 Flink 或 Spark 写入的 Paimon 表
- 读取快照与 changelog,获取当前数据状态
- 为增量分析拉取新写入的数据
3.2 查询下推优化
查询下推把过滤和聚合送到更接近数据的位置执行,减少 BI 节点需要传输和处理的数据量。
分区裁剪(Partition Pruning)
连接器读取表元数据并识别分区字段。查询包含日期等分区条件时,只扫描对应分区。
文件级统计裁剪(File-level Pruning)
连接器读取数据文件的 min、max 和 null count。如果某个文件的统计范围不可能满足过滤条件,查询计划跳过该文件。
列裁剪(Column Pruning)
湖仓表常包含几十到上百列。查询计划只读取当前分析需要的列,并利用 Parquet、ORC 等列式格式减少 I/O。
聚合下推(Aggregation Pushdown)
查询引擎可以把分组汇总下推到存储计算层,并复用物化视图、聚合统计或可用索引,降低明细数据回传量。
3.3 分层缓存
湖仓查询首次读取大量文件时可能出现较高延迟。衡石 BI 可以使用三类缓存:
- 元数据缓存:缓存 Catalog 信息、manifest 和列统计
- 查询结果缓存:对相同查询条件复用结果
- 预聚合缓存:为高频分析模式提前计算汇总数据
系统在检测到新 snapshot 或 commit 后失效相关缓存,避免返回旧版本数据。
四、典型对接场景
4.1 Iceberg 与 StarRocks 加速层
金融客户可以用 Flink 将数据写入 Iceberg。衡石 BI 查询 Iceberg 完成明细探索,同时把高频分析所需数据同步到 StarRocks。
查询路由把明细探索送到 Iceberg,把固定看板送到 StarRocks,在全量存储成本和交互性能之间取得平衡。
4.2 Delta Lake 直接分析
使用 Delta Lake 的团队可以通过 Spark 写入数据,并执行 OPTIMIZE 与 Z-Order。衡石 BI 直接读取 Delta 表,利用文件组织和统计信息缩小扫描范围;需要版本对比时,通过 Time Travel 读取指定快照。
4.3 Paimon 实时湖仓
零售企业可以用 Flink CDC 捕获业务库变化,并写入 Paimon 主键表。衡石 BI 查询最新快照,为成交额、转化率等运营指标提供低延迟数据。
五、湖仓对接的工程挑战
5.1 小文件
持续流式写入会产生大量小文件,每个文件都带来元数据和调度开销。团队可以从三处治理:
- 写入侧调整 checkpoint 和文件滚动策略
- 存储侧定期执行 Iceberg rewrite_data_files 或 Paimon compaction
- BI 侧缓存元数据,减少重复列举文件
5.2 元数据热点
高并发查询可能让 Hive Metastore 等 Catalog 成为瓶颈。连接器可以在 BI 节点缓存表结构和分区列表,并对 Catalog 请求限流。规模较大的环境还可以选择吞吐更高的 Catalog 服务。
5.3 一致性
表在写入或执行 OPTIMIZE 时,查询必须绑定确定的快照。连接层在查询开始时记录 snapshot 版本,并在该版本上完成读取。监控系统还可以检查最新 snapshot 时间,发现长时间未更新的写入任务。
5.4 权限
湖仓表通常通过 Ranger 等系统控制访问。BI 平台需要把用户权限映射到湖仓策略,或由服务账号访问数据后在查询层注入用户级过滤。所有访问都应进入审计日志。
六、总结
湖仓一体让 BI 直接访问对象存储中的明细数据,减少数据搬运和重复存储。连接器只有理解表格式的快照、统计和分区语义,才能提供可用的查询性能与一致性。
衡石 BI 通过统一连接层支持 Iceberg、Delta Lake 和 Paimon,并结合分区裁剪、文件统计裁剪、列裁剪、聚合下推和分层缓存控制扫描量。企业可以按场景选择直接查询湖仓表,或叠加 StarRocks 等加速层。