订单 CDC 已经落到对象存储,Spark 也能查。业务第二天要看订单当前状态,团队才发现同一个order_id有十几个版本:CREATED、PAID、CANCELLED全在,哪一条才是现在,得由每条 SQL 临时判断。
把数据写成 Parquet 不难,给文件补上 Schema、分区和 Snapshot 也不难。真正难的是:不把对象存储改造成数据库,仍能持续接收 Update/Delete,让批查询读到可信当前状态,同时让流任务从提交边界继续消费。
Paimon 的关键差异不是又管理了一批文件,而是用 Bucket 内 LSM 承接主键更新,再用原子 Snapshot 把当前状态、批量重算和流式增量接到同一张湖表上。
同一批订单落湖,四条路径承担的责任不同
先固定共同前提:MySQL 订单持续产生 Insert、Update、Delete;同一主键可能乱序;分钟级可见即可;Flink 继续消费变化,Spark 需要随时重算;底层使用共享文件系统或对象存储。机器规模、实际吞吐和 SLA 未提供,因此这里只比较执行路径,不比较跑分。
| 路径 | 更新怎样保存 | 当前状态由谁计算 | 批流如何对齐 | 主要代价 |
|---|---|---|---|---|
| 普通 Parquet 追加目录 | 每次变化继续追加 | 每个查询按主键排序去重 | 另建文件发现、位点和重放规则 | 逻辑分散,删除、乱序和恢复容易各算一套 |
| Paimon Append Table | 以追加记录为主 | 上游或查询负责状态化 | Snapshot 提供提交边界,可流式读追加数据 | 不定义主键时,不能直接靠 Upsert 接收完整 Changelog |
| Paimon Primary Key Table | 同一 Bucket 内形成多组有序文件 | LSM 读取与 Merge Engine 合并同主键记录 | 批读与流读围绕同一组 Snapshot 推进 | Compaction、Bucket、Changelog 和保留策略必须治理 |
| 服务型数据库 | 数据库内部完成主键更新 | 数据库返回当前行 | 增量通常通过日志、订阅或导出链路提供 | 在线服务强,但存储成本、多引擎重算和历史开放性是另一套取舍 |
这张表没有绝对冠军。只保存不可变日志,Paimon Append Table 更简单;需要毫秒级点查和高并发接口,服务型数据库仍应保留;需要同一份低成本数据同时承接持续更新、流式消费和批量重算,Primary Key Table 才体现出差异。
LSM 没有消灭更新成本,只把随机改写变成顺序追加与后台合并
Paimon 2.0 Primary Key Table 文档 明确说明:一张表或一个分区会被拆成多个 Bucket,每个 Bucket 内部是一棵 LSM Tree。Bucket 是最小读写单元,也限制最大处理并行度。
新记录不是去对象存储里找到旧行并原地修改。写入先进入内存缓冲,排序后形成新的 Sorted Run。不同 Sorted Run 的主键范围可以重叠,也可以出现同一个主键;读取时必须按照 Merge Engine 和记录顺序合并。
CDC RowKind / 业务版本 → Partition 与 Bucket 路由 → 内存缓冲并按主键排序 → Checkpoint 刷出 L0 Sorted Run → Commit 汇总 Manifest 变化 → Snapshot 文件提交成功后全局可见 → 读取时合并,或由 Compaction / Deletion Vector 提前消化旧版本这条链把对象存储不擅长的随机更新,转换成追加新文件、提交元数据和后续合并。收益是写入可以保持顺序 I/O,代价则分散到三处:Writer 的 Flush 与提交、后台 Compaction、Reader 的多路归并。
Table Mode 文档 因此才区分 MOR、COW 和 MOW。MOR 把更多成本留给读取,COW 把全量合并放进写入,MOW 用 Deletion Vector 改善读取,但 L0 文件仍有 Compaction 后才可见的边界。所谓实时更新不是免费更新,而是可以明确选择成本由谁、在什么时候支付。
真正的可见性开关不是文件写完,而是 Snapshot 提交成功
Paimon 的 Data File、Manifest 和 Snapshot 不是同一个层次。Data File 保存数据,Manifest 描述文件增删,Snapshot 引用一组能够共同解释表状态的元数据。
Snapshot 格式说明 给出一个关键事实:每次 Commit 生成一个连续编号的 Snapshot;写入方抢占下一个 Snapshot ID,Snapshot 文件成功写入后,本次提交才可见。
固定到release-2.0.0,最短源码链是:
StoreSinkWrite / StoreSinkWriteImpl → Committable / ManifestCommittable → StoreCommitter → FileStoreCommitImpl → SnapshotCommit → RenamingSnapshotCommit 或 CatalogSnapshotCommitFileStoreCommitImpl会在提交前检查待删除文件和修改范围冲突,再通过SnapshotCommit完成原子提交。文件系统提交路径中的RenamingSnapshotCommit最终尝试原子写入snapshot-<id>;成功后更新LATEST提示。
LATEST只是提示文件,可能不准确。读取端在提示不可信时仍会扫描 Snapshot 文件确定边界。因此不能把目录里出现新 Data File、LATEST被更新,或 Flink 某个 Subtask 已完成,当作整张表已提交的充分证据。
一组乱序订单,能同时验证三层正确性
下面是一套最小实验设计,不是生产性能报告。
Paimon:2.0.0 计算引擎:Flink 1.20 存储:本地文件系统,仅用于隔离语义 表模式:Primary Key Table,默认 deduplicate / MOR 变化变量:同一 order_id 的输入顺序 未验证:对象存储延迟、并发 Writer、故障恢复和生产吞吐先创建一张订单当前状态表。实验使用单调业务版本source_version,避免把处理时间误当成业务顺序。
CREATETABLEorders(order_idBIGINT,statusSTRING,amountDECIMAL(18,2),source_versionBIGINT,PRIMARYKEY(order_id)NOTENFORCED)WITH('bucket'='4','sequence.field'='source_version');对同一主键依次提交三条数据,最后一条故意迟到:
INSERTINTOordersVALUES(1001,'CREATED',100.00,1);INSERTINTOordersVALUES(1001,'PAID',100.00,3);INSERTINTOordersVALUES(1001,'CANCELLED',100.00,2);批读的预期结果仍应是版本 3 的PAID。若删除sequence.field后重复实验,默认合并顺序依赖输入顺序,迟到的版本 2 可能成为最终行。这一对照证明的是 Sequence 与 Merge Engine 如何决定表内当前状态,不证明 Paimon 比其他系统更快。
-- 观察对象:表内最终业务状态;只读。-- 正常信号:1001 只返回一行,source_version = 3。-- 异常信号:出现多行,或最终版本不是 3。SELECT*FROMordersWHEREorder_id=1001;接着核对提交边界:
-- 观察对象:Snapshot 提交序列;只读。-- 正常信号:snapshot_id 连续,最近写入产生 APPEND 类提交。-- 它不能证明:订单最终状态正确、下游获得完整 UPDATE_BEFORE。SELECTsnapshot_id,commit_user,commit_identifier,commit_kind,commit_time,total_record_count,changelog_record_countFROMorders$snapshotsORDERBYsnapshot_id;最后检查物理文件:
-- 观察对象:当前 Snapshot 引用的数据文件;只读。-- 判断目标:同一主键可能仍存在于不同 Sorted Run,逻辑一行不等于物理一行。-- 注意:生产大表应限定 Snapshot、分区或抽样范围,避免高成本元数据扫描。SELECT*FROMorders$files;这套实验能够证明三件事:业务版本控制乱序更新、Snapshot 是提交可见性边界、逻辑当前行与物理旧版本可以同时存在。它不能证明流式下游一定得到完整撤回,也不能证明某种 Bucket 或 Compaction 配置适合生产。
表里是 100,Snapshot 也成功,下游仍可能算错
批读得到一行正确的PAID,只证明 Primary Key Table 的当前状态正确。下游若按门店汇总金额,更新 100 → 80 时需要知道旧值 100,才能先撤回再加入 80。
Paimon 默认changelog-producer=none不额外保存完整旧值变化;Flink 可能需要 Normalize State 记住每个主键的旧值。input、lookup和full-compaction可以在不同前提下提供更完整的 Changelog,却会增加文件、状态或 Compaction 成本。
因此需要分开验收:
源事件完整 ≠ Checkpoint 对应 Snapshot 已提交 ≠ 表内当前状态正确 ≠ 下游收到完整 Changelog ≠ 故障恢复后业务结果连续完整 Changelog 的选择和代价会在系列第 03 篇单独展开。本篇只保留边界:表内状态正确不能替下游计算正确作证。
Paimon 值得进入架构的信号只有三个
第一,同一份数据既有 Update/Delete,又要留在低成本共享存储中。第二,Flink 需要持续消费变化,Spark 或其他引擎又要基于同一提交点批量重算。第三,团队愿意治理 Bucket、Compaction、Snapshot 保留和 Changelog,而不是把它们当成默认参数。
不满足这些条件时,选择应更简单:
- 只有不可变事件,优先评估 Append Table;
- 只有在线主键查询和事务写入,保留 OLTP 或服务型数据库;
- 指标完全固定且要求极低读延迟,预计算和缓存可能更直接;
- 无法承担 Compaction、文件数和 Snapshot 生命周期治理,不要只因实时湖仓四个字引入 Paimon。
Paimon 替掉的不是 Flink、Spark 或数据库,而是同一份更新事实为了流计算、批量重算和历史存储被迫维护多套不一致副本的那部分复杂度。
文件能被很多引擎读取,只说明它足够开放;更新能在同一个 Snapshot 序列里被提交、合并、重算和继续消费,才是 Paimon 真正值钱的地方。
面试表达主线
面试时不要停在“Paimon 支持流批一体”。先说 Primary Key Table 如何在 Bucket 内以 LSM 接收更新,再说 Snapshot 如何形成原子可见边界,最后主动补上 Compaction、Changelog 与在线点查不是免费能力。这样回答的是机制、收益与代价,而不是产品口号。
Java 把一次写入真正提交成 Snapshot
依赖paimon-flink-1.20:2.0.0,参数传本地或对象存储 Warehouse。程序等待TableResult,失败会直接抛出,而不是把 SQL 已提交当成数据已可见。
importorg.apache.flink.table.api.EnvironmentSettings;importorg.apache.flink.table.api.TableEnvironment;importorg.apache.flink.table.api.TableResult;publicfinalclassPaimonSnapshotWrite{publicstaticvoidmain(String[]args)throwsException{if(args.length!=1)thrownewIllegalArgumentException("warehouse is required");TableEnvironmentt=TableEnvironment.create(EnvironmentSettings.newInstance().inBatchMode().build());t.executeSql("CREATE CATALOG p WITH ('type'='paimon','warehouse'='"+args[0]+"')");t.executeSql("USE CATALOG p");t.executeSql("CREATE DATABASE IF NOT EXISTS demo");t.executeSql("CREATE TABLE IF NOT EXISTS demo.orders (id BIGINT, status STRING, ver BIGINT, "+"PRIMARY KEY (id) NOT ENFORCED) WITH ('bucket'='4','sequence.field'='ver')");TableResultwrite=t.executeSql("INSERT INTO demo.orders VALUES (1001,'PAID',3),(1001,'CANCELLED',2)");write.await();t.executeSql("SELECT * FROM demo.orders$snapshots ORDER BY snapshot_id").print();}}executeSql(INSERT)进入 Flink Sink,最终沿StoreSinkWrite → StoreCommitter → FileStoreCommitImpl → SnapshotCommit提交;await()异常是写入失败信号,$snapshots有新行才证明形成可见提交。该示例验证提交与乱序语义,不代表生产吞吐。
Paimon 不是普通 Parquet 目录加元数据。Primary Key Table 在 Bucket 内用 LSM 把随机更新转成有序追加,由 Merge Engine 得到当前状态,再通过原子 Snapshot 统一批读与流读边界。代价是 Compaction、Bucket、Changelog 和 Snapshot 生命周期必须治理,它也不替代 OLTP 与在线查询数据库。
官方资料
- Apache Paimon 2.0 Documentation
- Primary Key Table
- Append Table
- Table Mode
- Sequence Field and RowKind
- Snapshot Specification
- System Tables
- Apache Paimon
release-2.0.0