You need to enable JavaScript to run this app.
文档中心
流式计算 Flink版

流式计算 Flink版

复制全文
下载 pdf
Paimon Compaction
Paimon Compaction 概览
复制全文
下载 pdf
Paimon Compaction 概览

文档背景

Apache Paimon的Compaction是一种在后台自动合并数据文件的过程,其主要目的是通过减少文件数量、优化数据布局来提升查询性能,并在数据持续写入时维持稳定的读写效率。

  1. 常规 Compaction:Flink 任务写入 Paimon 的时候默认会自动的进行 Compaction,用户无需额外配置。是最常见的 Compaction 模式。
  2. 批作业 Compaction: 可以通过 Flink Batch 作业触发全量的 Compaction 任务,常用于对于实时性要求不高,但是对于查询性能要求极高的场景。
  3. 分离式 Compaction 作业 (Dedicated Compaction Job):为了避免Compaction在数据写入时争抢资源(可能引起写入吞吐量波动或提交冲突),用户可以配置 'write-only' = 'true' 来让写入作业跳过Compaction。然后,通过一个独立的Flink作业(专用Compaction作业)来专门负责Compaction。这种方式使得数据写入可以持续进行,尤其适用于多作业同时写入的场景。Flink 分离式 Compaction 策略

本文档适用于在火山引擎(Volcano Engine) 上使用 EMR、LAS 等服务的 Paimon 用户,涵盖以下场景:

  • 使用 Flink 进行实时流式写入和 Compaction 管理
  • 使用 Spark 进行批量写入和离线 Compaction
  • 使用 TOS 对象存储作为底层存储的分布式部署
  • 需要对历史分区、多并发写入等场景进行专项优化

目标读者:大数据平台工程师、数据仓库开发工程师、DataOps 工程师,建议具备基本的 Flink/Spark 使用经验和 SQL 基础。

1. 概述

本节主要介绍 Paimon Compaction 的机制和核心技术要点。通过介绍 Paimon Compaction 的技术要点为用户提供更全面的 Compaction 能力,帮助用户更好地理解 Compaction 的技术细节,并提供更灵活的 Compaction 配置方式,以满足不同场景的需求。

1.1 核心要点

  • Compaction 是 Paimon 保障读写性能的核心机制:无论主键表(LSM Tree + Universal Compaction)还是 Append 表(Coordinator/Worker 拓扑),Compaction 都是控制小文件数量、提升查询效率的关键手段。
  • 主键表与 Append 表的 Compaction 机制本质不同:默认情况下,主键表 Compaction 与写入互相影响,Write Stall 可能造成反压,可以通过异步 Compaction 解耦;Append 表(Unaware-Bucket)通过独立 Coordinator + Worker 拓扑实现永不反压的异步合并。
  • 分离式 Compaction 是多 Job 并发写入的首选方案:通过 write-only = true 将写入端与 Compaction 解耦,消除多 Job 竞争同一分区时的合并冲突,适合高并发实时写入场景。

Image

  • Flink 与 Spark 的 Compaction 能力互补:Flink 支持流式自动 Compaction、异步 Compaction 和分离式 Compaction Job;Spark 不支持流式自动 Compaction,但其批式 CALL 存储过程灵活支持 Sort Compact、分区过滤等高级能力,适合离线 ETL 和历史分区整理。
  • 如果需要多任务同时写入 / Compaction 一张 Paimon 表,火山引擎环境有特有配置要求
    • 方案一:在 TOS 对象存储场景下需要保证写入目标桶是 HNS 类型
    • 方案二: TOS FNS 桶需要开启 Rename 原子性,Flink 配置自定义参数 flink.hadoop.fs.tos.rename.enabled: true

Image

1.2 Compaction 的背景与必要性

Apache Paimon 是一种面向流批一体场景的数据湖存储格式,其主键表采用 LSM Tree(Log-Structured Merge Tree) 结构存储数据。LSM Tree 将写入操作转化为顺序追加,极大提升了写入吞吐,但随着数据的持续写入,存储层会积累大量层级各异、大小不一的小文件(Sorted Run),给读取查询带来严重的 I/O 放大问题。
Image

Compaction(合并压缩) 是 LSM Tree 架构下的必要维护手段,其核心目标包括:

  • 减少文件数量:将多个小 Sorted Run 合并为更少、更大的文件,降低查询时的文件扫描开销。
  • 消除重复数据:对于主键表,合并过程中进行去重和排序,确保每个主键只保留最新版本。
  • 提升查询效率:减少 Merge-on-Read 的代价,接近 Copy-on-Write 的读取性能。
  • 控制存储空间放大:防止过期数据和冗余文件长期占用存储。

对于 Append 表(仅追加表),Compaction 的主要目的是将写入过程中产生的大量小文件合并为较少的大文件,提升下游读取的扫描效率,同时为 Sort Compact 等高级特性提供基础。

2. Paimon 表类型与 Compaction 机制差异

2.1 主键表(Primary Key Table)的 Compaction 机制

LSM Tree 存储结构与 Sorted Run 概念

主键表以 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 生产和快照过期清理等附加工作。

Universal Compaction 策略

Paimon 主键表默认采用 Universal Compaction 策略(类似 RocksDB 的 Universal Compaction)。该策略基于以下规则触发合并:

  1. 当 Sorted Run 总数超过 num-sorted-run.compaction-trigger(默认 5)时,触发一次合并。
  2. 每次合并选取满足条件的 Sorted Run 子集,从最旧的文件开始合并,直到剩余 Sorted Run 数量降至阈值以下。
  3. 合并操作在 Writer 内部异步执行,不阻塞正常写入(但可能触发 Write Stall)。

Lookup Compaction

当使用 changelog-producer = lookup 时,Paimon 会在 Compaction 前执行 Lookup 操作以生成 Changelog,通过主键查询数据更新前后的值,在 Compaction 的过程中产生 -U 和 +U 的真实 changelog。这种 changelog-producer 会增加 Compaction 的成本。

  1. Paimon 上游数据并非 CDC 数据,比如来自 Kafka 的 JSON 数据。需要主键进行更新后,产出 changelog。
  2. 在 Partial-Update 和 Aggregation 表的条件下,必须使用 lookup 或者 full-compaction。

触发条件与 Write Stall 机制

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 影响写入吞吐。

关键参数表

参数名

默认值

说明

num-sorted-run.compaction-trigger

5

Sorted Run 数量超过此值时触发 Compaction

num-sorted-run.stop-trigger

compaction-trigger + 3 (即 8)

Sorted Run 数量超过此值时触发 Write Stall

changelog-producer

none

Changelog 生产方式:none / input / full-compaction / lookup

write-buffer-size

256MB

Write Buffer 大小

target-file-size

128MB 主键表
256MB APPEND 表

目标文件大小

write-only

false

开启后跳过 Compaction(用于分离式方案)

2.2 Append 表(Append-Only Table)的 Compaction 机制

Append 表不支持主键更新,数据仅追加不修改。根据是否指定 Bucket 数量,Append 表分为两种子类型,其 Compaction 机制存在本质差异。

Unaware-Bucket 表(bucket = -1)的 Compaction

bucket = -1 时,表为 Unaware-Bucket 模式。Flink 写入任务中会自动启动一个 Compact Coordinator 和若干 Compact Worker 算子,形成独立的 Compaction 拓扑:

  • Compact Coordinator:扫描各分区的文件状态,根据触发规则(文件数、文件大小)生成合并任务,分发给 Compact Worker。
  • Compact Worker:接收合并任务,执行实际的文件合并操作,完成后提交结果。

这种拓扑设计的核心优势是:Compaction 完全异步执行,永不对写入造成反压。Compaction Worker 的失败不会影响写入链路,写入与合并天然解耦。

Bucketed Append 表(bucket = N)的 Compaction

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 的触发规则如下:

  • 正常触发:当某个分区/Bucket 内文件数 ≥ compaction.min.file-num,且这些文件的总大小 ≥ target-file-size 时触发。另外根据放大率等参数,判断是否触发合并。

参数名

默认值

说明

compaction.min.file-num

5

触发 Compaction 的最小文件数

target-file-size

128MB 主键表
256MB APPEND 表

目标合并文件大小

2.3 主键表 vs Append 表 Compaction 对比总结

对比维度

主键表(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、更新类数据

日志、事件流、追加类数据

Image

Flink 写入 Paimon 表时,默认在写入算子(Sink Writer)内部自动执行 Compaction。这是最简单的使用方式,无需额外配置,适合写入吞吐适中、对延迟不敏感的场景。
默认行为说明:

  • 每次 Checkpoint 完成后,Writer 会检查当前 Sorted Run 数量,若超过 num-sorted-run.compaction-trigger(默认 5),则触发一次 Compaction。
  • Compaction 在 Writer 线程内异步执行,但若 Sorted Run 数量超过 num-sorted-run.stop-trigger(默认 8),则触发 Write Stall,暂停写入,直到 Compaction 赶上。
  • Write Stall 期间,上游数据积压,可能触发反压,进而影响整个写入链路的吞吐。

适用场景: 单一 Flink 写入 Job、写入并发不高、数据量适中的在线实时入库场景。

当写入吞吐较高,同步 Compaction 模式下频繁 Write Stall 影响稳定性时,可开启异步 Compaction 模式。异步 Compaction 通过禁用 Write Stall 并放宽 Compaction 触发阈值,让写入尽可能不受 Compaction 影响。

3 个核心参数详解

参数 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 的生成会有一定延迟,下游消费者可能需要等待更长时间才能看到最新变更。

完整 SQL 配置示例

-- 创建开启异步 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'
);

适用场景: 高吞吐实时写入(单表每秒数百万条记录)、对写入延迟敏感但对读取稍有容忍的场景。

3.3 独立 Compaction 任务(分离式 Compaction)

概念与原理

分离式 Compaction 是将 Compaction 工作从写入 Job 中完全剥离,由专用的独立 Compaction Job 承担的架构模式。
写入端通过设置 write-only = true 跳过所有 Compaction 操作(同时也会跳过快照过期清理),专注于数据写入。独立 Compaction Job 负责扫描表的文件状态,按需执行合并。
核心优势:

  • 消除多 Job 竞争:多个 Flink Job 并发写同一分区时,若各 Job 均自行执行 Compaction,会产生文件锁竞争和重复合并。分离式方案将 Compaction 收归一处,彻底消除竞争。
  • 资源独立:Compaction 消耗的 CPU/内存资源不影响写入 Job 的稳定性,可独立扩缩容。
  • 写入最简化:写入 Job 只做最基本的 flush 操作,逻辑简单,故障域缩小。

注意:write-only = true 会同时跳过快照过期清理,需要在 Compaction Job 中或通过其他方式定期执行快照过期,避免存储空间持续增长。

方式一:Flink SQL Procedure

使用 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' -- 任务级新增表参数(可选)
);

方式二:Flink Action Jar

通过提交 Paimon Action Jar 的方式执行 Compaction:

Image

##  --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

关键参数配置

参数名

说明

推荐值

write-only

写入端跳过 Compaction

true(写入 Job 中设置)

num-sorted-run.stop-trigger

禁用 Write Stall

2147483647

sort-spill-threshold

Compaction 排序溢写阈值

10

changelog-producer.lookup-wait

不等待 Lookup Compaction

false

continuous.discovery-interval

Compaction Job 扫描新分区的时间间隔

1 min(流式模式,如果不敏感的话可以降低频率)

火山引擎环境注意事项

在火山引擎 TOS(对象存储) 场景下,由于 TOS FNS 不支持原生文件系统的并发锁,必须采取以下方案确保多 Job 写入安全:

  • 方案一:在 TOS 对象存储场景下需要保证写入目标桶是 HNS 类型
  • 方案二: TOS FNS 桶需要开启 Rename 原子性,Flink 配置自定义参数 flink.hadoop.fs.tos.rename.enabled: true

3.4 批式 Compaction 任务

full vs minor 策略对比

策略

适用模式

说明

full

仅 Batch 模式

将表(或指定分区)内所有文件合并为最少数量的文件,彻底消除小文件

minor

Streaming 和 Batch 均支持

基于 LSM 层级规则,仅合并满足条件的文件子集,对写入影响最小

在流式 Compaction 中,compact_strategy 默认为 minor;在批式 Compaction 中,默认为 full
建议:大规模批式 Compaction 建议采用 Spark 进行。

历史分区 Compaction(partition_idle_time)

对于分区表,历史分区(不再写入新数据的分区)可能积累大量小文件。可通过 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(zorder / hilbert / order)

Sort Compact 是一种特殊的 Compaction 模式,在合并文件的同时按指定列对数据进行全局排序,以优化特定查询模式下的扫描效率(数据聚集性)。
排序策略说明:

策略

说明

适用场景

order

按指定列普通排序(字典序)

等值查询、范围查询单维度

zorder

多维 Z-order 曲线排序,多列联合优化

多维度范围查询(如经纬度、时间+地域)

hilbert

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 中实现定时执行。

4. Spark Compaction 机制

Spark 对湖仓表使用说明--E-MapReduce-火山引擎

4.1 Spark 自动 Compaction

Spark 写入 Paimon 的自动 Compaction 与 Flink 共用同一套 MergeTreeWriter + UniversalCompaction 底层引擎,核心区别在于 Spark batch 的 prepareCommit() 硬编码了 waitCompaction=true,确保每次提交前会同步等待所有已触发的 Compaction 都执行完毕,因此 Spark batch 不存在 Write Stall 的概念 —— 它天然是 "同步等待" 的。
Spark 写入过程:写完→等待 Compaction→提交 snapshot。

4.2 Spark 手动 Compaction

Spark 手动 Compaction 通过 CALL sys.compact(...) 存储过程执行,功能灵活,支持多种高级选项。

完整参数说明表

参数名

类型

必填

说明

table

STRING

表的完整三段式路径(catalog.database.table)

partitions

STRING

分区过滤表达式,如 'dt=2025-01-01'

order_by

STRING

Sort Compact 排序列(新版本参数名,旧版用 order_columns

order_strategy

STRING

排序策略:order / zorder / hilbert

compact_strategy

STRING

合并策略:full / minor

where

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 环境特有配置

在火山引擎 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、历史数据整理

5. 最佳实践与场景推荐

5.1 场景一:大规模实时写入(高吞吐优先)

问题描述: 单张主键表每秒写入数百万条记录,默认同步 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 数据,使用前请确认数据类型。
  • 单个 Bucket 数据量建议控制在 1GB 左右,可根据此原则估算 Bucket 数量。
  • sink.parallelism 不应超过 Bucket 数量,否则多个 Writer 写同一 Bucket 会造成锁竞争。

5.2 场景二:多 Job 并发写入同一表

问题描述: 多个 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' -- 任务级新增表参数(可选)
);

注意事项:

  • TOS 必须选择 HNS 桶,或者开启原子 RENAME 功能。
  • write-only = true 会跳过快照过期清理,需在 Compaction Job 中或单独配置快照过期任务。

5.3 场景三:历史分区定期整理

问题描述: 实时写入表中,历史分区(如按天分区的前几天数据)积累大量小文件,影响下游批量读取查询性能。
推荐方案: 定时批式 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"'
);

5.4 场景四:查询性能优化(Sort Compact)

问题描述: Append 表(Unaware-Bucket)写入的日志数据,下游按 event_dateuser_id 做多维度范围查询,扫描效率低下。
推荐方案: 定期执行 zorder Sort Compact
前提条件:

  • 表类型必须为 Unaware-Bucket(bucket = -1)的 Append 表
  • 执行模式必须为 Batch
-- 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' -- 任务级新增表参数(可选)
);

注意事项:

  • Sort Compact 会对表中所有数据(或指定分区)进行重新排序和重写,计算资源消耗大,建议在业务低峰期执行。
  • zorder 适合多维度等值或范围查询;若仅单维度查询,order 策略更简单高效。

5.5 关键参数速查表(汇总)

主键表核心参数

参数名

默认值

说明

调优建议

num-sorted-run.compaction-trigger

5

Compaction 触发阈值(Sorted Run 数)

高吞吐场景可适当提高至 8-10

num-sorted-run.stop-trigger

8

Write Stall 触发阈值

异步场景设为 2147483647

sort-spill-threshold

Compaction 排序溢写阈值

建议 10

changelog-producer.lookup-wait

true

是否等待 Lookup Compaction

高吞吐场景设 false

write-buffer-size

256MB

Write Buffer 大小

根据内存适当调大

write-buffer-spillable

none
(1.1 版本行为)

是否允许 Write Buffer 溢写

高吞吐场景建议 true

local-merge-buffer-size

本地预聚合缓冲区大小

非 CDC 场景建议 64MB+

write-only

false

跳过 Compaction(分离式方案)

多 Job 并发写时写入端设 true

target-file-size

128MB

目标文件大小

通常无需调整

Append 表核心参数

参数名

默认值

说明

调优建议

compaction.min.file-num

5

触发 Compaction 的最小文件数

小文件较多时可降低至 3

target-file-size

256MB

目标文件大小

根据查询模式调整

compaction.max-size-amplification-percent

200

空间放大率上限(%)

通常无需调整

通用 Compaction 参数

参数名

默认值

说明

compact_strategy

minor(流式)/ full(批式)

Compaction 策略

partition_idle_time

历史分区静止时间(仅 Batch 模式)

continuous.discovery-interval

10s

流式 Compaction Job 扫描间隔

最近更新时间:2026.07.09 22:03:27
这个页面对您有帮助吗?
有用
有用
无用
无用