Apache Paimon 流式数据湖深度实战:从主键表 LSM 结构、Changelog 生成与 Merge Engine 到 Bucket 分桶与 Lookup Join 的工程全解
如果你把 Iceberg 或 Delta Lake 直接接到 Flink 流式写入上,几乎一定会撞上同一堵墙:每次 commit 都要重写 manifest、列出上万文件、事务冲突靠乐观重试,分钟级提交尚可,秒级 Upsert 就是灾难。Iceberg 的设计前提是"批式写入、大文件、低频提交",它把一致性做在了元数据侧;而流式湖仓真正需要的是"高频写入、主键去重、可读 changelog、秒级可见"。这正是 Apache Paimon(原 Flink Table Store)存在的理由——它不是又一个 table format,而是把 LSM-Tree 的写入模型搬到了对象存储上。
一、核心设计:LSM 语义 + 湖存储格式
Paimon 的目录结构与 Iceberg 形似而神不同:
warehouse/db.db/orders/
├── schema/
│ └── schema-0
├── manifest/
│ ├── manifest-list-xxx
│ └── manifest-xxx
├── snapshot/
│ ├── LATEST
│ └── snapshot-3
├── bucket-0/
│ ├── data-8f3a2b1c.parquet
│ └── changelog-xxxx.parquet
└── bucket-1/
关键差异在三点:
- Snapshot 是 LSM 的"版本视图",每次 commit 生成一个 snapshot,指向一组 manifest,manifest 记录文件的元数据与统计信息(min/max、null count、row count),用于谓词下推与文件裁剪。
- Bucket 是一级物理分区,主键表按
bucket = hash(主键) % num_buckets强制分布。这一点极其重要:主键的 Upsert 只需要落到一个确定的 bucket 内做归并,跨 bucket 无需协调,从而把分布式 Upsert 的代价从"全局 join"降到"单桶局部归并"。 - Bucket 内部是 LSM 的 sorted run 结构,L0 文件由每次 checkpoint 刷出,后台 compaction 把多层文件归并。读取时对同一个 key 的多个版本按
sequence number排序,取最新一条。
所以 Paimon 主键表本质上是一个"跑在 S3 上的分布式 LSM-Tree",而 append-only 表(无主键)则退化为普通的湖格式,只做小文件合并与分区管理。
二、主键表与 Merge Engine:Upsert 的三种语义
建一张主键表:
CREATE TABLE paimon_orders (
order_id BIGINT,
user_id BIGINT,
status STRING,
amount DECIMAL(18,2),
update_time TIMESTAMP(3),
dt STRING,
PRIMARY KEY (order_id, dt) NOT ENFORCED
) PARTITIONED BY (dt)
WITH (
'bucket' = '16',
'merge-engine' = 'deduplicate',
'sequence.field' = 'update_time',
'changelog-producer' = 'lookup',
'compaction.max.file-num'= '50'
);
merge-engine 决定了相同主键多条记录如何合并,这是 Paimon 最有工程价值的抽象:
- deduplicate:默认。保留按
sequence.field排序最大的一条。没有sequence.field时按输入顺序,最后写入的胜出。生产环境强烈建议显式配置 sequence field,否则乱序数据(比如两条 binlog 落到不同 Flink 并行度)会产生不可复现的结果。 - partial-update:按列做部分更新,非 null 字段覆盖、null 字段保留旧值。这是构建"宽表"的杀手级能力——十几个业务流各自写自己的列,各自独立 checkpoint,最终在主键上拼成一张大宽表,彻底替代"多流 join + 状态爆炸"的老方案。
- aggregation:对指定列做聚合(sum/max/last_value 等),写入即增量计算,相当于把物化视图的下推到了存储层。
用 partial-update 做宽表拼装:
CREATE TABLE user_profile_wide (
user_id BIGINT,
nickname STRING,
city STRING,
last_login TIMESTAMP(3),
total_amt DECIMAL(18,2),
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'merge-engine' = 'partial-update',
'fields.total_amt.aggregate-function' = 'sum',
'fields.last_login.sequence-group' = 'login_time'
);
注意 sequence-group:它让 last_login 只与同组的 sequence 比较,避免一个来源的旧数据污染另一个来源的字段。这是 partial-update 在生产里最容易被忽略、也最容易出线上故障的配置项。
三、Changelog 生成:让下游能正确消费更新流
湖表存的是"最终状态",但下游 Flink 作业往往需要的是"+I/-U/+U/-D"变更流。Paimon 通过 changelog-producer 提供四种策略,取舍非常明确:
| 模式 | 原理 | 延迟 | 写放大 | 适用场景 |
|---|---|---|---|---|
| none | 不产出 changelog | - | 无 | 下游只做批查/全量读 |
| input | 直接透传上游 changelog | 最低 | 无 | 上游本身是 CDC 流 |
| lookup | 写前查旧值生成 -U/+U | 低 | 中(每次写有 lookup) | 主键表通用 Upsert |
| full-compaction | 每次 full compaction 对比前后 | 高 | 低 | 高吞吐、容忍延迟 |
input 模式下 Paimon 会把上游的 changelog 落盘成 changelog-*.parquet,下游读取时无需回查历史,代价是存储放大;lookup 模式则在写入时通过本地 RocksDB/内存缓存查旧值,要求 bucket 内数据可被单一 writer 独占,因此开启 lookup 后 write-only 与多作业并发写要格外小心。
一个典型的坑:changelog-producer=none 的下游做 SUM(amount) 聚合时,如果上游发生 Upsert,状态会被重复累加。正确姿势是下游用 sum + retract,或者在上游开启 lookup 保证 changelog 完整。
四、Compaction 与读写分离
流式写入每 10 秒一个 checkpoint,L0 文件会迅速膨胀,读放大随之飙升。Paimon 提供两条独立可调的旋钮:
- Minor Compaction:合并小文件、消除同一 key 的多版本,默认在写入线程内触发,由
compaction.min.file-num/compaction.max.file-num控制。 - Full Compaction:把某个 bucket 全部文件归并成一到多个,一般由
full-compaction.delta-commits周期性触发。
生产上强烈推荐的形态是读写分离(write-only):
-- 写入作业:关闭自动 compaction,只负责刷文件
ALTER TABLE paimon_orders SET (
'write-only' = 'true'
);
-- 独立 compaction 作业,dedicated 资源,可随时调并发
CALL sys.compact(
`table` => 'db.paimon_orders',
partitions => 'dt=20261001',
order_strategy => 'order'
);
这样做的好处是写入延迟不受 compaction 抖动影响(避免 checkpoint 超时引发的连锁反压),compaction 可以单独扩缩容,且失败可重试而不阻塞主链路。
五、Bucket 分桶:静态、动态与跨分区 Upsert
bucket 的取值决定了表的可扩展性:
- 固定值(如 16):简单可控,但业务量涨 10 倍后单 bucket 过大,重新分桶需要全量重写。
- bucket = -1(动态分桶):Paimon 根据数据量自动扩 bucket,新增 bucket 时旧数据不动,新数据写入新 bucket。读取时对同一 key 跨 bucket 归并,代价是读侧多一层合并。适合增量为主、主键不跨历史更新的场景。
- 跨分区 Upsert:主键必须包含全部分区字段(如上面的
PRIMARY KEY (order_id, dt)),否则同一个order_id落在两个分区会变成两条独立记录。如果业务确实需要"主键跨分区唯一"(比如订单改期导致 dt 变化),必须加'cross-partition-upsert.enabled' = 'true',但要知道它会显著增加读侧代价。
六、Lookup Join:把 Paimon 当维表
Paimon 表可以直接作为 Flink 维表使用,这是流式湖仓最实用的能力之一——维表不再是静态的 Hive 分区,而是持续更新的湖表。
CREATE TABLE dim_user (
user_id BIGINT,
city STRING,
level STRING,
PRIMARY KEY (user_id) NOT ENFORCED
) WITH ('bucket' = '8', 'changelog-producer' = 'lookup');
SELECT o.order_id, d.city, d.level, o.amount
FROM kafka_orders AS o
LEFT JOIN dim_user FOR SYSTEM_TIME AS OF o.proc_time AS d
ON o.user_id = d.user_id;
调优要点:
- 开启 changelog-producer=lookup,否则维表更新无法下推,Per-Job 缓存会读到陈旧数据。
- 维表 bucket 数不要太大,lookup 会按 bucket 建立本地缓存,bucket 过多会显著增加内存。
- 大维表用
'lookup.cache' = 'rocksdb'并配置rocksdb.memory.managed,避免 OOM;小维表用partial缓存命中率更高。 - 用
lookup.join.cache.ttl控制刷新周期,它决定"维表更新到下游可见"的时延上界。
七、生产运维清单
文件与快照过期。Paimon 不会自动删旧数据,必须配置:
ALTER TABLE paimon_orders SET (
'snapshot.time-retained' = '1 h',
'snapshot.num-retained.min' = '10',
'snapshot.num-retained.max' = '100',
'changelog.time-retained' = '1 d',
'partition.expiration-time' = '30 d',
'partition.expiration-check-interval' = '1 d'
);
snapshot.time-retained 太小会让正在运行的流式读作业读不到旧 snapshot 而失败,一般给 1~2 小时;num-retained.min 是兜底,防止写入慢时被误清理。
Tag 与 Branch。Paimon 的 tag 可以给某个 snapshot 打上持久标签(如每日零点),用于可复现的批查询与时间旅行;branch 则支持独立的写入分支,做"先写 branch 校验再 merge 到 main"的灰度发布。
-- 创建每日 tag,保留 7 天
CALL sys.create_tag(`table` => 'db.paimon_orders', tag => 'dt-20261001');
CALL sys.expire_tags(`table` => 'db.paimon_orders', expire_time => '7 d');
小文件治理。核心是 checkpoint 间隔与文件大小目标的平衡:
| 参数 | 建议值 | 说明 |
|---|---|---|
| checkpoint interval | 30s~2min | 太短则文件爆炸 |
| target-file-size | 256MB(默认 128MB) | 对象存储上大文件更划算 |
| write-buffer-size | 256MB | 写缓存,与 target 配合 |
| num-sorted-run.compaction-trigger | 5 | 触发 minor compaction 的 run 数 |
| num-sorted-run.stop-trigger | 10 | 达到则强制写入限速,防写放大失控 |
stop-trigger 是一道保险:当 sorted run 堆积到阈值,Paimon 会主动降低写入速度,避免读放大无限恶化。如果线上频繁触发,说明 compaction 资源不足,而不是该把阈值调大。
八、选型判断
Paimon 不是 Iceberg 的替代品,两者定位不同:Iceberg 适合以批为主、多引擎(Spark/Trino/Flink)共享、强调 schema evolution 与时间旅行的开放湖仓;Paimon 适合以流为主、需要主键 Upsert、秒级可见、下游要消费 changelog 的实时湖仓。如果你的场景是"Flink 实时写入 + 秒级查询 + 下游还要接着算",Paimon 几乎是目前唯一工程上站得住的选择;如果场景是"Spark 批量 ETL + Trino 即席分析 + 多团队协作",老老实实用 Iceberg。
一个务实的混合架构是:Paimon 承接实时层(ODS/DWD 的主键表与宽表),Iceberg 承接离线层(DWS/ADS 的批表),两者通过 Flink 定期同步,各自发挥所长。试图用一种格式通吃流批两头,最终往往两头都不讨好。
结语
Paimon 的设计哲学可以概括为一句话:把数据库的主键语义,用 LSM 的方式实现在对象存储上。理解了 bucket 是分片的 LSM、snapshot 是版本视图、merge-engine 是合并策略、changelog-producer 是变更流出口这四件事,剩下的所有参数调优都是在"写放大、读放大、可见延迟"这个三角里做取舍。生产落地时优先解决三件事——显式配置 sequence field、读写分离跑 compaction、配好 snapshot 过期——就能避开 80% 的线上事故。

发表评论 取消回复