Apache Paimon的Compaction是一种在后台自动合并数据文件的过程,其主要目的是通过减少文件数量、优化数据布局来提升查询性能,并在数据持续写入时维持稳定的读写效率。
'write-only' = 'true' 来让写入作业跳过Compaction。然后,通过一个独立的Flink作业(专用Compaction作业)来专门负责Compaction。这种方式使得数据写入可以持续进行,尤其适用于多作业同时写入的场景。Flink 分离式 Compaction 策略 本文档适用于在火山引擎(Volcano Engine) 上使用 EMR、LAS 等服务的 Paimon 用户,涵盖以下场景:
目标读者:大数据平台工程师、数据仓库开发工程师、DataOps 工程师,建议具备基本的 Flink/Spark 使用经验和 SQL 基础。
本节主要介绍 Paimon Compaction 的机制和核心技术要点。通过介绍 Paimon Compaction 的技术要点为用户提供更全面的 Compaction 能力,帮助用户更好地理解 Compaction 的技术细节,并提供更灵活的 Compaction 配置方式,以满足不同场景的需求。
write-only = true 将写入端与 Compaction 解耦,消除多 Job 竞争同一分区时的合并冲突,适合高并发实时写入场景。flink.hadoop.fs.tos.rename.enabled: true。Apache Paimon 是一种面向流批一体场景的数据湖存储格式,其主键表采用 LSM Tree(Log-Structured Merge Tree) 结构存储数据。LSM Tree 将写入操作转化为顺序追加,极大提升了写入吞吐,但随着数据的持续写入,存储层会积累大量层级各异、大小不一的小文件(Sorted Run),给读取查询带来严重的 I/O 放大问题。
Compaction(合并压缩) 是 LSM Tree 架构下的必要维护手段,其核心目标包括:
对于 Append 表(仅追加表),Compaction 的主要目的是将写入过程中产生的大量小文件合并为较少的大文件,提升下游读取的扫描效率,同时为 Sort Compact 等高级特性提供基础。
主键表以 LSM Tree 为底层存储结构。每次写入操作首先进入内存的 Write Buffer,当 Write Buffer 满时执行 flush,将数据写入磁盘形成一个 Sorted Run(有序文件集合)。随着写入的持续进行,磁盘上会积累越来越多的 Sorted Run,形成多个 Level。
每个 Sorted Run 内部按主键排序,但不同 Sorted Run 之间的主键范围可能存在重叠。查询时需要对多个 Sorted Run 执行 Merge-on-Read,Sorted Run 数量越多,查询性能越差。
Compaction 的核心任务就是将多个 Sorted Run 合并为更少的 Sorted Run,同时完成主键去重、Changelog 生产和快照过期清理等附加工作。
Paimon 主键表默认采用 Universal Compaction 策略(类似 RocksDB 的 Universal Compaction)。该策略基于以下规则触发合并:
num-sorted-run.compaction-trigger(默认 5)时,触发一次合并。当使用 changelog-producer = lookup 时,Paimon 会在 Compaction 前执行 Lookup 操作以生成 Changelog,通过主键查询数据更新前后的值,在 Compaction 的过程中产生 -U 和 +U 的真实 changelog。这种 changelog-producer 会增加 Compaction 的成本。
Write Stall(写入停顿) 是 Paimon 主键表的一种自我保护机制。当 Sorted Run 数量过多时,若不加限制地继续写入,会导致读取性能急剧恶化。Write Stall 会在 Sorted Run 超过 num-sorted-run.stop-trigger 时暂停写入,等待 Compaction 赶上进度。num-sorted-run.stop-trigger 的默认值为 num-sorted-run.compaction-trigger + 3,即默认 8。
在异步 Compaction 场景下,通常需要将此参数设为一个极大值(如 2147483647),以避免 Write Stall 影响写入吞吐。
参数名 | 默认值 | 说明 |
|---|---|---|
| 5 | Sorted Run 数量超过此值时触发 Compaction |
| compaction-trigger + 3 (即 8) | Sorted Run 数量超过此值时触发 Write Stall |
|
| Changelog 生产方式:none / input / full-compaction / lookup |
| 256MB | Write Buffer 大小 |
| 128MB 主键表 | 目标文件大小 |
| false | 开启后跳过 Compaction(用于分离式方案) |
Append 表不支持主键更新,数据仅追加不修改。根据是否指定 Bucket 数量,Append 表分为两种子类型,其 Compaction 机制存在本质差异。
当 bucket = -1 时,表为 Unaware-Bucket 模式。Flink 写入任务中会自动启动一个 Compact Coordinator 和若干 Compact Worker 算子,形成独立的 Compaction 拓扑:
这种拓扑设计的核心优势是:Compaction 完全异步执行,永不对写入造成反压。Compaction Worker 的失败不会影响写入链路,写入与合并天然解耦。
当 bucket = N(N > 0)时,每个 Bucket 内的数据按写入顺序有序排列。Compaction 在 Writer 内部的 Bucket 级别执行,将同一 Bucket 内的多个小文件合并为更大的文件,保证 Bucket 内的数据顺序性。
由于 Compaction 在 Writer 内部执行,在极端情况下可能影响写入吞吐(但通常影响较小)。
对比维度 | Unaware-Bucket(bucket=-1) | Bucketed Append(bucket=N) |
|---|---|---|
Compaction 执行位置 | 独立 Coordinator + Worker 算子 | Writer 内部 |
反压风险 | 无(永不反压) | 极小 |
数据有序性 | 不保证全局有序 | 保证 Bucket 内有序 |
Sort Compact 支持 | 支持(Batch 模式) | 支持(Batch 模式) |
并发写入安全性 | 高 | 需注意 Bucket 并发写 |
适用场景 | 大规模非有序写入 | 有序消费场景 |
Append 表 Compaction 的触发规则如下:
compaction.min.file-num,且这些文件的总大小 ≥ target-file-size 时触发。另外根据放大率等参数,判断是否触发合并。参数名 | 默认值 | 说明 |
|---|---|---|
| 5 | 触发 Compaction 的最小文件数 |
| 128MB 主键表 | 目标合并文件大小 |
对比维度 | 主键表(Primary Key Table) | Append 表(Append-Only Table) |
|---|---|---|
存储结构 | LSM Tree(多级 Sorted Run) | 扁平文件集合 |
Compaction 触发机制 | Sorted Run 数量超阈值 | 文件数量 / 文件大小超阈值 |
反压影响 | 可能触发 Write Stall | Unaware-Bucket 永不反压 |
合并语义 | 主键去重 + 版本合并 | 文件合并(追加语义保留) |
Changelog 生产 | 与 Compaction 深度耦合 | 无 Changelog 需求 |
分离式 Compaction(write-only=true) | 支持 | 支持 |
适用数据模式 | CDC、更新类数据 | 日志、事件流、追加类数据 |
Flink 写入 Paimon 表时,默认在写入算子(Sink Writer)内部自动执行 Compaction。这是最简单的使用方式,无需额外配置,适合写入吞吐适中、对延迟不敏感的场景。
默认行为说明:
num-sorted-run.compaction-trigger(默认 5),则触发一次 Compaction。num-sorted-run.stop-trigger(默认 8),则触发 Write Stall,暂停写入,直到 Compaction 赶上。适用场景: 单一 Flink 写入 Job、写入并发不高、数据量适中的在线实时入库场景。
当写入吞吐较高,同步 Compaction 模式下频繁 Write Stall 影响稳定性时,可开启异步 Compaction 模式。异步 Compaction 通过禁用 Write Stall 并放宽 Compaction 触发阈值,让写入尽可能不受 Compaction 影响。
参数 1:num-sorted-run.stop-trigger = 2147483647
将 Write Stall 触发阈值设为 Integer.MAX_VALUE,实际上等同于禁用 Write Stall。写入端将不再因 Sorted Run 数量过多而暂停,确保写入链路的吞吐稳定性。代价是短期内 Sorted Run 数量可能较多,读取时 Merge 开销增大。
参数 2:sort-spill-threshold = 5
在 Compaction 合并多个 Sorted Run 时,若参与合并的 Sorted Run 数量超过此阈值,则将中间排序结果溢写到磁盘而非保留在内存,防止大规模合并时发生 OOM。建议根据单 TaskManager 可用内存和 Sorted Run 的平均大小调整,初始值建议为 5。
参数 3:changelog-producer.lookup-wait = false (Paimon 1.1+ 直接设置 lookup-wait = false)
当 changelog-producer = lookup 时,默认情况下写入端需要等待 Lookup Compaction 完成后才能提交 Checkpoint(lookup-wait = true)。将此参数设为 false 后,写入端无需等待 Lookup Compaction 完成即可提交 Checkpoint,大幅降低 Checkpoint 延迟,提升写入吞吐。
注意:
lookup-wait = false时,Changelog 的生成会有一定延迟,下游消费者可能需要等待更长时间才能看到最新变更。
-- 创建开启异步 Compaction 的主键表 CREATE TABLE orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(18, 2), order_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'bucket' = '8', 'changelog-producer' = 'lookup', -- 异步 Compaction 核心参数 'num-sorted-run.stop-trigger' = '2147483647', 'sort-spill-threshold' = '5', 'changelog-producer.lookup-wait' = 'false', -- 写入性能优化 'write-buffer-spillable' = 'true', 'write-buffer-size' = '256MB', -- 快照相关保留周期 'snapshot.expire.execution-mode' = 'async', 'snapshot.time-retained' = '1d' );
适用场景: 高吞吐实时写入(单表每秒数百万条记录)、对写入延迟敏感但对读取稍有容忍的场景。
分离式 Compaction 是将 Compaction 工作从写入 Job 中完全剥离,由专用的独立 Compaction Job 承担的架构模式。
写入端通过设置 write-only = true 跳过所有 Compaction 操作(同时也会跳过快照过期清理),专注于数据写入。独立 Compaction Job 负责扫描表的文件状态,按需执行合并。
核心优势:
注意:
write-only = true会同时跳过快照过期清理,需要在 Compaction Job 中或通过其他方式定期执行快照过期,避免存储空间持续增长。
使用 Paimon 内置的存储过程执行 Compaction:
-- 运行环境 Flink 1.17 流式任务 -- 创建 Catalog,按需选择 LAS/FileSystem Catalog CREATE CATALOG paimon_catalog WITH ( 'type' = 'paimon', 'warehouse' = 'tos://paimon-dedicated-compaction/filesystem_catalog' ); USE CATALOG paimon_catalog; CALL sys.compact( 'paimon_db.write_only_paimon_0_8_2_20250911', -- 表名, 格式为 db.table '', -- 分区(可选,不填则不限制分区) '', -- 排序策略(可选) '', -- 排序字段(可选) 'consumer.mode=at-least-once,consumer-id=your-consumer-id,sort-spill-threshold=5,num-sorted-run.stop-trigger=2147483647' -- 任务级新增表参数(可选) );
通过提交 Paimon Action Jar 的方式执行 Compaction:
## --table_conf 为追加表级参数,视具体场景进行调整 compact --warehouse tos://paimon-dedicated-compaction/filesystem_catalog --database paimon_db --table write_only_paimon_0_8_2_20250911 --table_conf consumer.mode=at-least-once --table_conf consumer-id=your-consumer-id --table_conf sort-spill-threshold=5 --table_conf num-sorted-run.stop-trigger=2147483647
参数名 | 说明 | 推荐值 |
|---|---|---|
| 写入端跳过 Compaction |
|
| 禁用 Write Stall |
|
| Compaction 排序溢写阈值 |
|
| 不等待 Lookup Compaction |
|
| Compaction Job 扫描新分区的时间间隔 |
|
在火山引擎 TOS(对象存储) 场景下,由于 TOS FNS 不支持原生文件系统的并发锁,必须采取以下方案确保多 Job 写入安全:
flink.hadoop.fs.tos.rename.enabled: true。策略 | 适用模式 | 说明 |
|---|---|---|
| 仅 Batch 模式 | 将表(或指定分区)内所有文件合并为最少数量的文件,彻底消除小文件 |
| Streaming 和 Batch 均支持 | 基于 LSM 层级规则,仅合并满足条件的文件子集,对写入影响最小 |
在流式 Compaction 中,compact_strategy 默认为 minor;在批式 Compaction 中,默认为 full。
建议:大规模批式 Compaction 建议采用 Spark 进行。
对于分区表,历史分区(不再写入新数据的分区)可能积累大量小文件。可通过 partition_idle_time 参数在批式 Compaction 中只处理"已静止"的分区:
# 仅对超过 1 天未写入的分区执行 full Compaction flink run \ -D execution.runtime-mode=batch \ /path/to/paimon-flink-action-*.jar \ compact \ --warehouse tos://your-bucket/warehouse \ --database my_database \ --table my_table \ --partition_idle_time '1 d' \ --compact_strategy full
partition_idle_time参数仅在 Batch 模式下有效。在流式 Compaction 中该参数会被忽略。
Sort Compact 是一种特殊的 Compaction 模式,在合并文件的同时按指定列对数据进行全局排序,以优化特定查询模式下的扫描效率(数据聚集性)。
排序策略说明:
策略 | 说明 | 适用场景 |
|---|---|---|
| 按指定列普通排序(字典序) | 等值查询、范围查询单维度 |
| 多维 Z-order 曲线排序,多列联合优化 | 多维度范围查询(如经纬度、时间+地域) |
| Hilbert 曲线排序,多维度聚集性优于 zorder | 高基数多维度查询 |
-- 执行 Sort Compact(zorder 策略) SET 'execution.runtime-mode' = 'batch'; CALL sys.compact( `table` => 'my_catalog.my_database.events', `compact_strategy` => 'full', `order_by` => 'event_date,user_id', `order_strategy` => 'zorder' );
参数版本注意:Paimon 0.8.x 中使用
order_columns参数,新版本(master)中已改为order_by。请根据实际使用的 Paimon 版本选择正确参数名。
对于历史分区的定期整理,推荐通过以下方式实现定时调度(参考 jar 作业调度方案执行相关 flink jar):
#!/bin/bash # 每日凌晨 2 点对前一天的历史分区执行 full Compaction YESTERDAY=$(date -d "yesterday" +%Y-%m-%d) WAREHOUSE="tos://your-bucket/warehouse" flink run \ -D execution.runtime-mode=batch \ -D parallelism.default=8 \ /path/to/paimon-flink-action-*.jar \ compact \ --warehouse ${WAREHOUSE} \ --database my_database \ --table events \ --partition "dt=${YESTERDAY}" \ --compact_strategy full
可将上述脚本配置在 火山引擎 DataLeap 调度系统或 Crontab 中实现定时执行。
Spark 对湖仓表使用说明--E-MapReduce-火山引擎
Spark 写入 Paimon 的自动 Compaction 与 Flink 共用同一套 MergeTreeWriter + UniversalCompaction 底层引擎,核心区别在于 Spark batch 的 prepareCommit() 硬编码了 waitCompaction=true,确保每次提交前会同步等待所有已触发的 Compaction 都执行完毕,因此 Spark batch 不存在 Write Stall 的概念 —— 它天然是 "同步等待" 的。
Spark 写入过程:写完→等待 Compaction→提交 snapshot。
Spark 手动 Compaction 通过 CALL sys.compact(...) 存储过程执行,功能灵活,支持多种高级选项。
参数名 | 类型 | 必填 | 说明 |
|---|---|---|---|
| STRING | 是 | 表的完整三段式路径(catalog.database.table) |
| STRING | 否 | 分区过滤表达式,如 |
| STRING | 否 | Sort Compact 排序列(新版本参数名,旧版用 |
| STRING | 否 | 排序策略: |
| STRING | 否 | 合并策略: |
| STRING | 否 | 分区过滤谓词(SQL WHERE 语法) |
-- 示例 1:对整张表执行 full Compaction CALL sys.compact( table => 'paimon_catalog.my_database.orders' ); -- 示例 2:对指定分区执行 Compaction CALL sys.compact( table => 'paimon_catalog.my_database.orders', partitions => 'dt=2025-01-01' ); -- 示例 3:使用 WHERE 谓词过滤多个分区 CALL sys.compact( table => 'paimon_catalog.my_database.events', where => 'dt >= "2025-01-01" AND dt <= "2025-01-07"' ); -- 示例 4:Sort Compact(zorder 策略,需为 Append + bucket=-1 表) CALL sys.compact( table => 'paimon_catalog.my_database.events', order_strategy => 'zorder', order_by => 'event_date,user_id' ); -- 示例 5:指定 full 策略 + 分区过滤 CALL sys.compact( table => 'paimon_catalog.my_database.orders', compact_strategy => 'full', where => 'dt = "2025-01-01"' );
在火山引擎 EMR Serverless Spark(内置 Spark 3.5.1 + Paimon 0.8.1)环境中,CALL 存储过程默认不被支持,需要额外开启自定义 SQL 解析功能:
-- 在 Spark Session 中开启自定义 SQL 解析(EMR Serverless 必须) SET spark.sql.extensions = org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions; SET emr.serverless.spark.custom.parse.enabled = true;
或在提交 Spark 作业时通过配置参数传入:
spark-submit \ --conf spark.sql.extensions=org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions \ --conf emr.serverless.spark.custom.parse.enabled=true \ ...
能力维度 | Flink | Spark |
|---|---|---|
流式自动 Compaction | 支持 | 不支持 |
批式手动 Compaction | 支持(Action Jar / SQL Procedure) | 支持(CALL sys.compact) |
分离式 Compaction Job | 支持 | 不支持 |
Sort Compact | 支持(Batch 模式) | 支持(手动 CALL) |
历史分区 Compaction | 支持(partition_idle_time) | 支持(WHERE 分区过滤) |
Checkpoint 依赖 | 是 | 否 |
多库 Compaction | 支持(compact_database) | 不支持 |
实时增量 Compaction | 支持(combined 模式自动感知) | 不支持 |
适合场景 | 实时流式写入、多 Job 并发写 | 批量导入、离线 ETL、历史数据整理 |
问题描述: 单张主键表每秒写入数百万条记录,默认同步 Compaction 模式下频繁触发 Write Stall,导致写入链路不稳定。
推荐方案: 开启异步 Compaction + 适当放大 Write Buffer
-- 表配置 CREATE TABLE high_throughput_table ( id BIGINT, content STRING, ts TIMESTAMP(3), dt STRING, PRIMARY KEY (id, dt) NOT ENFORCED ) PARTITIONED BY (dt) WITH ( 'bucket' = '16', 'changelog-producer' = 'none', -- 异步 Compaction 'num-sorted-run.stop-trigger' = '2147483647', 'sort-spill-threshold' = '5', 'changelog-producer.lookup-wait' = 'false', -- 写入优化 'write-buffer-spillable' = 'true', 'write-buffer-size' = '256MB', 'local-merge-buffer-size' = '64MB' );
关键要点:
num-sorted-run.stop-trigger 设为最大值,彻底禁用 Write Stall。write-buffer-spillable,允许 Write Buffer 溢写磁盘,防止内存压力。local-merge-buffer-size(建议 64MB 起调)可在 Writer 本地预聚合,减少 Sorted Run 生成,但不支持 CDC 数据,使用前请确认数据类型。sink.parallelism 不应超过 Bucket 数量,否则多个 Writer 写同一 Bucket 会造成锁竞争。问题描述: 多个 Flink 写入 Job 同时写入同一张 Paimon 表(如多数据源合并写入),各 Job 自行执行 Compaction 导致文件锁冲突。
推荐方案: 分离式 Compaction 架构
架构图:
写入 Job 1 (write-only=true) ──┐ 写入 Job 2 (write-only=true) ──┼──> Paimon Table <── 独立 Compaction Job 写入 Job 3 (write-only=true) ──┘ (sys.compact / Action Jar)
写入 Job 配置:
-- 所有写入 Job 均设置 write-only=true INSERT INTO paimon_catalog.my_db.orders /*+ OPTIONS('write-only' = 'true') */ SELECT * FROM kafka_source;
独立 Compaction Job 配置(TOS 环境):
-- 运行环境 Flink 1.17 流式任务 -- 创建 Catalog,按需选择 LAS/FileSystem Catalog CREATE CATALOG paimon_catalog WITH ( 'type' = 'paimon', 'warehouse' = 'tos://paimon-dedicated-compaction/filesystem_catalog' ); USE CATALOG paimon_catalog; CALL sys.compact( 'paimon_db.write_only_paimon_0_8_2_20250911', -- 表名, 格式为 db.table '', -- 分区(可选,不填则不限制分区) '', -- 排序策略(可选) '', -- 排序字段(可选) 'consumer.mode=at-least-once,consumer-id=your-consumer-id,sort-spill-threshold=5,num-sorted-run.stop-trigger=2147483647' -- 任务级新增表参数(可选) );
注意事项:
write-only = true 会跳过快照过期清理,需在 Compaction Job 中或单独配置快照过期任务。问题描述: 实时写入表中,历史分区(如按天分区的前几天数据)积累大量小文件,影响下游批量读取查询性能。
推荐方案: 定时批式 Compaction + partition_idle_time 过滤
#!/bin/bash # 脚本:每日凌晨对超过 1 天未写入的分区执行 full Compaction WAREHOUSE="tos://your-bucket/warehouse" DATABASE="my_database" TABLE="orders" flink run \ -D execution.runtime-mode=batch \ -D parallelism.default=8 \ /path/to/paimon-flink-action-*.jar \ compact \ --warehouse ${WAREHOUSE} \ --database ${DATABASE} \ --table ${TABLE} \ --catalog-conf metastore=hive \ --catalog-conf uri=thrift://hive-metastore:9083 \ --catalog-conf lock.enabled=true \ --partition_idle_time '1 d' \ --compact_strategy full echo "历史分区 Compaction 完成:$(date)"
Spark 替代方案(使用 WHERE 谓词过滤):
-- 对 30 天前的分区批量执行 Compaction CALL sys.compact( table => 'paimon_catalog.my_database.orders', compact_strategy => 'full', where => 'dt < "2024-12-01"' );
问题描述: Append 表(Unaware-Bucket)写入的日志数据,下游按 event_date 和 user_id 做多维度范围查询,扫描效率低下。
推荐方案: 定期执行 zorder Sort Compact
前提条件:
bucket = -1)的 Append 表-- Flink SQL 执行 Sort Compact SET 'execution.runtime-mode' = 'batch'; CALL sys.compact( 'paimon_db.write_only_paimon_0_8_2_20250911', -- 表名, 格式为 db.table '', -- 分区(可选,不填则不限制分区) 'zorder', -- 排序策略(可选) 'event_date,user_id', -- 排序字段(可选) 'consumer.mode=at-least-once,consumer-id=your-consumer-id,sort-spill-threshold=5,num-sorted-run.stop-trigger=2147483647' -- 任务级新增表参数(可选) );
注意事项:
zorder 适合多维度等值或范围查询;若仅单维度查询,order 策略更简单高效。参数名 | 默认值 | 说明 | 调优建议 |
|---|---|---|---|
| 5 | Compaction 触发阈值(Sorted Run 数) | 高吞吐场景可适当提高至 8-10 |
| 8 | Write Stall 触发阈值 | 异步场景设为 2147483647 |
| 无 | Compaction 排序溢写阈值 | 建议 10 |
| true | 是否等待 Lookup Compaction | 高吞吐场景设 false |
| 256MB | Write Buffer 大小 | 根据内存适当调大 |
| none | 是否允许 Write Buffer 溢写 | 高吞吐场景建议 true |
| 无 | 本地预聚合缓冲区大小 | 非 CDC 场景建议 64MB+ |
| false | 跳过 Compaction(分离式方案) | 多 Job 并发写时写入端设 true |
| 128MB | 目标文件大小 | 通常无需调整 |
参数名 | 默认值 | 说明 | 调优建议 |
|---|---|---|---|
| 5 | 触发 Compaction 的最小文件数 | 小文件较多时可降低至 3 |
| 256MB | 目标文件大小 | 根据查询模式调整 |
| 200 | 空间放大率上限(%) | 通常无需调整 |
参数名 | 默认值 | 说明 |
|---|---|---|
| minor(流式)/ full(批式) | Compaction 策略 |
| 无 | 历史分区静止时间(仅 Batch 模式) |
| 10s | 流式 Compaction Job 扫描间隔 |