Skip to content

S1 第 5 期:3 行数据写进去,文件里只有 2 行? ​

Apache Paimon 源码学习

作者 X老师(DaemonforY),Paimon 2.0 / master 源码,按 CC BY-NC-SA 4.0 发布。配套实验和代码在 GitHub。

Yui和Kai展示写入合并与文件落盘过程
Yui和Kai展示写入合并与文件落盘过程AI 生成配图

TL;DR:Paimon 主键表的写入先进内存里的写缓冲,不会立刻生成文件;提交前刷盘时,缓冲里的数据按“主键 + 写入顺序号”排好序,同一主键的多条记录交给合并引擎合成一条再写进 L0 文件。所以一次写 3 行、其中两行主键相同,文件里就只有 2 条。这种合并只发生在同一个缓冲内——分两次提交,旧版本就会留在磁盘上,要靠读取时归并或 compaction 去掉。

实验准备 ​

一张最简单的主键表,1 个桶:

sql
CREATE TABLE wb_one (
  order_id BIGINT,
  status   STRING,
  PRIMARY KEY (order_id) NOT ENFORCED
) WITH ('bucket' = '1');

INSERT INTO wb_one VALUES (1, 'A'), (2, 'B'), (1, 'C');

订单 1 在同一条 INSERT 里出现了两次:先 A,后 C。

完整 SQL 在仓库 labs/sql/lab04/compare-write-buffer.sql,一条命令复现:cd labs && ./run.sh sql/lab04/compare-write-buffer.sql

现象:写了 3 行,文件里只有 2 条 ​

一条 INSERT 写 3 行,数据文件里只有 2 条记录

(节选自 labs/logs/ep5-flink2.2/compare.log,省略了中间的 INFO 日志和异常堆栈):

Flink SQL> SELECT snapshot_id, commit_kind, total_record_count, delta_record_count FROM `wb_one$snapshots`;
+----------------------+--------------------------------+----------------------+----------------------+
|          snapshot_id |                    commit_kind |   total_record_count |   delta_record_count |
+----------------------+--------------------------------+----------------------+----------------------+
|                    1 |                         APPEND |                    2 |                    2 |
+----------------------+--------------------------------+----------------------+----------------------+
1 row in set

Flink SQL> SELECT level, record_count, min_sequence_number, max_sequence_number FROM `wb_one$files`;
+-------------+----------------------+----------------------+----------------------+
|       level |         record_count |  min_sequence_number |  max_sequence_number |
+-------------+----------------------+----------------------+----------------------+
|           0 |                    2 |                    1 |                    2 |
+-------------+----------------------+----------------------+----------------------+
1 row in set
  • 只有 1 个数据文件,里面 2 条记录;
  • sequence 范围是 1 ~ 2——seq 0 那条去哪了?

再用 audit_log 读一下快照 1(这张表唯一的历史版本):

sql
SELECT * FROM `wb_one$audit_log` /*+ OPTIONS('scan.snapshot-id' = '1') */;
+--------------------------------+----------------------+--------------------------------+
|                        rowkind |             order_id |                         status |
+--------------------------------+----------------------+--------------------------------+
|                             +I |                    1 |                              C |
|                             +I |                    2 |                              B |
+--------------------------------+----------------------+--------------------------------+
2 rows in set

订单 1 的 'A' 不在任何文件、任何快照里,时间旅行也查不到它。它从来没有落过盘。

原因:写缓冲在刷盘时按主键合并 ​

写缓冲按主键排序并合并记录后写入文件
写缓冲按主键排序并合并记录后写入文件AI 生成配图

3 行数据在写缓冲里经历了什么

主键表每个桶有一个 MergeTreeWriter,写入分两步:

第 1 步:写入只进内存。 MergeTreeWriter.write(第 166~168 行)给每条记录分配一个递增的 sequence number,然后 writeBuffer.put(...) 放进内存排序缓冲。此时磁盘上一个数据文件都没有。

写入顺序主键值seq
1订单 1A0
2订单 2B1
3订单 1C2

第 2 步:提交前刷盘。 批作业在输入结束时、流作业在每次 checkpoint 前,都会调用 prepareCommit → flushWriteBuffer(第 254~255 行)。刷盘时(第 226 行 writeBuffer.forEach):

  1. 缓冲按主键升序、再按 seq 升序排序(SortBufferWriteBuffer 第 92~93 行把 sequence 字段追加到排序字段末尾),同一主键的记录挨在一起:(1,A,0)、(1,C,2)、(2,B,1);
  2. 同一主键的一组记录交给表的合并引擎。默认的 DeduplicateMergeFunction 每来一条就覆盖上一条(第 54 行 latestKv = kv),留下的就是 seq 最大的那条;
  3. 合并结果写进一个 L0 文件:(1,C,2)、(2,B,1)——正好是 $files 里的 2 条、seq 1~2。

想亲眼看这个过程,可以按实验 4 的断点清单在 DeduplicateMergeFunction.java:59 return latestKv; 打断点:订单 1 命中时 latestKv 的值是 C、seq 是 2。

对照:同样 3 行,分两次写 ​

sql
INSERT INTO wb_two VALUES (1, 'A'), (2, 'B');
INSERT INTO wb_two VALUES (1, 'C');

一次写入 vs 分两次写入

(节选自 labs/logs/ep5-flink2.2/compare.log,省略了中间的 INFO 日志和异常堆栈):

Flink SQL> SELECT snapshot_id, commit_kind, total_record_count, delta_record_count FROM `wb_two$snapshots`;
+----------------------+--------------------------------+----------------------+----------------------+
|          snapshot_id |                    commit_kind |   total_record_count |   delta_record_count |
+----------------------+--------------------------------+----------------------+----------------------+
|                    1 |                         APPEND |                    2 |                    2 |
|                    2 |                         APPEND |                    3 |                    1 |
+----------------------+--------------------------------+----------------------+----------------------+
2 rows in set

Flink SQL> SELECT level, record_count, min_sequence_number, max_sequence_number FROM `wb_two$files`;
+-------------+----------------------+----------------------+----------------------+
|       level |         record_count |  min_sequence_number |  max_sequence_number |
+-------------+----------------------+----------------------+----------------------+
|           0 |                    1 |                    2 |                    2 |
|           0 |                    2 |                    0 |                    1 |
+-------------+----------------------+----------------------+----------------------+
2 rows in set

这次是 2 个文件、3 条记录:第一次提交刷出 seq 0~1,第二次提交只有 seq 2。两个缓冲各自刷盘,订单 1 的 A 和 C 不在同一个缓冲里,谁也合并不了谁。

但 SELECT * FROM wb_two 的结果和 wb_one 一模一样:订单 1 = C,订单 2 = B。区别只在磁盘上:旧版本 A 还在,要等读取时归并(每次查询都做一遍)或 compaction(重写一次,以后都不用再做)去掉——这正是第 4 期讲的内容。

所以结论要说准确:写入时合并只发生在同一个写缓冲内。一次提交里同一主键写了很多次(比如一个 checkpoint 周期里同一个用户更新了 100 次状态),落盘时只会剩 1 条;跨提交的多个版本,写缓冲管不了。

“合并”不等于“丢掉一行” ​

合并引擎可聚合数值而非简单丢弃旧记录
合并引擎可聚合数值而非简单丢弃旧记录AI 生成配图

合并方式取决于合并引擎

刷盘时调用的是表配置的合并引擎,不一定是“留最新”。换成聚合引擎试试:

sql
CREATE TABLE wb_agg (
  order_id BIGINT,
  amount   INT,
  PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
  'bucket' = '1',
  'merge-engine' = 'aggregation',
  'fields.amount.aggregate-function' = 'sum'
);

INSERT INTO wb_agg VALUES (1, 10), (2, 20), (1, 5);

数据文件同样只有 2 条记录,查询结果是订单 1 = 15、订单 2 = 20——两行没有“丢一行”,而是在刷盘时加在了一起。

那如果下游确实需要每一条原始变更呢?给表加上 'changelog-producer' = 'input':

(节选自 labs/logs/ep5-flink2.2/compare.log,省略了中间的 INFO 日志和异常堆栈):

Flink SQL> SELECT level, record_count FROM `wb_input$files`;
+-------------+----------------------+
|       level |         record_count |
+-------------+----------------------+
|           0 |                    2 |
+-------------+----------------------+
1 row in set

Flink SQL> SELECT * FROM `wb_input$audit_log` /*+ OPTIONS('incremental-between' = '0,1', 'incremental-between-scan-mode' = 'changelog') */;
+--------------------------------+----------------------+--------------------------------+
|                        rowkind |             order_id |                         status |
+--------------------------------+----------------------+--------------------------------+
|                             +I |                    1 |                              A |
|                             +I |                    1 |                              C |
|                             +I |                    2 |                              B |
+--------------------------------+----------------------+--------------------------------+
3 rows in set

$ echo '=== wb_input 的文件 ===' && find ${warehouse_dir}/default.db/wb_input -name '*.parquet' | sed 's|.*/wb_input/||' | sort
=== wb_input 的文件 ===
bucket-0/changelog-b529b708-6fbe-40f4-90d4-3d51926ec55d-0.parquet
bucket-0/data-b529b708-6fbe-40f4-90d4-3d51926ec55d-1.parquet

(文件名中的 UUID 每次运行都不同。)

表目录下多了一个 changelog-*.parquet:数据文件还是 2 条,changelog 文件原样保留 3 条。源码里,刷盘遍历缓冲时每读出一条原始记录,都会先交给 rawConsumer(SortBufferWriteBuffer 第 324~325 行)写进 changelog,然后才进入合并(MergeTreeWriter 第 218~221 行只有 changelog-producer = input 时才创建这个 changelog writer)。

写缓冲装不下怎么办? ​

上面的数据只有 3 行,缓冲当然装得下。如果一次提交写入的数据比缓冲大呢?做个极端实验:50000 行,只有 10 个不同的主键,写缓冲压到最小的 256 kb,并关闭自动 compaction(write-only)以免干扰。

sql
-- 两张表只差一个参数
'write-buffer-size' = '256 kb', 'page-size' = '64 kb', 'write-only' = 'true'
-- wb_nospill 额外加上:
'write-buffer-spillable' = 'false'

INSERT INTO wb_spill   SELECT id % 10, id FROM gen;   -- gen:id 从 1 到 50000
INSERT INTO wb_nospill SELECT id % 10, id FROM gen;

溢写 vs 提前刷盘

(节选自 labs/logs/ep5-flink2.2/buffer-full.log,两条查询分别输出,省略了中间的 INFO 日志和异常堆栈):

Flink SQL> SELECT 'spill' AS t, COUNT(*) AS files, SUM(record_count) AS records FROM `wb_spill$files`;
+--------------------------------+----------------------+----------------------+
|                              t |                files |              records |
+--------------------------------+----------------------+----------------------+
|                          spill |                    1 |                   10 |
+--------------------------------+----------------------+----------------------+
1 row in set

Flink SQL> SELECT 'nospill' AS t, COUNT(*) AS files, SUM(record_count) AS records FROM `wb_nospill$files`;
+--------------------------------+----------------------+----------------------+
|                              t |                files |              records |
+--------------------------------+----------------------+----------------------+
|                        nospill |                   20 |                  200 |
+--------------------------------+----------------------+----------------------+
1 row in set
  • 默认(可溢写):内存满了,BinaryExternalSortBuffer 把排好序的数据溢写到本地临时磁盘,提交前再把内存和磁盘上的有序数据统一归并。结果还是 1 个文件、10 条记录,5 万行在写入时就合并完了。
  • 关闭溢写:writeBuffer.put 返回 false,MergeTreeWriter 只能提前刷盘(第 168~171 行),每满一次写一个 L0 文件。从每个文件的 seq 范围(2510~2519、5030~5039……)能看出缓冲每装满 2520 行就刷一次,一共 20 个文件,每个 10 条——每个文件只能在自己这批数据里合并。

两张表的查询结果完全相同(10 行,k=1 的值是 49991,即最后写入的那条),差别还是在磁盘上:关闭溢写时,每个主键都在 20 个文件里各有一个版本。

溢写默认开启,磁盘用量上限 write-buffer-spill.max-disk-size 默认不限。生产里如果看到一次提交就产生大量 L0 小文件,可以先检查这两个参数和 write-buffer-size。

同一主键的旧版本,在三个时刻被“合并掉” ​

三个合并时刻

时刻合并的范围本系列
写缓冲刷盘同一次提交、同一个桶、同一个缓冲内的记录本期
读取时归并查询涉及的所有文件,每次查询都做第 4 期
compaction把多个文件重写成一个,旧版本从最新快照中消失第 4 期

生产启示 ​

  1. checkpoint 间隔影响写入时的合并效果:流作业每次 checkpoint 都会刷盘提交。间隔越长,同一主键在一个缓冲里被合并的机会越多,落盘的记录越少;代价是数据可见延迟变长。
  2. 想要每条原始变更,用 changelog-producer = input:数据文件照常合并,原始记录写进 changelog 文件给流读下游用。
  3. 不要关闭写缓冲溢写,除非你确认内存足够:关闭后缓冲一满就提前刷盘,小文件变多,合并压力变大。

下期预告:Paimon 时间旅行——查一小时前的数据,到底读的是哪些文件?

自测题:流作业里,同一个订单在两次 checkpoint 之间更新了 3 次,又在下一个 checkpoint 周期里更新了 2 次。如果中间没有发生 compaction,磁盘上这个订单有几条记录?(答案:2 条。每个 checkpoint 刷盘一次,各自合并成 1 条;跨 checkpoint 的两条要等读取时归并或 compaction。前提是缓冲没有提前刷盘。)


代码示例在页面里运行时使用 HiveGPT 的模型接口。延伸阅读来自 JavaGuide(Apache-2.0),版权归原作者。