湖仓存储系统设计剖析和性能优化
编者荐语:数据湖相关资料
以下文章来源于 DataFunTalk,作者毕岩

DataFunTalk 专注于大数据、人工智能技术应用的分享与交流。致力于成就百万数据科学家。定期组织技术分享直播,并整理大数据、推荐/搜索算法、广告算法、NLP 自然语言处理算法、智能风控、自动驾驶、机器学习/深度学习等技术应用文章。
导读
本文将从元数据、Merge-On-Read 等核心设计方面来剖析湖仓存储系统,对比当前流行的三种格式 Delta Lake、Hudi、Iceberg 的设计差异和优缺点,并分享我们在读写性能优化方面所做的思考和实际优化案例,主要包括以下内容:
全文目录
- 湖仓系统
- 核心设计
- 性能优化
分享嘉宾:毕岩 阿里巴巴开源大数据平台 技术专家
编辑整理:蒋长强 平安银行
出品社区:DataFun
01 湖仓系统:阿里云 EMR 湖仓系统

相较于传统的数仓、数据湖来讲,湖仓系统是一种新的数据管理系统。上图展示了阿里云 EMR 湖仓系统的整体架构,它是围绕着 Delta Lake、Iceberg、Hudi 等开源数据湖格式构建的,它同时具备数仓的高性能和数据湖的低成本、开放性。这些数据湖格式基于开源的 Parquet 和 ORC 构建,能够在 AWS S3、阿里 OSS 等低成本存储系统上运行,它还具备 ACID 事务、批流一体以及 Upsert 等能力,可以对接多种商业或开源的查询计算引擎。这些能力使得湖仓体系逐步成为了一种趋势。
湖仓系统有一定的学习成本,比如合理配置、小文件、清理策略、性能调优等等。下面将从湖仓系统设计上入手,了解三种格式的差异。以 Spark 计算引擎为例,去分析读写计算过程中的一些主要的链路和影响性能的关键点。
02 核心设计
Delta Lake、Iceberg、Hudi 三个数据湖格式在功能、特性、支持程度上基本一致,但是在具体设计上各有利弊和权衡,这些设计形成了支撑湖格式特性的基石,下文将主要分析元数据、MOR 读取这两块核心设计。
1. 元数据
元数据由 schema、配置、有效的数据文件列表三个主要部分构成。传统数仓系统有单独服务来管理原数据和事务,三个数据湖格式都是将自己的元数据以自定义的数据结构持久化到了文件系统中,放置在表的路径下,但又和表数据分开存储。Hive 表路径或分区路径下的所有数据文件都是有效的,而数据湖格式引入了多版本的概念,所以当前版本的有效数据文件列表需要从元数据中挑选出来。三个湖格式都封装了自身元数据的加载和更新的能力,这些可以方便的嵌入到不同的引擎,由各个引擎 Plan 和 Execute 自己的查询。
以 Spark 为例来看看 Delta Lake 的元数据设计,它的元数据算是三个系统中最简洁的:每次对 Delta Lake 的写操作,或者添加字段等 DDL 操作,会生成一个新版本的 json deltalog 文件,这里面会记录元数据的变更,包含一些 schema 的配置和 file 的信息,多次 commit 之后,会自动产生一个 checkpoint 的 parquet 文件,这个 parquet 文件会包括前面所有版本的元数据信息,用于优化查询加载。

Delta Lake 元数据加载流程:
- 定位最新的 checkpoint 元数据文件
- List 后面的 delta log json 文件
- 按版本号依次解析,得到表的 schema、配置和有效数据文件列表
Iceberg 也有一个统一的元数据集,与 Delta Lake 不同的是,Iceberg 是三层的架构。

其中 metadata 文件很像 Delta Lake 的 checkpoint 文件,包含了全部的信息,但是不同的是,metadata 文件还包括了前几个快照的信息,并且 Iceberg 是三层架构,其 manifest file 能够对局部的数据文件做统计信息收集,因此也能用于分区之下、文件之上的裁剪。
Iceberg 元数据加载流程:
- 定位到当前 metadata 文件,得到表的 schema 和配置,和当前数据文件快照 snapshot 的 manifest list 文件
- 解析 manifest list 文件,得到一组 manifest 文件
- 解析 manifest 文件,得到有效数据文件列表
Hudi 和前面两个很不一样,其一是它没有统一的元数据结构,其二 Hudi 会对数据文件进行分组,并对文件名进行编码,这是 Hudi 特有的 file group 概念。Hudi 的数据必须有主键,主键可以映射到一个 file group,后续对于这些主键所有的更新,都会写到这个 file group,直到显式的调用修改表的文件布局。也就是说,一个 file group 随着多次的 commit,会产生多个版本。获取当前有效数据文件列表时,会先列出当前分区下的所有文件,按照 file group 分组,取出每个 group 的最新文件,再按照 timeline 筛选掉已经被删除的 group,最后得到一个有效的文件列表。

Hudi 元数据加载流程:
- 解析 hoodie.Properties 得到表的 schema 和配置
- 获取有效文件列表
- 未开启 metadata:List filesystem + timeline
- 开启 metadata:读取 metadata 表
Delta Lake、Iceberg、Hudi 元数据对比如下:

2. Merge-On-Read(MOR)
在一般情况下,在更新一个 Copy-On-Write 表时,即使我们只想执行一条更新操作,也需要将所有涉及到的数据文件加载进来,然后应用更新表达式,再将所有数据一起写出,这里就包含一些没有更新的数据,这就是写放大的现象。为了解决写放大现象,三个数据湖格式中 Hudi 第一个实现了 Merge-On-Read 表。
MOR 的设计思想是,只持久化需要写出的数据,再通过某种方式标识出来,原来的数据文件里的数据成为过期数据,读的时候进行合并;为了提高效率,会定期进行合并(Compaction),通常是按照 Copy-On-Write 的方式写一遍。不同的 MOR 实现的写入、合并策略会有所不同。
Hudi 定义了一个 filegroup 概念,每个 group 包括最多 1 个原始数据文件和多个日志文件,数据文件运行不存在。

数据通过主键映射到 filegroup,如果更新数据将会追加写到映射到 filegroup 内的日志文件,如果是删除,则只需要做一个主键记录,在合并的时候,首先读取原始的数据,然后按照这部分数据的主键去判断在增量日志中有没有相关的记录,如果存在就做合并。
Delta Lake 通过 Deletion Vector 的设计解决写放大的问题。Delta Lake 将需要更新、删除的数据在原数据文件中的 offset 标识出来,写入一个辅助文件。Iceberg V2 表有两个 MOR 的实现,其中基于 position 的设计和 Delta Lake 的 DV 是基本一样,仅在具体实现上有些区别。DV 的写入仅两步:1)根据 update 或者 delete 的 condition,找到文件中匹配的记录,记录他们在文件中的 offset。持久化到一个 bin 文件中;2)将更新后的数据,写到普通的一个新文件中。
相较于 Hudi 允许存在多个日志文件,Delta Lake 在查询性能做了权衡,一个普通数据文件只允许伴随至多一个 DV 文件,当对一个已存在 DV 文件的数据文件再做一个更新的时候,最终写出时会把两个 DV 合并的。由于 DV 文件中 offset 信息是通过位图(RoaringBitMap)来保存的,合并操作是比较高效的。另外 Hudi 是将更新的数据也写入日志文件,Delta Lake 是直接写入普通的 parquet 文件,然后在 bin 文件中做一个标记。以及 hudi 日志文件是行存格式,Delta Lake 的 DV 采用自定义的格式,而数据使用的 parquet 的列存。

Delta Lake 在 DV 下的查询方式我们可以直接看 LogicalPlan,会更加清晰,即将原本的 DeltaScan 转换成 Project + Filter + Delta Scan 的组合。在 Parquet Scan 某个数据文件时,追加了 _skip_row 的辅助字段,上层应用 _skip_row = false 的过滤,然后通过 Project 的投影保证仅了无辅助字段的额外输出。
显然核心就是 _skip_row 的标记,Delta Lake 自定义了 ParquetFileFormat,在读取 parquet 文件后,对每个数据判断 roaringBitMap 是否包含该 offset,有标记为 true,没有就是 false。这样就完成了 Deletion Vector 模式下 Delta Lake 表的查询。
Delta Lake、Hudi、Iceberg 的 MOR 实现对比:
基于 offset 或者 position 的 MOR 实现,由于会通过扫描文件来确定位置,因此写性能上会慢于 iceberg 的 equality 或 hudi mor 的实现,而由于该方案不需要类似 hash join 的读时合并策略,查询性能会好一些。

03 性能优化
1. 查询

以 Spark 为例,一个完整的 query 链路如下:
其中与数据湖相关的有三点:元数据加载、优化 plan、Table Scan。
(1)元数据加载
元数据加载包括获取 schema、构造 fileindex 等,分为单点加载和分布式加载两种。单点加载的代表有 Iceberg、Delta Lake-Standalone,适用于小表,有内存压力。分布式加载的实现有 Hudi、Delta Lake,适用于大表,需提交 Spark Job。
这里我们给出 LHBench 做了一个测试结果:使用的是 TPC-H 的 store-sales 表,设置的 filesize 都是 10MB,单个图表内表的文件数量从 1 千到 20W,三个图表对应的仅读取一行,读取一个分区,和普通字段作为过滤条件的三个场景。从趋势上三个场景是一致的。我们分析第一个,最左侧小表 Iceberg 的单点模式要稍好些,但与 Delta Lake 的分布式元数据加载差距不大。最右侧单点方式整个 Query 的执行时间中蓝色执行部分基本一致,很明显被红色 startup 部分限制,这部分在 plan 是一致的情况下可等同于元数据加载的时间。可见,如何智能的选择合适的加载模式,是一个可选的优化方向。

以下是两个 EMR 的优化案例。
案例 1、EMR Manifest,是一个无服务化的优化元数据加载方案。
阿里云 EMR 一客户的核心 ODS 表,通过 Spark Streaming 写入 Delta Lake,每天增量数据 3TB,目前全表 2.2PB,1500 万个数据文件,仅元数据 10GB。在正常情况下使用 Hive/Presto 查询,使用的 Delta-Standalone 的单机加载方式,会完全卡住或者需要超高内存。

该优化方案是将数据文件的元数据按照分区结构提前持久化到一个 manifest 文件,同时记录 manifest 的元数据版本。在用户查询的时候,根据 filter 做分区裁剪,直接去读分区下面的 manifest 文件,解析出本次查询的有效数据文件,跳过了所有元数据加载的步骤。如果 emr_manifest 的版本有滞后,我们也会拿到滞后的元数据,合并得到正确的数据快照。另外 manifest 文件中还会保存一些 size、stats 这些信息,会应用于一些文件级别的 data-skipping 优化。该方案在该表体量仅为 300TB 时提供,当时需要 10GB 内存 90s 加载完整的元数据,优化后可以实现秒级返回,且内存不需要额外调整。
案例 2、EMR DataLake Metastore,有服务的中心化的优化元数据加载方案。
这里不得不和 HMS 做一个对比,HMS 仅存储表的分区级信息,查询普通表时根据分区位置去 list 路径,拿到有效的数据文件列表。而数据湖格式具备多版本概念,所以针对湖格式的 Metastore 设计必须到文件级别。另外 Hudi、Iceberg 基于分布式锁来实现事务性,而 Delta 则是基于所在文件系统的原子性和持久性,在某些场景下无法提供更强的一致性保障,这也是我们实现 DataLake Metastore 的另外一个原因。
EMR DataLake Metastore 目前已经支持了 Hudi 格式,完全兼容社区版本和元数据协议,同时支持元数据双写。在提供了文件快照和事务的同时,后续也将继续拓展 data profiling 和行级索引为查询提供更精准的裁剪优化。
(2)优化 plan

数据湖的第二点优化是 plan 优化,在前面的 Spark 查询例子里,调用 Spark sql 内核结合各种的状态统计信息去优化 logic plan。比如根据运行时的一些统计信息,通过 AQE 去改变 join 的方式,基于表的静态信息和一些 cost 模型来优化的 CBO 等。
统计信息只有表级别和 column 级别两类,其中表级别有 statistics 信息包括 bytes、rows 等信息,字段级别有 min、max 等。

对于一个简单的 join,仅使用表级别的 rowcount 信息就能够去调整 join 顺序,使得整体的 plan 过程中需要 join 的数据量达到最少。
但是如果 sql 带有 aggregate 或者 filter,就需要结合列的信息去估算整个 scan 的 cost。
关于通过统计指标来优化查询,有几下几点思考:
- 如何更好的利用 statistics 来优化查询。比如 Delta Lake 中 count 整表的 sql 可以直接通过元数据层面的统计信息得到 count 值,将原本的 aggregate 操作转换成 LocalRelation,避免 Scan 全表。
- 如何高效收集/实时更新 statistics。
- 如何打通湖格式自身和 Spark 需要的 statistics。spark 会从 table properties 获取统计信息,而一般情况湖格式的统计信息都是记录到自己的元数据中,这部分需要更好的打通和复用。
(3)Table Scan
在完整 Logical Plan 的优化并转化为 Phyiscal Plan 后,下一步就要执行具体的数据文件读取,即 Table Scan 阶段。
这里小文件对查询的影响是大家所熟悉的。三个湖格式都提供了相应的合并小文件的功能,关键问题其一是目标文件设置多大是合适的,其二合并操作执行的时机和方式,这点我们后续再展开。默认情况下 parquet 在各引擎都是 128M 左右,但这是否是最优的还要看表的整体规模。过多的文件,将会导致元数据规模变大,影响元数据的读写,databricks 和阿里云 emr 都提供了自动调教 file size 的功能。这里给出 databricks 的设置文件大小和表规模的映射关系。为了控制元数据的规模,对于 10TB 的表建议文件大小设置为 1G。

在得到表的全部有效数据文件后,我们还是要根据查询条件来尽可能的进行裁剪,以减少最终读取 Parquet 或者 ORC 文件的体量。Table scan 优化分为两类:一是元数据裁剪,二是行级索引。元数据裁剪又可以分为分区裁剪、Manifest 裁剪(Iceberg 特有)、File 裁剪等不同粒度。

三个数据湖格式都支持 Z-Order 和 DataSkipping:先将 min-max 信息写到元数据当中,然后结合查询过滤条件去实现过滤。
针对点查,传统数据库更多是通过行级索引来实现。Databricks 商业版 Delta Lake 支持 BloomFilter,使用单独的目录文件,保存数据文件的索引信息;Hudi 也支持 Multi-Modal Index 多模索引。EMR 也计划在 DataLake Metastore 中嵌入 Index System 来支持点查加速。
最后就是实际的 Reader 来读取数据文件了。以 Parquet 为例,各引擎集成 parquet 之后,读写性能已经非常不错。但是有一些具体的湖格式场景会关闭一些优化参数,相关的比如谓词下推,向量化等。另外文件的压缩格式和压缩比也会影响文件的加载,以及目前一些 Native 框架支持的 Native Parquet Reader(如 Arrow,Presto,Velox 等)。
2. 写入
以一个带 where 条件的 update SQL 为例,首先要能找到匹配这些查询条件的所在的数据文件,然后加载文件,对匹配到的数据应用表达式完成更新,并最终写出。在完成提交后,也许还要执行一些表的 table service,包括 clean、checkpoint、compaction 等。

优化写操作,选择适合当前场景的表类型(如 MOR 表),并配置合适的参数。在湖格式场景下,这种后置的追加在湖格式写入后面的 table service 其实是逐渐成为了影响效率或者整体作业稳定性的关键性因素。要解决这样的问题,比较通用的做法就是拆成两个链路:一个链路就是正常写入,另外一个链路是用离线的方式通过调度一个作业来执行 table service 任务,包括像 clean、Hudi 表管理等等。
在阿里云 EMR 场景下,有些同时使用 Flink 和 Spark 引擎的客户,会采用 Flink Hudi 入湖和 Spark 离线管理的方案。另外 EMR 也提供了中心化的自动化湖表管理,是结合阿里云的 DLF 服务实现的。如图所示,湖格式自动将 commit 信息同步到 DLF 湖表服务,该服务将根据预定义的策略和实际表的状态及配置判断是否要产生湖表管理类的操作。

这里的核心其实是采取怎样的策略来更好的管理湖表。EMR 目前上线了两大类 5 个策略。在同一类中,又设置优先级避免多条策略的命中而造成计算资源的浪费。举个例子,大多数入湖的表是以时间分区的,DLF 湖表管理能够判断到新分区的到达,而在第一时间对已完成写入的分区执行小文件合并或者 zorder 的操作,以便快速提供这部分数据的高效查询。

end
