ARTICLE · 1095056
Apache Paimon 内核源码深度剖析——LSM 写入路径、Compaction 调度与 Changelog 生产机制
承接上一篇:第 09 篇我们站在架构师视角,把实时湖仓从需求到上线串成了一条带"闸门"的流水线,其中反复提到四个词:存储统一、Compaction、迟到更新、两阶段提交。本篇不再讲"为什么要这么设计系统",而是钻进 Apache Paimon 的代码,看这四个词在一个 LSM 存储引擎里到底如何落地。 技术基线:Apache Paimon(近年捐赠 ASF、成为顶级项目的湖仓流存储,原名 Flink Table Store)。本文类名与机制以主线版本为准,具体类名随版本演进,落地时对照所用版本源码与官方文档。
1.【问题背景与技术演进】
第 09 篇留下一个判断:"流批一体的本质是存储统一,而不是语言统一"。判断谁来统一?在流批一体架构里,真正承担"既能被 Flink 持续流式写入、又能被 Spark 批量重算、还能被 StarRocks 等分析引擎读取"的组件,就是数据湖表格式。第一代湖表格式(以 Iceberg、Hudi 为代表)主要围绕批量数据文件和快照提交建立;当业务把 CDC 主键表、实时维表、流计算结果集都持续写入湖表时,"实时更新写入"和"批量查询读放大"之间的矛盾就暴露出来。
Apache Paimon 的出现,本质是把 LSM-Tree 这套为"高频更新、读时合并"设计的存储思想,引入开放湖表存储。它不再把主键表简单看成"一个个互相独立的批文件",而是通过 Bucket、Sorted Run 和 Compaction 组织成一棵持续合并的 LSM 树。这个设计选择,决定了它的写入路径、Compaction 调度、Changelog 生产、故障恢复,都和传统以不可变数据文件为核心的湖表存储存在明显差异。
这里需要特别强调一个容易产生误解的地方:Paimon 使用 LSM 模型,并不意味着 Paimon 的物理数据文件就是 RocksDB 意义上的 SST 文件。 当前 Paimon 默认数据文件格式可以是 Parquet,主键数据文件中还会保存用于版本合并的 Row Kind、Sequence Number 等信息。因此,LSM 描述的是 Paimon 的数据组织、写入和合并模型,而不是直接等同于 RocksDB 的文件格式。
为什么今天必须讲源码级?因为第 09 篇里说的"Compaction 积压""小文件膨胀""Changelog 延迟",到了 Paimon 这里不再是抽象指标,而是具体到"哪几个组件、哪个阶段、哪个参数"出了问题。讲不清内核,排障就只能靠重启和赌运气。
2.【读完本文你将获得什么】
- 讲清 Paimon 主键表一条数据从 Sink 进入到生成数据文件、形成 Sorted Run、提交快照的完整源码链路,能在脑子里画出关键类与状态流转;
- 说清 LSM 写路径(内存缓冲 → 排序 → 数据文件 → Sorted Run)和 Compaction 为什么是"写放大换读优化"的权衡;
- 理解 Changelog 四种生产模式(none/input/lookup/full-compaction)各自的代价与适用场景,以及它为什么是"流批一体"成立的关键;
- 掌握主键表、无主键表、Bucket、快照提交(两阶段提交)在内核层面如何工作;
- 一套 Compaction 积压、小文件、Changelog 延迟、快照膨胀的源码视角排障路径;
- 能在面试中把"Paimon 为什么适合实时更新"讲成机制,而不是简单比较谁"更快"。
3.【核心概念与架构定位】

先建立内核层面的整体认知。Paimon 表主要可以分为:主键表(Primary Key Table)和追加表(Append Table)。主键表使用 LSM-Tree,支持按主键进行更新、删除和版本合并;追加表主要面向只追加数据,不需要按主键维护更新状态。二者共享 Snapshot、Manifest 等湖表元数据体系,但写入和查询路径差异很大——这是理解所有问题的前提。
在逻辑分层上,一张 Paimon 主键表可以理解为:表(Table)→ 分区(Partition)→ Bucket → Sorted Run(有序运行集合)→ Data Files。Bucket 是 LSM Tree 的基本组织单元,每个 Bucket 拥有自己的 LSM Tree。固定 Bucket 模式下,数据按照 Bucket Key 哈希到固定 Bucket;当前 Paimon 同时支持 Dynamic Bucket 和 Postpone Bucket,因此不能简单概括为"建表时按 hash 固定、运行期不可变"。写入时按照 Bucket 分布策略落到对应 Bucket,Compaction 和读取主要围绕 Bucket 内的 LSM Tree 展开。
在进程角色上,Paimon 作为 Flink 的 Sink,由 TaskManager 上的写入算子承载。它的内核可以粗分为三层:写入层(接收记录、进入写缓冲、排序并生成数据文件)、Compaction 层(后台把多个 Sorted Run 合并,减少后续读取时的 Merge 工作)、提交与元数据层(基于 Snapshot 和 Manifest 把一次写入原子地对外发布)。它不自己负责整个 Flink 作业的资源调度,而是与 Flink 的 Checkpoint 机制结合完成流式写入提交。

4.【底层原理深度解析】
4.1 写入路径:内存缓冲如何变成 Sorted Run
一条记录进入主键表后,并不是直接修改已经存在的数据文件,而是首先进入写入缓冲。内核里,写入算子维护用于排序和写文件的内存结构,新写入的记录按照主键和相应的版本信息组织。当达到相应的 Flush 条件时,数据会被排序并写成新的数据文件,形成新的 Level 0 Sorted Run。

这里需要纠正一个非常容易出现的源码理解错误:Paimon 主键表的物理数据文件不能直接称为 RocksDB SST 文件。 Paimon 的 Data File 默认可以采用 Parquet 格式,文件中还会保存 Row Kind、Sequence Number 等用于后续 Merge 的信息。更准确的源码链路应该理解为:
Record → Write Buffer → Sort → Data File → Level 0 Sorted Run
而不是:
Record → RocksDB MemTable → SST
数据文件内部有序这一点至关重要:它让后续 Compaction 能够进行高效的多路归并,而不是每次都对全部数据进行无序重排。
为什么要先在内存中排序再落 L0?因为 LSM 的收益建立在 Sorted Run 有序的基础之上。如果落盘文件完全无序,后续 Compaction 和查询都会付出更高代价。代价是:写入时需要缓冲、排序和 Flush,而在 Flink 流式写入场景中,数据的提交可见性还受到 Checkpoint 的影响——因此 Checkpoint 间隔会直接参与流式数据从写入到可见之间的延迟权衡,但不能简单理解成"只有 Checkpoint 才会触发文件 Flush"。
4.2 分层 LSM:为什么读要合并、后台要 Compaction

每个 Bucket 的数据按照 Sorted Run 组织。Level 0 通常包含较新的数据文件,不同 Run 之间可能存在重叠;随着 Compaction,多个 Run 会被合并成新的 Run,从而减少查询时需要处理的数据文件和版本。
查询主键时,最新数据可能分布在多个 Run 中,因此需要读时合并(merge on read)——读取候选数据后,根据 Primary Key、Sequence Number、Row Kind 以及配置的 Merge Engine 得到最终逻辑结果。
读时合并很灵活,但代价是读放大:Sorted Run 越多,查询要打开、读取和归并的数据就越多,点查和扫描都可能受到影响。这就是 Compaction 存在的根本原因:它是一个后台异步任务,把多个 Sorted Run 按策略进行归并,顺便根据 Merge Engine 处理"同主键多版本合并""删除记录处理"等逻辑,让后续读取需要处理的版本数量下降。
需要注意,Compaction 并不意味着永远"不阻塞写入"。正常情况下 Paimon 可以在后台进行 Compaction,但当 Sorted Run 数量达到停止阈值,或者某些模式需要 Compaction 完成后才能提交时,写入可能受到反压或等待;例如 lookup Changelog Producer 默认会等待相应的 Lookup Compaction。
这里有一个关键的设计哲学:Compaction 是用"写放大 + 后台开销"换"读性能"。写入路径把大量合并工作推迟到后台,换来了更好的持续写入能力;代价就是 Compaction 必须跟得上写入速度,否则 Sorted Run 越积越多,读放大、内存压力和文件数量都会上升。第 09 篇说的"Compaction 是主键表运维核心",内核原因就在这里。
4.3 Changelog:为什么 LSM 表还要生产变更流

一个天然的疑问:LSM 文件已经能够通过 Merge Engine 得到"最新结果",为什么 Paimon 还要专门生产 Changelog?因为下游流计算(另一个 Flink 作业、维表 Join、实时特征)需要的不是"当前全量最新状态",而是"状态发生了什么变化"(+I/-U/+U/-D 等 Changelog 类型)。直接读取最终表状态,并不天然包含每一次更新的 before image。
Paimon 通过 changelog-producer 参数解决,主流模式有四种:
- none:不额外生成完整 Changelog 文件,但流式读取仍然可以读取增量变化;下游如果能够处理 Upsert,或者由 Flink 等下游进行 Normalize,则可以正常使用。它不是"只能批读快照",也不是完整 CDC 审计日志;
- input:将上游输入的 Changelog 写入独立的 Changelog 文件并提供给流式读取,要求上游本身已经提供所需的变更语义,例如包含 before/after 的 CDC 记录;Paimon 不负责补齐缺失的 before image;
- lookup:通过 Lookup Compaction 读取已有状态,再应用新数据生成相应的 Changelog。它适用于上游没有 before image、但下游需要完整变更语义的场景,代价是额外的 Lookup、Compaction、缓存和 I/O 开销;
- full-compaction:在 Full Compaction 前后比较表状态差异并生成 Changelog。它描述的是两个 Full Compaction 状态之间的变化,中间发生的多次更新可能被折叠,并不是每一条源事件的完整复制。
这四种模式不是简单的"零成本 vs 语义完整"连续谱,而是不同的消费者契约与计算/存储代价。选型原则:下游能处理 Upsert 用 none;上游已经提供完整 Changelog 用 input;没有 before image、但下游需要完整变化用 lookup;下游可以接受周期性延迟、希望通过 Full Compaction 生成状态差异时用 full-compaction。这也是面试里最能区分"是否真用过"的点。
4.4 快照提交:两阶段提交如何保证精确一次

Paimon 的一次写入不是"写完文件就可见",而是通过快照(Snapshot)发布已经提交的表状态。每个 Checkpoint 周期中,写入算子会准备待提交的数据文件和提交信息;当 Flink Checkpoint 成功完成后,Paimon Writer 的提交链路继续完成 Commit,新的 Snapshot 被发布,对外可见。
这里需要区分三个概念:
Flink Checkpoint ≠ Paimon Snapshot。
Checkpoint 是 Flink 的状态一致性机制;Snapshot 是 Paimon 的表状态元数据。两者通过 Paimon Sink 的提交协议连接起来。当前官方文档明确说明,Streaming Writer 会随着 Checkpoint 完成提交数据,而读取端查询的是已经提交的 Snapshot。
这套机制带来三个重要性质:其一,读写隔离,查询读取的是已经提交的表状态,不会直接看到尚未发布的数据文件;其二,在满足 Flink Source、Checkpoint、Sink 等端到端语义要求的情况下,可以实现精确一次的数据提交语义;其三,时间旅行,保留的 Snapshot 可以支持历史状态查询。但代价是——快照及其引用的数据文件在仍然需要时不能删除,Snapshot 保留越久,存储和元数据管理成本越高。
5.【核心链路与关键流程】

把上面四条线串成一条最重要的路径:一次 Flink 流式写入到 Paimon 主键表并生产 Changelog。
写入算子从 Kafka/Fluss 读到一条 upsert 记录 → 按 Bucket 分布策略路由到目标 Bucket → 进入该 Bucket 当前活跃的写缓冲,按主键及版本信息组织 → 若配置 lookup 模式,则由 Lookup Compaction 读取已有状态并生成相应 Changelog → 达到 Flush 条件后,数据排序并写成新的数据文件,形成 L0 的 Sorted Run → 后台 Compaction 按策略把多个 Sorted Run 归并,顺便根据 Merge Engine 合并多版本、处理删除等语义 → Flink Checkpoint 完成后,Paimon Commit 发布对应的 Snapshot → 下游流读作业从已提交 Snapshot 开始继续读取增量变化。
另一条关键路径是故障恢复:作业失败重启后,Flink 从最近成功 Checkpoint 恢复状态和消费位置,Paimon 继续从已经成功提交的表状态向后写入;尚未进入有效 Snapshot 的文件不会成为正常表状态的一部分,并由 Paimon 后续文件清理机制处理,而不是简单地"任务失败立即删除"。已提交 Snapshot 保持不变,下游不会因为一次失败直接读取到一个半提交状态。理解这条链,就理解了 Paimon 为什么能够在湖表上与 Flink Checkpoint 协作实现端到端一致性。

6.【真实业务场景与工程实战】
以电商"实时用户宽表"为例:上游把用户基础属性、订单、行为三股 CDC 流写入 Kafka/Fluss,一个 Flink 作业把它们按 user_id 聚合写入 Paimon 主键表(部分列更新),下游实时推荐和风控直接读这张表做维表 Join。
工程上有几个关键决策。
其一,Bucket 模式与数量:按写入并行度、数据倾斜、更新频率和文件大小规划。固定 Bucket 需要提前规划,但当前 Paimon 也支持 Dynamic Bucket 和 Postpone Bucket,因此不能再简单认为"Bucket 建表后不可改";如果确实需要调整固定 Bucket 布局,可以通过官方提供的 Rescale Bucket 流程处理。
其二,changelog-producer 选型:本场景如果上游已经提供完整 CDC Changelog,优先评估 input;如果上游没有 before image、而下游需要完整 -U/+U,可以考虑 lookup;如果下游可以接受周期性状态变化,则可以评估 full-compaction。
其三,快照保留:生产环境按"需要回溯多久、流作业恢复窗口、批查询耗时"设置快照保留,例如保留数天支撑数据回溯与恢复,不建议无限保留。
建表时围绕这几点给出关键参数即可(具体取值对照所用版本官方文档):指定主键、Bucket 模式与数量、Bucket Key、合并引擎、Sequence Field、Changelog 模式、快照保留策略、Compaction 相关参数。
这里不堆大段 SQL,因为真正决定成败的是上面的选型判断,而非参数本身。
7.【性能、成本与技术选型】

Paimon 主键表的性能特征,全部能从"LSM 写入、Merge on Read、后台 Compaction"推出来。
写延迟受到写入缓冲、Flush、Checkpoint、Commit 以及 Compaction 等多个环节影响,因此不能简单归因于 Checkpoint 间隔;低延迟场景要综合权衡 Checkpoint 间隔、文件生成频率和提交开销。
读性能取决于 Sorted Run 数量、文件布局、Merge on Read 以及 Compaction 是否跟得上,Compaction 积压时读延迟可能明显恶化;
存储成本由数据文件、Changelog 文件、Compaction 重写、Snapshot 保留等因素共同决定。
什么时候该用主键表:高频 CDC 主键更新、实时维表、需要按主键维护状态的流结果。
什么时候不该用:纯追加的海量日志,可以优先考虑 Append Table,避免不必要的主键 Merge 语义。在与 Iceberg/Hudi 的分工上:以 Flink 实时更新、主键状态维护和流式增量消费为主,Paimon 的 LSM 模型更契合;纯离线大规模批分析、变更稀少等场景则应结合具体查询引擎、生态和数据管理需求评估 Iceberg/Hudi 等方案。所有吞吐、压缩比、目标文件大小的具体数值都是环境强相关的经验起点,必须以自身硬件、报文和并发压测为准,不能照搬。
8.【生产问题与排障方法】
故障一:Compaction 积压、小文件爆炸。 现象:Bucket 下 Sorted Run / 数据文件数量持续上涨,查询延迟明显增加。按"现象→指标→原理→定位→修复→预防":先看每个 Bucket 的 Sorted Run 数量和 Compaction 是否在正常推进;原理上这是写入速度超过了后台合并速度;常见根因是 Compaction 资源不足、Compaction 任务报错、Sorted Run 触发阈值配置不合理,或者存在大量 Lookup/Full Compaction 工作。修复上,可以根据业务场景增加 Compaction 资源、使用独立 Compaction Job,或在业务低峰执行针对性的 Full Compaction;预防上,按峰值写入速度反推 Compaction 算力,不要让 Compaction 和业务写入长期挤同一份资源。
故障二:Changelog 延迟或语义异常。 现象:下游流作业看到的 -U/+U 对缺失、延迟过高或语义不符合预期。原理上这和 changelog-producer 选型直接相关:input 模式依赖上游提供完整 Changelog;lookup 模式依赖 Lookup Compaction,异步 Compaction 时 Changelog 可用时间可能进一步延后;full-compaction 模式则只能在相应 Full Compaction 完成后产生状态差异。定位时先确认用的哪种模式、上游记录是否自带 before image、下游需要的是 Upsert 还是完整 Changelog;修复上按消费者契约切换模式,而不是简单通过降低上游乱序来解决。
故障三:快照膨胀导致存储上涨。 现象:写入量稳定但存储持续增长。原理是 Snapshot 保留过久、Changelog 文件增加、Compaction 重写以及历史 Snapshot 仍引用数据文件等因素共同作用;修复上按业务需要调整 Snapshot Retention,检查 Tag/Branch 等长期引用,并通过 Paimon 自身的文件生命周期管理和孤儿文件清理机制回收不再需要的数据。重要边界:不要直接到对象存储/HDFS 里手删 Paimon 数据文件,那会破坏 Snapshot/Manifest 与物理文件之间的引用关系,可能直接导致表读取失败——这是 Paimon 运维中必须严格避免的低级事故。
9.【高级面试与架构思考】
题一:为什么 Paimon 适合实时更新?讲成机制。
要点:Paimon 主键表采用 LSM 存储模型,写入先进入缓冲并形成新的 Sorted Run,不需要每次更新都立即重写整个历史数据文件;后台 Compaction 再逐步合并多个版本。因此它把高频更新路径与历史整理路径解耦,更适合 CDC 和持续更新型工作负载。不要简单回答"Paimon 一定比 Iceberg 快",因为最终性能取决于具体更新比例、文件布局、查询模式和资源配置。
题二:changelog-producer 四种模式怎么选?
要求说清"消费者语义 vs 额外成本"的权衡:下游能处理 Upsert 用 none;上游已经提供完整 Changelog 用 input;需要从已有状态中补出 before image 用 lookup;能够接受周期性状态变化、并希望通过 Full Compaction 生成 Changelog 用 full-compaction。
题三:Compaction 为什么会积压?怎么压测出来?
要点:本质是写入速度与后台合并速度的赛跑;压测方法是等比复现峰值写入,监控 Sorted Run 数量、Compaction backlog 和查询延迟随时间是否收敛,不收敛就说明当前 Compaction Capacity 无法覆盖写入压力,需要继续检查 CPU、IO、网络、文件大小和 Compaction 策略。
题四:作业失败重启,下游为什么不会读到重复或跳变数据?
要点:Flink Checkpoint 与 Paimon Commit 协作,只有已经成功提交并发布的 Snapshot 才作为正常表状态对外可见;作业恢复后从最近成功 Checkpoint 恢复状态和消费位置,未成功提交的数据文件不会直接成为新的表状态。需要注意,"端到端精确一次"并非 Paimon 单独保证,而是由 Source、Flink Checkpoint、Sink Commit 以及下游消费语义共同决定。
10.【总结、知识提炼与下一篇】

本篇把第 09 篇提出的四个词落进了 Paimon 的代码:
存储统一体现为一张主键表既能通过 LSM 维护持续变化的逻辑状态,又能基于 Snapshot 进行批量读取和增量消费;
Compaction 是后台把多个 Sorted Run 归并、减少读时 Merge 工作、用写放大换读性能的核心机制;
迟到更新通过 Sequence、Merge Engine 和多版本合并体现;
两阶段提交则通过 Flink Checkpoint 与 Paimon Commit 协作,把已经完成的数据发布为新的 Snapshot。对照第二章:我们画清了写缓冲→数据文件→Sorted Run→提交的源码链路,讲清了四种 Changelog 模式的取舍,掌握了 Bucket、Compaction、快照和 Changelog 三类故障的源码视角排障,也把"Paimon 为什么适合实时更新"讲成了机制。
一句话结论:Paimon 的内核就是"用 LSM 的持续写入与后台合并,把湖表从单纯的批文件集合变成可持续更新、可流式读取、可版本管理的状态存储——它的优势和运维痛点,都源自这个模型选择。
下一篇(第 11 篇)将钻进与 Paimon 提交强绑定的 Flink:Checkpoint 与状态后端源码深度剖析——Barrier 如何对齐传播、Unaligned Checkpoint 为何能够缓解反压下的快照超时、RocksDB/ForSt 状态后端和增量快照如何工作、Savepoint 与故障恢复的全流程。今天讲的"Checkpoint 参与写入提交两阶段提交",下一篇会在 Flink 侧找到它们真正的源头。

