Flink 与 Paimon 实时数仓
概述
Apache Paimon 是专为流式湖仓设计的表格式:Flink 写入实时更新,同时保留湖的开放与低成本。核心能力包括 Changelog Producer(变更日志生成)、Merge Engine(更新模式)与流读批读。本文逐一讲透。
一、Paimon 是什么
| 特性 | 说明 |
|---|---|
| 定位 | 流式湖仓存储 |
| 表格式 | 基于 LSM + 列式文件 |
| 引擎 | Flink 原生支持 |
| 核心 | 流写 + 流读 + 更新 |
Paimon 表 = LSM 文件(数据 + 变更日志)
+ Snapshot(快照)
+ 索引1.1 与 Iceberg/Delta/Hudi 差异
| 维度 | Paimon | 三大湖格式 |
|---|---|---|
| 流读 | 原生(changelog) | 弱 |
| 更新 | Merge Engine | 表类型 |
| 定位 | Flink 流优先 | 通用 |
Paimon 为 Flink 实时场景优化:
upsert、流读、增量二、表结构
2.1 建表
sql
CREATE TABLE orders (
order_id BIGINT PRIMARY KEY NOT ENFORCED,
user_id BIGINT,
amount DECIMAL(10, 2),
ts TIMESTAMP(3)
) WITH (
'connector' = 'paimon',
'path' = 'hdfs:///warehouse/orders',
'bucket' = '4'
);| 参数 | 说明 |
|---|---|
| path | 表存储路径 |
| bucket | 桶数(并行粒度) |
| primary-key | 主键(更新依据) |
2.2 存储模型
LSM 结构:
Level 0:新写入的小文件
Level N:合并后的文件
主键 → 最新记录| 特性 | 说明 |
|---|---|
| LSM | 写快、合并 |
| 主键索引 | 更新定位 |
| 快照 | 版本管理 |
三、Changelog Producer
3.1 作用
生成表的变更日志(changelog):
+I(插入)、+U/-U(更新)、-D(删除)
→ 下游可以流式消费变更3.2 配置
sql
WITH (
'changelog-producer' = 'input' -- 或 lookup / full-compaction
)| 模式 | 说明 |
|---|---|
| input | 记录输入变更(默认,更新不透出) |
| lookup | 查询模式生成完整变更 |
| full-compaction | 全量合并时生成 |
3.3 流式读取变更
sql
-- 流读:持续消费变更
SELECT * FROM orders /*+ OPTIONS('scan.mode'='latest') */;
-- 增量读
SELECT * FROM orders
WHERE _snapshot_id > 100;| 能力 | 说明 |
|---|---|
| 流读 | 持续输出变更 |
| 增量 | 从指定快照 |
| 回溯 | 时间/快照 |
四、Merge Engine
4.1 作用
决定同主键多记录的合并方式:
相同主键 → 按引擎合并| 引擎 | 合并逻辑 |
|---|---|
| deduplicate(默认) | 保留最新 |
| partial-update | 字段级部分更新 |
| aggregation | 聚合(sum/max/min) |
4.2 配置
sql
WITH (
'merge-engine' = 'deduplicate'
)| 场景 | 引擎 |
|---|---|
| 实时更新主表 | deduplicate |
| 多字段分步写入 | partial-update |
| 实时聚合累计 | aggregation |
4.3 示例:聚合引擎
sql
CREATE TABLE order_agg (
user_id BIGINT PRIMARY KEY NOT ENFORCED,
cnt BIGINT,
gmv DECIMAL(10,2)
) WITH (
'merge-engine' = 'aggregation',
'fields.cnt.aggregate-function' = 'sum',
'fields.gmv.aggregate-function' = 'sum'
);
-- 实时累计聚合
INSERT INTO order_agg
SELECT user_id, 1, amount FROM orders;五、写入模式
5.1 批量写
sql
INSERT INTO orders SELECT ... FROM source;5.2 流式写(upsert)
Flink 流写:
CDC → Paimon(实时更新)
主键变更 → 更新/删除sql
-- Flink CDC 实时入 Paimon
INSERT INTO orders
SELECT order_id, user_id, amount, ts FROM cdc_orders;| 能力 | 说明 |
|---|---|
| 精确一次 | Flink checkpoint |
| 更新删除 | 主键语义 |
| 背压 | 流式反馈 |
5.3 写参数
| 参数 | 说明 |
|---|---|
write.buffer-size | 写缓冲 |
write.buffer-spillable | 溢出 |
| 桶数 | 与并行度匹配 |
六、实时数仓中的位置
6.1 典型链路
CDC → Kafka → Flink → Paimon(ODS/DWD)
→ Flink 流读 → 实时计算 → Doris/大屏
→ Spark 批读 → 离线加工| 层 | Paimon 角色 |
|---|---|
| ODS | CDC 明细落湖 |
| DWD | 实时宽表 |
| DWS | 实时聚合(aggregation) |
6.2 与 Doris 配合
Paimon:明细 + 更新(湖)
Doris:高并发查询(OLAP)
Paimon → Doris:批量/实时同步七、运维要点
| 要点 | 说明 |
|---|---|
| 小文件 | 自动/手动合并 |
| 快照清理 | 定期过期 |
| 桶数规划 | 一劳永逸,后续改需重建 |
| 状态 | Flink 状态与 Paimon 配合 |
7.1 合并
sql
CALL paimon.sys.compact('orders');7.2 快照管理
| 参数 | 说明 |
|---|---|
| snapshot.num-retained | 保留数 |
| snapshot.time-retained | 保留时间 |
八、Paimon vs 场景选型
| 场景 | 推荐 |
|---|---|
| Flink 实时数仓 | Paimon |
| 通用湖仓多引擎 | Iceberg |
| 更新/CDC 密集 | Paimon / Hudi |
| Spark 生态 | Delta |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 更新不生效 | 检查主键与 merge-engine |
| 流读无变更 | changelog-producer 配置 |
| 写慢 | 调桶数/缓冲 |
| 快照膨胀 | 清理策略 |
| 与 Doris 同步慢 | 批量 + 幂等 |