核心概念
1. 两种表类型:主键表 vs Append-only 表
1.1 主键表(Primary Key Table,最常用)
适合场景:业务库 CDC、维度表、明细表、汇总表、需要按主键 Update/Delete 的表。
建表示例:
CREATE TABLE ods_order_info (
id BIGINT,
user_id BIGINT,
total_amount DECIMAL(18,2),
order_status TINYINT,
create_time TIMESTAMP(3),
update_time TIMESTAMP(3),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'bucket' = '4',
'changelog-producer' = 'full-compaction'
);
主键表底层是 LSM 树结构:MemTable + SST 文件分层,写入后自动 Merge。主键唯一性在 Flink 写入端保证(跨 Checkpoint 多流拼接时用同一主键保证幂等)。
1.2 Append-only 表(日志/埋点专用)
适合场景:日志、埋点、行为数据、物联网数据,只 Insert,绝不 Update/Delete。
建表示例:
CREATE TABLE dwd_user_behavior_log (
user_id BIGINT,
page_id STRING,
action STRING,
ts TIMESTAMP(3),
WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
) PARTITIONED BY (dt STRING)
WITH (
'bucket' = '-1', -- 动态 bucket,自动分裂
'sequence.field' = 'ts' -- 乱序时按 ts 取最新,防重复
);
⚠️ 千万不要把日志类表建成主键表——主键表要做 LSM Merge + 索引,写放大高;日志场景就 Append-only + sequence.field 防重,性能是主键表的 2~5 倍。
2. 分区 Partition(时间维度必选)
跟 Hive 分区完全一致,按 dt / date / 小时分区就行。
CREATE TABLE dwd_order_detail (...)
PARTITIONED BY (dt STRING, hr STRING)
WITH (
'partition.expiration-time' = '180 d', -- 分区过期自动删(对象存储省钱)
'partition.expiration-check-interval' = '1 h'
);
强烈建议:除了小型维表(< 1000 万行),其他全部按天(或天+小时)分区。理由:
- 下游 BI 查询能按分区裁剪扫描量;
- 数据 TTL 过期删除方便;
- Compaction 是按分区 + Bucket 粒度做的,分区能把小文件的"脏区"控制在更小的范围。
3. Bucket(分桶,数据倾斜/Join 性能核心)
3.1 Bucket 是什么?
一张主键表 = N 个分区 × M 个 Bucket,每个 Bucket 是一个独立的 LSM 树。相同主键的行永远在同一个 Bucket 里,保证 Upsert 不会跨文件。
3.2 Bucket 数怎么选?(经验值)
| 单分区数据量 / 表规模 | 推荐 bucket |
|---|---|
| 小维表,< 10 GB | bucket = '1' |
| 中等事实表,10~200 GB | bucket = '4' ~ bucket = '16' |
| 大表,> 200 GB | bucket = '32' ~ bucket = '128' |
如果一开始估不准,直接上动态 bucket:
'bucket' = '-1',
'bucket-key' = 'id' -- 分桶 key,不填默认就是主键
动态 bucket 会根据写入数据量自动从 1 开始按 2 倍指数分裂(1→2→4→8...),避免一开始 bucket 写得很大后期倾斜。
⚠️ 一旦固定 bucket 数(比如
bucket = '4'),以后就不能改了! 要改只能建新表INSERT OVERWRITE迁数据。所以拿不准就用-1动态 bucket,运维成本最低。
4. Snapshot 快照(时间旅行核心)
每次 Flink Checkpoint 提交成功 / 每次 Batch Insert 成功都会产生一个新的 Snapshot。Paimon 会把快照链保留起来:
Snapshot-1 → Snapshot-2 → Snapshot-3 → ... → Snapshot-N
09:00 09:05 09:10 12:00
4.1 查询旧快照(时间旅行)
-- Flink SQL
SELECT * FROM dwd_order_detail /*+ OPTIONS('scan.mode'='snapshot','scan.snapshot-id'='123') */;
SELECT * FROM dwd_order_detail FOR SYSTEM_TIME AS OF TIMESTAMP '2026-09-04 10:00:00';
-- Spark SQL
SELECT * FROM paimon.db.dwd_order_detail VERSION AS OF 123;
SELECT * FROM paimon.db.dwd_order_detail TIMESTAMP AS OF '2026-09-04 10:00:00';
4.2 快照保留策略(默认保留 1 小时 ~ 1 天,可改大)
CREATE TABLE ... WITH (
'snapshot.time-retained' = '7 d', -- 快照最长保留 7 天
'snapshot.num-retained.min' = '10', -- 至少保留最近 10 个快照
'snapshot.num-retained.max' = '1000' -- 最多保留 1000 个
);
5. Tag(长期版本标记,重要节点必打)
Snapshot 会过期删掉,但 Tag 是用户显式打上、永远不会自动删的版本标记。适合:
- 每天 T+1 跑批完成,打
daily_20260904标签; - 月底/年底结算,打
monthly_202609永久保留; - 数据出现问题,快速在
daily_xxx上回滚。
-- 创建 Tag
CALL sys.create_tag('db.dwd_order_detail', 'daily_20260904', 123);
CALL sys.create_tag('db.dwd_order_detail', 'daily_20260904', TIMESTAMP '2026-09-04 23:59:59');
-- 读取 Tag 版本
SELECT * FROM dwd_order_detail /*+ OPTIONS('scan.tag-name'='daily_20260904') */;
-- 删除 Tag(手动清理)
CALL sys.delete_tag('db.dwd_order_detail', 'daily_20260904');
6. Changelog Producer(下游流读关键)
Paimon 作为湖仓表,不仅能让下游批读,还能像 Kafka 一样被 Flink 流读——这靠的就是 changelog-producer 参数。
| 可选值 | 适用场景 | 说明 |
|---|---|---|
none(默认) | 只有上游会 CDC 写入,下游不需要流读 | 仅产生 Insert / After 记录,不产生 Before,下游批读 OK,流读拿不到撤回 |
input | 上游就是 Flink CDC / Flink 自己的 Upsert 流(推荐,写性能最好) | 直接复用上游流里自带的 Before/After,额外计算为 0 |
full-compaction | 上游是纯批写入 / Spark/Hive/Upsert 但链路不带 Before(最通用) | 做 Full Compaction 时自己生成 Before/After,需要攒一攒才发一次 changelog,延迟略高但兼容所有写入方 |
lookup | 上游不是 CDC,你又希望延迟低,愿意用 Lookup Join 查旧值 | 每写一条查一次主键旧值,性能低,不推荐 |
最推荐的组合:
-- ODS 业务表:上游就是 Flink CDC,直接用 input 最简单
CREATE TABLE ods_order_info (...) PRIMARY KEY (id) NOT ENFORCED
WITH ('changelog-producer' = 'input');
-- DWD/DWS 汇总:上游是你自己跑的 Flink SQL聚合(带 Upsert),也用 input
CREATE TABLE dws_user_day_agg (...) PRIMARY KEY (user_id, dt) NOT ENFORCED
WITH ('changelog-producer' = 'input');
如果以后要用 Spark 往 Paimon DWD 里写批量修复 SQL,下游 Flink 流着读 DWD,那把 DWD 改成
full-compaction就能同时兼容。
7. Compaction(小文件治理,Paimon 自动做)
流式写入默认每次 Checkpoint 会产出若干小文件,Paimon 自带 3 类 Compaction:
| 类型 | 触发时机 | 用途 |
|---|---|---|
changelog-compaction | Checkpoint 后低优先级异步跑 | 合并 changelog 文件成 base 文件,缓解小文件 |
full-compaction | 按 full-compaction.delta-commits = 5(默认攒 5 个快照做一次) | 把 base 文件进一步合并,下游 changelog 产出时用 |
| 手动 Compaction | 调 CALL sys.compact('db.tbl', 'dt=2026-09-04') | 批跑 T+1 修复后手动强制做一次 |
参数调节(一般保持默认就够了,大表才调):
WITH (
'compaction.max.file-num' = '50', -- 单次最多合并 50 个文件
'compaction.min.file-num' = '5', -- 至少 5 个文件才合并,避免频繁小合并
'full-compaction.delta-commits' = '5', -- 每 5 个快照做一次全量合并
'target-file-size' = '128 MB' -- 合并后单文件目标大小
);
8. Catalog 与元数据集成
8.1 Paimon Catalog(Flink/Spark 直接用)
Flink SQL 建 Catalog:
CREATE CATALOG paimon WITH (
'type' = 'paimon',
'metastore' = 'hive', -- 核心参数:复用现有 Hive Metastore
'uri' = 'thrift://master1:9083', -- 你 HDP 3.3.2 的 HMS URI
'warehouse' = 'hdfs:///user/hive/warehouse/paimon' -- 和 Hive 共用 warehouse 或单独建目录
);
USE CATALOG paimon;
8.2 HMS 集成的优点
- Hive / Spark / Trino / Doris / StarRocks 全部能通过相同 HMS 看到 Paimon 表,不用在每个引擎单独建元数据;
- 权限直接复用 Ranger/Hive 权限,不用再给 Paimon 单独做一套。