实时看板要求订单数据从一分钟内可见缩短到十秒,团队最直接的动作是把 Flink Checkpoint 从 60 秒调到 10 秒。延迟指标很快变漂亮了,但几个小时后,对象存储请求、Manifest 和小文件数量一路上涨,Spark 扫描同一分区反而越来越慢。
这不是一次简单的“Checkpoint 太频繁”。在 Paimon 流写里,每次 Checkpoint 都会推动缓存数据 Flush 并形成可提交文件;当它再与活跃 Bucket 和 Changelog 文件相乘,十秒新鲜度就可能变成持续增长的文件债。
在 Paimon 流写中,Checkpoint 不只是保存 Flink 状态,它还推动数据 Flush 与 Snapshot 提交;间隔缩短会把同样数据切成更多提交和文件。
下文的写入路径与参数边界固定在Apache Paimon 2.0.0,源码基线为release-2.0.0/604e6d5e...。
十秒新鲜度为什么会换来一地小文件
Checkpoint 频率 × 活跃 Bucket 数 × 文件类型 → L0 Data File / Changelog File / Manifest 增长 → Compaction 与元数据扫描负担 → 对象存储请求和查询 Split 增长只说十秒一次 Checkpoint 并不能估算文件数。128 个 Bucket、完整 Changelog 与 4 个 Bucket、changelog-producer=none面对的是不同数量级的文件生命周期。
Changelog Producer 官方文档 特别说明:短 Checkpoint(示例为 30 秒)叠加大量 Bucket,每个 Snapshot 会产生许多小 Changelog 文件;precommit-compact=true可以合并 Changelog 文件,但它会把 Compact Coordinator 和 Worker 加进写入拓扑,并非零成本开关。
第一步不是调参,而是确认文件到底由谁产生
生产检查先只读:
SELECTsnapshot_id,commit_kind,total_record_count,changelog_record_countFROMorders$snapshotsORDERBYsnapshot_idDESC;SELECT*FROMorders$files;按时间计算 Snapshot 增速,按 Bucket 统计文件数、文件大小与 Level。若 Changelog 文件占主导,检查changelog-producer和precommit-compact;若 L0 Data File 持续增长,检查 Checkpoint、Bucket 和 Compaction;若只有少数 Bucket 增长,先查 Key 倾斜。
不要只看仓库目录总文件数。Snapshot 过期前旧文件不会立即物理删除,Compaction 生成新文件也会暂时让新旧文件共存。
止血不能只调大目标文件,否则只是推迟暴露
先建立故障时间线:Checkpoint 从何时缩短、文件增速何时改变、查询何时退化。没有时间对齐,只能证明表里文件多,不能证明是本次配置造成。
处理顺序应是:
- 确认业务真正需要的数据首次可见时间,而不是把 Checkpoint 间隔当 SLA。
- 在可接受新鲜度内适当拉长 Checkpoint,观察每分钟文件增量。
- 缩减过多 Bucket,或修复只写少数 Bucket 的倾斜。
- 只有完整 Changelog 确有需要时才保留对应 Producer。
- 对 Changelog 小文件评估
precommit-compact,同时验证 Checkpoint 尾延迟。 - 持续观察多个 Compaction 周期,要求文件产生速度长期低于消化速度。
调大target-file-size不能让每个短 Checkpoint 凭空攒够数据;增加 Compaction 并行度则可能放大对象存储读写并与 Writer 抢资源。
最小止血可以把 Checkpoint 恢复到变更前区间,但要明确数据新鲜度会下降;根因修复应在影子表按 10、30、60 秒三档复放相同输入,并固定 Bucket、Producer 与资源。验收至少覆盖两个完整 Compaction 周期:每分钟净增文件不再上升、Checkpoint P99 稳定、查询扫描文件数回落、业务哨兵行仍满足新鲜度 SLO。任一项恶化就停止继续缩短间隔。
本案例没有真实吞吐、Bucket 和对象存储账单,不能给出通用秒数。它证明的是乘法关系与观测方法,不证明 30 秒适用于所有表。
Paimon 的低延迟不是把 Checkpoint 无限缩短,而是在可见性、文件尺寸和合并能力之间找到能长期偿债的节奏。
用 Java 连续采样,算出文件债增长斜率
固定版本源码如何把 Checkpoint 变成文件提交
在release-2.0.0/604e6d5e...中,paimon-core/src/main/java/org/apache/paimon/table/sink/TableWriteImpl.java的prepareCommit(boolean, long)把本轮写入交给底层 FileStore Write;paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreWrite.java的prepareCommit(...)收集各分区和 Bucket 的新文件与 Compaction 结果。对应源码:TableWriteImpl、AbstractFileStoreWrite。
Java 连续采样$snapshots与$files,才能观察“提交次数上升—活跃 Bucket 被反复刷出—小文件斜率上升”。一次文件数快照只能证明当前存量,不能把因果直接归到 Checkpoint 周期。
以下只读程序按paimon-flink-1.20:2.0.0API 编写,需要定时保存多次输出才能计算斜率;本环境未运行全文示例。
importorg.apache.flink.table.api.EnvironmentSettings;importorg.apache.flink.table.api.TableEnvironment;publicfinalclassSmallFileAudit{publicstaticvoidmain(String[]args){if(args.length!=1)thrownewIllegalArgumentException("warehouse is required");Stringwarehouse=args[0].replace("'","''");TableEnvironmentt=TableEnvironment.create(EnvironmentSettings.newInstance().inBatchMode().build());t.executeSql("CREATE CATALOG p WITH ('type'='paimon','warehouse'='"+warehouse+"')");t.executeSql("USE CATALOG p");t.executeSql("SELECT snapshot_id,commit_time,total_record_count,changelog_record_count "+"FROM demo.orders$snapshots ORDER BY snapshot_id DESC LIMIT 120").print();t.executeSql("SELECT bucket,COUNT(*) files,AVG(file_size_in_bytes) avg_bytes "+"FROM demo.orders$files GROUP BY bucket ORDER BY files DESC").print();}}连续运行并保存输出,才能计算文件产生斜率;一次结果不能证明 Checkpoint 是根因。Checkpoint 驱动 Writer Flush 与 Commit,Changelog Producer 还可能增加独立文件,最终均进入 Manifest/Snapshot 提交链。
把 Checkpoint 调快,只能让数据更早提交;如果每次提交都切出一批小文件,低延迟最终会用 Compaction 和查询延迟来还账。
如果准备通过write-only把合并移出 Writer,下一步应先看独立 Compaction 没有真正接管时,查询为什么会一周比一周慢,避免把小文件问题从写入侧搬到读取侧。
官方资料
- Write Performance
- Changelog Producer
- Understand Files
- Snapshot Specification