深入理解ApachePaimon湖存储文件操作
Apache Paimon 流式数据湖存储技术解析
以下文章来源于 Apache Paimon,作者齐典

Apache Paimon
Apache Paimon 中文社区官微,PMC 成员维护
摘要
Apache Paimon 作为一项流式数据湖存储技术,在写入和查询过程中包含许多对文件的操作,而这些操作会对 Paimon 底层文件产生影响。熟悉底层文件创建和更新的流程有助于更好地使用和管理 Paimon。本文将通过具体的例子来展示用户执行的 SQL 语句和 Flink 命令如何影响底层文件。通过深入探讨提交快照(snapshot)和小文件合并(compact)的原理,来帮助读者更好地了解底层文件的创建和更新。
前置阅读建议
在进一步阅读本文之前,读者可以在 Paimon 官网中阅读以下文章中的相关概念以更好地理解本文接下来的一系列操作:
- 基础概念 [1]
- 文件布局 [2]
- 如何在 Flink 中使用 Paimon [3]
在本机安装 Flink 1.17 并部署 Paimon 的相关依赖后,使用 ./bin/start-cluster.sh 启动本地 Flink 集群。后续操作如果涉及到更改 flink-conf.yaml 配置文件,则需要先运行 ./bin/stop-cluster.sh 停止集群后重新启动才会生效。
01 创建 Catalog
进入 <FLINK_HOME> 通过 ./bin/sql-client.sh 启动 Flink SQL 客户端,并逐一执行以下语句来创建 Paimon Catalog。Paimon Catalog 可以方便用户管理 Paimon 元数据,以及使用 Table API 和 SQL 语句进行查询。
sql
CREATE CATALOG paimon WITH (
'type' = 'paimon',
'warehouse' = 'file:///tmp/paimon'
);
USE CATALOG paimon;
创建 Paimon Catalog 后,会在本机创建对应的物理目录,可以在 /tmp/paimon 查看新创建的空目录。
02 创建 Paimon 表
继续执行以下建表语句,创建 3 个字段的 Paimon 表:
sql
CREATE TABLE T (
id BIGINT,
a INT,
b STRING,
dt STRING COMMENT 'timestamp string in format yyyyMMdd',
PRIMARY KEY(id, dt) NOT ENFORCED
) PARTITIONED BY (dt);
这将在路径 /tmp/paimon/default.db/T 下创建 Paimon 表 T,它的表结构(schema)存储在 /tmp/paimon/default.db/T/schema/schema-0 中,读者可以自行查看文件内容了解 Paimon 表的元数据。
03 向表中插入记录
在 Flink SQL 中执行下面的 insert 语句,注意该语句不是同步执行完成的,而是会提交一个 Flink 任务到 Flink 集群。如果想查看 Flink 任务拓扑和日志信息,可以在 ./conf/flink-conf.yaml 中添加如下配置:
yaml
rest.port: 8081
insert 语句如下:
sql
INSERT INTO T VALUES (1, 10001, 'varchar00001', '20230501');
在 Flink 任务完成后,一条记录通过一次成功的提交(commit)写入 Paimon 表,对用户可见。读者可以通过运行 SELECT * FROM T; 来验证表内数据。同时,该次成功提交创建了一个快照(snapshot),对应的路径为 /tmp/paimon/default.db/T/snapshot/snapshot-1。以下是该快照对应的文件布局:

snapshot-1 的文件内容记录了快照的元信息,比如 manifest list 和 schema id(具体的文件名是随机的,每次实验可能不同):
json
{
"version": 3,
"id": 1,
"schemaId": 0,
"baseManifestList": "manifest-list-4ccc-c07f-4090-958c-cfe3ce3889e5-0",
"deltaManifestList": "manifest-list-4ccc-c07f-4090-958c-cfe3ce3889e5-1",
"changelogManifestList": null,
"commitUser": "7d758485-981d-4b1a-a0c6-d34c3eb254bf",
"commitIdentifier": 9223372036854775807,
"commitKind": "APPEND",
"timeMillis": 1684155393354,
"logOffsets": {},
"totalRecordCount": 1,
"deltaRecordCount": 1,
"changelogRecordCount": 0,
"watermark": -9223372036854775808
}
manifest 和 manifest list 的概念在前置文章(文件布局)有详细描述。简单来说,manifest list 包括对应快照的所有变更,记录的方式为 基线(baseManifestList)+ 变更(deltaManifestList)。每个 manifest list 又包含了一个或多个 manifest 记录来描述具体对数据文件的操作。以本次写入操作为例,第一次提交创建快照(snapshot-1),也创建了 2 个 manifest 记录,和 1 个 manifest 文件。我们来具体看一下快照内容:
manifest-list-4ccc-c07f-4090-958c-cfe3ce3889e5-0是快照-1 的基线(baseManifestList),对应图-1 中的manifest-list-1-base。因为是首次写入,基线中没有任何内容。manifest-list-4ccc-c07f-4090-958c-cfe3ce3889e5-1是快照-1 的变更 (deltaManifestList),对应图-1 中的manifest-list-1-delta。它包含一个 manifest 记录:manifest-2b833ea4-d7dc-4de0-ae0d-ad76eced75cc-0是快照1 对数据文件的操作记录,对应图-1 中的manifest-1-0。在本例中,它记录了对数据文件的添加操作,体现在图-1 中指向数据文件的 ADD 箭头。
接下来我们尝试一次添加一批记录,对应多个分区。在 Flink SQL 中执行以下语句:
sql
INSERT INTO T VALUES
(2, 10002, 'varchar00002', '20230502'),
(3, 10003, 'varchar00003', '20230503'),
(4, 10004, 'varchar00004', '20230504'),
(5, 10005, 'varchar00005', '20230505'),
(6, 10006, 'varchar00006', '20230506'),
(7, 10007, 'varchar00007', '20230507'),
(8, 10008, 'varchar00008', '20230508'),
(9, 10009, 'varchar00009', '20230509'),
(10, 10010, 'varchar00010', '20230510');
Flink 任务执行完成后,第二次提交成功。读者可以通过 SELECT * FROM T; 验证表内共有 10 条记录。该次提交创建了一个新的快照,snapshot-2。
我们进入表目录递归查看所有文件:
bash
% ls -atR .
./T:
dt=20230501
dt=20230502
dt=20230503
dt=20230504
dt=20230505
dt=20230506
dt=20230507
dt=20230508
dt=20230509
dt=20230510
snapshot
schema
manifest
./T/snapshot:
LATEST
snapshot-2
EARLIEST
snapshot-1
./T/manifest:
manifest-list-9ac2-5e79-4978-a3bc-86c25f1a303f-1 # 快照-2 变更
manifest-list-9ac2-5e79-4978-a3bc-86c25f1a303f-0 # 快照-2 基线
manifest-f1267033-e246-4470-a54c-5c27fdbdd074-0 # 快照-2 manifest 记录
manifest-list-4ccc-c07f-4090-958c-cfe3ce3889e5-1 # 快照-1 变更
manifest-list-4ccc-c07f-4090-958c-cfe3ce3889e5-0 # 快照-1 基线
manifest-2b833ea4-d7dc-4de0-ae0d-ad76eced75cc-0 # 快照-1 manifest 记录
./T/dt=20230501/bucket-0:
data-b75b7381-7c8b-430f-b7e5-a204cb65843c-0.orc
...
# 20230502 - 20230510 每个分区均包含一个数据文件
...
./T/schema:
schema-0
对应的文件布局发生如下改变:

可以看到,snapshot-2 的基线变为 snapshot-1 的全量数据,snapshot-2 的变更为在分区 20230502 至 20230510 的 bucket-0 内分别创建一个数据文件。
04 从表中删除记录
现在我们删除一些分区的数据。假设我们想把 20230503 及之后的分区数据全部删除,在 Flink SQL 中执行如下语句:
sql
DELETE FROM T WHERE dt >= '20230503';
执行完成后,第三次提交成功并且创建了对应的快照,snapshot-3。进入 Paimon 目录(/tmp/paimon/default.db/ 查看发现分区并没有被删除;相反地,20230503 至 20230510 每个分区下面都多了一个数据文件:
bash
./T/dt=20230510/bucket-0:
data-b93f468c-b56f-4a93-adc4-b250b3aa3462-0.orc # DELETE 语句创建的数据文件
data-0fcacc70-a0cb-4976-8c88-73e92769a762-0.orc # INSERT 语句创建的数据文件
这很合理,因为对文件的更新和删除是比较重的操作,虽然逻辑上 20230510 这个分区已经没有数据了,但是为了提高写入性能,Paimon 不会立即对该分区内记录执行合并(Compact),而是新增一条对应的删除操作记录:
- 首次写入时,该分区内新增记录:
+I[10, 10010, 'varchar00010', '20230510'] - 后续删除时,该分区内新增记录:
-D[10, 10010, 'varchar00010', '20230510']
在查询的时候,Paimon 会根据记录的写入顺序计算数据被更新/删除后的结果,所以执行删除操作后,SELECT * FROM T; 会返回以下 2 条记录,没毛病。
sql
+I[1, 10001, 'varchar00001', '20230501']
+I[2, 10002, 'varchar00002', '20230502']
快照-3 对应的文件布局如下:

需要注意的是,manifest 记录(manifest-3-0)中包含 8 个 ADD 操作,对应在分区 20230502 至 20230510 创建的 8 个数据文件(data-delete-0)。
05 表的合并(Compact)
从上面的例子可以看到,不管逻辑上是新增/更新/删除记录,每次执行 SQL 语句都会导致文件数量的增加。Paimon 表中小文件的数量会随着快照次数的增加而增加,这会导致读取性能下降。因此,需要定期执行小文件合并(Compact)以减少小文件的数量。
现在让我们触发全量合并(full-compact)。在 flink-conf.yaml 中添加条目将 Flink 执行模式设置为批处理模式:
yaml
execution.runtime-mode: batch
然后通过 flink run 运行一个全量合并作业,语法如下:
bash
<FLINK_HOME>/bin/flink run
/path/to/paimon-flink-action-0.5-SNAPSHOT.jar
compact
--warehouse <warehouse-path>
--database <database-name>
--table <table-name>
[--partition <partition-name>]
[--catalog-conf <paimon-catalog-conf> [--catalog-conf <paimon-catalog-conf> ...]]
在本例中,进入 <FLINK_HOME> 执行如下命令:
bash
./bin/flink run
./lib/paimon-flink-action-0.5-SNAPSHOT.jar
compact
--path file:///tmp/paimon/default.db/T
执行成功后,Paimon 表中所有的数据文件会被合并,同时产生一个新的快照,snapshot-4,快照元信息如下:
json
{
"version": 3,
"id": 4,
"schemaId": 0,
"baseManifestList": "manifest-list-be16-82e7-4941-8b0a-7ce1c1d0fa6d-0",
"deltaManifestList": "manifest-list-be16-82e7-4941-8b0a-7ce1c1d0fa6d-1",
"changelogManifestList": null,
"commitUser": "a3d951d5-aa0e-4071-a5d4-4c72a4233d48",
"commitIdentifier": 9223372036854775807,
"commitKind": "COMPACT",
"timeMillis": 1684163217960,
"logOffsets": {},
"totalRecordCount": 38,
"deltaRecordCount": 20,
"changelogRecordCount": 0,
"watermark": -9223372036854775808
}
和之前执行写入操作的快照信息有所不同,commitKind 字段是 COMPACT 表示这是一个合并产生的快照。
最新的文件布局如下:

在 manifest-4-0 记录中,出现了对分区数据文件的删除记录。虽然这是个好兆头,意味着表的数据文件不会一直递增下去,但是数据文件并不会被物理删除,而是要等到分区过期(snapshot expire)时通过判断后才会真正被物理上抹去。
我们数一下 manifest-4-0 记录包含的操作数量,共有 20 个操作记录,其中有 18 个 DELETE 操作和 2 个 ADD 操作:
- 分区
20230503至20230510,每个分区有 2 个 DELETE 操作,对应 2 个数据文件的删除。 - 分区
20230501至20230502,每个分区有 1 个 DELETE 操作和 1 个 ADD 操作。
需要注意的是,图中没有表示对分区 20230501 至 20230502 的文件操作,因为它们不会改动文件内容,只会改动数据文件的 level。感兴趣的读者可以进一步了解 LSM 树中的 level 和 sorted run 的概念。
06 变更表结构
如果需要变更 Paimon 表的一些配置信息,可以执行以下 Flink SQL 语句:
sql
ALTER TABLE T SET ('full-compaction.delta-commits' = '1');
对表的变更操作会创建一个新的 schema 文件,下次成功的快照会使用最新的文件。
07 快照过期
上文中我们提到,在 manifest 中被标记删除的记录不会立即被物理删除,而是要等到快照过期阶段判断可以安全删除才能和该快照一起被清理。读者可以在快照过期 [4] 文章中了解更多内容。
现在假设在创建快照-5 的时候执行快照过期,判断需要清理快照-1 到快照-4。过程如下:
- 对于所有被标记删除的数据文件,执行物理删除,记录删除的数据文件对应的分区和 bucket。
- 删除该快照下的所有 changelog 和符合条件的 manifest 文件。
- 删除快照文件,并更新最早的快照(snapshot 目录下的 EARLIEST 记录)。
快照-5 创建成功后,文件布局如下:

可以看到,分区 20230503 到 20230510 被物理删除。
08 Flink 流式写入
最后,我们以生产中常用的 CDC 导入为例,即通过 Flink CDC 流批一体读取 MySQL 全量和增量记录写入 Paimon,来串联以上提到的一系列文件操作。本节内容包括源端 CDC 数据的读取,Paimon 数据的写入和提交,异步小文件合并,和快照过期。
我们先看下 CDC 任务的流程和拓扑:

MySQL CDC Source流批一体读取 MySQL 全量和增量记录,其中SnapshotReader负责读取全量记录,BinlogReader负责读取增量记录。Paimon Sink负责写入数据到 Paimon 表,内部包括CompactManager负责异步触发 bucket 级别的小文件合并。Committer Operator是一个组件,负责最终提交写入使数据可见,以及快照过期清理。
接下来,我们以端到端的视角说明数据是如何写入到文件的。

首先,MySQL CDC Source 读取全增量数据,归一(normalize)后发往下游。

Paimon Sink 首先缓存 CDC 记录到内存中的 LSM 树,在缓存满或者 Flink 快照时 flush 数据到磁盘(比如 OSS、HDFS、或者本地文件系统)。目前为止,没有 manifest 文件和快照被创建。在 Flink checkpoint 触发之前,Paimon Sink 会发送 committable 信息到下游的 Committer Operator,由 Committer Operator 负责后续的提交。

Flink checkpoint 过程中,Committer Operator 获取上游发送的 committable 信息,创建对应的快照和 manifest list。最终 snapshot 包含该时刻表内的全量数据。

后续同步过程中,如果 CompactManager 判断某个 bucket 需要进行小文件合并,会在 Paimon Sink 异步执行小文件合并,并且将合并结果以 committable 信息发往下游。Committer Operator 收到信息后,解析出该次合并更改的文件并创建 Compact 类型的快照和对应的 manifest 文件。在这种情况下,Committer Operator 在一次 Flink checkpoint 中会产生两个快照,一个是数据写入对应的 Append 类型快照,另外一个是小文件合并对应的 Compact 类型快照。如果在 checkpoint 间隔窗口内,没有数据文件被写入 Paimon 表,那么可以通过配置决定是否生成 Append 类型快照。Committer Operator 也负责快照过期,物理删除被标记的数据文件。

end
