云计算百科
云计算领域专业知识百科平台

Apache Paimon 实时湖仓实战(第 1 篇):别再说 Paimon 只是表格式,真正值钱的是让更新进入数据湖

订单 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 或 CatalogSnapshotCommit

FileStoreCommitImpl 会在提交前检查待删除文件和修改范围冲突,再通过 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,避免把处理时间误当成业务顺序。

CREATE TABLE orders (
order_id BIGINT,
status STRING,
amount DECIMAL(18, 2),
source_version BIGINT,
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'bucket' = '4',
'sequence.field' = 'source_version'
);

对同一主键依次提交三条数据,最后一条故意迟到:

INSERT INTO orders VALUES (1001, 'CREATED', 100.00, 1);
INSERT INTO orders VALUES (1001, 'PAID', 100.00, 3);
INSERT INTO orders VALUES (1001, 'CANCELLED', 100.00, 2);

批读的预期结果仍应是版本 3 的 PAID。若删除 sequence.field 后重复实验,默认合并顺序依赖输入顺序,迟到的版本 2 可能成为最终行。这一对照证明的是 Sequence 与 Merge Engine 如何决定表内当前状态,不证明 Paimon 比其他系统更快。

— 观察对象:表内最终业务状态;只读。
— 正常信号:1001 只返回一行,source_version = 3。
— 异常信号:出现多行,或最终版本不是 3。
SELECT *
FROM orders
WHERE order_id = 1001;

接着核对提交边界:

— 观察对象:Snapshot 提交序列;只读。
— 正常信号:snapshot_id 连续,最近写入产生 APPEND 类提交。
— 它不能证明:订单最终状态正确、下游获得完整 UPDATE_BEFORE。
SELECT snapshot_id,
commit_user,
commit_identifier,
commit_kind,
commit_time,
total_record_count,
changelog_record_count
FROM orders$snapshots
ORDER BY snapshot_id;

最后检查物理文件:

— 观察对象:当前 Snapshot 引用的数据文件;只读。
— 判断目标:同一主键可能仍存在于不同 Sorted Run,逻辑一行不等于物理一行。
— 注意:生产大表应限定 Snapshot、分区或抽样范围,避免高成本元数据扫描。
SELECT * FROM orders$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 已提交当成数据已可见。

import org.apache.flink.table.api.EnvironmentSettings;
import org.apache.flink.table.api.TableEnvironment;
import org.apache.flink.table.api.TableResult;

public final class PaimonSnapshotWrite {
public static void main(String[] args) throws Exception {
if (args.length != 1) throw new IllegalArgumentException("warehouse is required");
TableEnvironment t = 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')");
TableResult write = 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
赞(0)
未经允许不得转载:网硕互联帮助中心 » Apache Paimon 实时湖仓实战(第 1 篇):别再说 Paimon 只是表格式,真正值钱的是让更新进入数据湖
分享到: 更多 (0)

评论 抢沙发

评论前必须登录!