Apache Paimon Changelog Producer 详解

Apache Paimon Changelog Producer 详解

本文内容由 AI 辅助生成,已经人工审核和编辑。

Apache Paimon Changelog Producer 详解

总结

  • 更新时间:2026-06-18
  • 本文档深入解析 Paimon 主键表的四种 Changelog Producer 模式(none / input / full-compaction / lookup)的工作机制与源码实现
  • none 模式:不产生 changelog 文件,下游流读依赖 Compaction 前后的 delta 数据文件感知变化,延迟低但信息有损(只能感知最终状态,无中间 BEFORE 值)
  • input 模式:在 WriteBuffer flush 时,将原始 RowKind 信息双写到 changelog 文件,零延迟、无额外 IO,但要求上游显式提供完整的 UPDATE_BEFORE / UPDATE_AFTER 对(适合 CDC 场景)
  • full-compaction 模式:等待 Full Compaction 完成后,通过对比 L(max-1) 层旧值(topLevelKv)与合并后新值,生成 INSERT / UPDATE_BEFORE / UPDATE_AFTER / DELETE changelog,有较大写入延迟(受 full-compaction.delta-commits 控制)
  • lookup 模式:在每次涉及 Level-0 的 Compaction 时,通过查询 LookupLevels 获取高层旧值,实时生成精确 changelog,延迟低于 full-compaction,但需要额外的本地磁盘 Lookup 文件
  • 核心选择原则:CDC 同步用 input;需要低延迟精确 changelog 用 lookup;可接受延迟但希望降低写放大用 full-compaction;不需要 changelog 用 none
  • 源码基于 Apache Paimon 1.3.1(release-1.3 分支)

一、基础概念

1.1 什么是 Changelog

Changelog 是一条带有变更类型(RowKind)的数据记录,描述某一行数据从旧值到新值的变化过程。

Paimon 使用 Flink 的 RowKind 枚举表达四种变更语义:

RowKind 含义 符号
INSERT 新插入一行 +I
UPDATE_BEFORE 更新前的旧值 -U
UPDATE_AFTER 更新后的新值 +U
DELETE 删除一行 -D

一次 UPDATE 操作在 changelog 中通常以 -U(旧值)+ +U(新值)的成对形式出现,下游算子(如 Flink 聚合)正是依赖这对信息来正确维护状态。

1.2 主键表为何需要 Producer

Paimon 主键表采用 LSM 树结构:数据以追加形式写入 Level-0,通过 Compaction 逐步合并到更高层级。同一主键的多条记录分散在不同层级的 SST 文件中,文件层面本身并不记录"旧值是什么"

因此,下游若要流式消费 Paimon 表的变更(如 Flink 流读、Lookup Join),必须有额外机制提供 UPDATE_BEFORE 信息——这就是 Changelog Producer 的核心价值:以不同的代价和精度,为下游提供完整的变更语义

1.3 四种模式一览

模式 生成时机 是否需要旧值 延迟 写放大 典型场景
none 不生成 仅追加、无需精确 changelog
input WriteBuffer flush 时 否(上游提供) 极低 低(多写一份) CDC 同步
full-compaction Full Compaction 完成后 是(从 LSM 最高层读取) 高(受 delta-commits 控制) 低(复用 Compaction IO) 可接受延迟的批同步
lookup 涉及 L0 的 Compaction 时 是(Lookup 查询高层) 中(额外本地 Lookup 文件) 需要低延迟精确 changelog

二、none 模式

2.1 概念

none 是默认模式,不产生任何 changelog 文件。表在每次写入 checkpoint 后只有数据文件(Data Files)和对应的 Manifest,没有额外的 changelog-* 文件。

下游流读(DataTableStreamScan)使用 DeltaFollowUpScanner 感知变化:只扫描 CommitKind = APPEND 的快照,读取这次 checkpoint 新增的数据文件(delta 文件)。

这意味着下游收到的记录只有追加语义,没有 UPDATE_BEFORE。如果同一主键在一个 checkpoint 内发生了 UPDATE,下游最终只看到最新的 +I 记录(因为 WriteBuffer 在 flush 前先做 MergeFunction 合并)。

2.2 读取端工作机制

DataTableStreamScan.createFollowUpScanner()
  └─ case NONE → new DeltaFollowUpScanner()

DeltaFollowUpScanner.shouldScanSnapshot(snapshot)
  └─ 只接受 CommitKind == APPEND 的快照(跳过 COMPACT 快照)

DeltaFollowUpScanner.scan(snapshot, snapshotReader)
  └─ snapshotReader.withMode(ScanMode.DELTA).withSnapshot(snapshot).read()
     └─ 读取本次 APPEND 快照新增的 Level-0 数据文件

关键点:COMPACT 快照被跳过。这意味着如果一条数据在写入后立即被 Compaction 合并,none 模式下的流读可能会错过这次 Compaction 带来的文件变化(但不会漏数据,因为 delta 扫描读的是原始追加文件)。

2.3 源码定位

  • Scanner 创建DataTableStreamScan.java:277
    case NONE:
        followUpScanner = new DeltaFollowUpScanner();
    
  • 快照过滤DeltaFollowUpScanner.java:34
    public boolean shouldScanSnapshot(Snapshot snapshot) {
        if (snapshot.commitKind() == Snapshot.CommitKind.APPEND) {
            return true;
        }
        // ...
        return false;
    }
    

2.4 适用限制

  • ❌ 下游无法获得 UPDATE_BEFORE,不能用于需要精确 changelog 的场景(如 Flink 聚合回撤)
  • ✅ 无额外存储开销,适合只需要最新状态的场景(如直接覆盖写入)
  • ✅ Flink Lookup Join 在 none 模式下使用 COMPACT_DELTA_MONITOR 策略,通过扫描 Compaction 前后的 diff 感知变化

三、input 模式

3.1 概念

input 模式的核心思想:上游数据本身已经携带 RowKind 信息(如来自 CDC),Paimon 只需如实透传,不做任何计算

写入时,在 WriteBuffer flush 到磁盘的同时,双写一份 changelog 文件(changelog-*)。这份 changelog 文件直接保存原始输入中的每一条带 RowKind 的记录(未经 MergeFunction 合并)。

3.2 触发时机与写入流程

WriteBuffer flush 发生在两种场景:

  1. WriteBuffer 内存满时(自动 flush)
  2. Flink checkpoint 触发时(prepareCommitflushWriteBuffer
MergeTreeWriter.flushWriteBuffer(waitForLatestCompaction, forcedFullCompaction)
  │
  ├─ // input 模式:创建 changelogWriter
  │  final RollingFileWriter<KeyValue, DataFileMeta> changelogWriter =
  │      changelogProducer == ChangelogProducer.INPUT
  │          ? writerFactory.createRollingChangelogFileWriter(0)  // 写到 Level-0
  │          : null;
  │
  ├─ final RollingFileWriter<KeyValue, DataFileMeta> dataWriter =
  │      writerFactory.createRollingMergeTreeFileWriter(0, FileSource.APPEND);
  │
  └─ writeBuffer.forEach(
         keyComparator,
         mergeFunction,
         changelogWriter == null ? null : changelogWriter::write,  // rawConsumer
         dataWriter::write                                          // mergedConsumer
     )

forEach 的双消费者机制是 input 模式的关键:

// WriteBuffer.java
void forEach(
    Comparator<InternalRow> keyComparator,
    MergeFunction<KeyValue> mergeFunction,
    @Nullable KvConsumer rawConsumer,   // 接收原始(未合并)记录 → changelog 文件
    KvConsumer mergedConsumer            // 接收合并后记录 → 数据文件
)

SortBufferWriteBuffer.MergeIterator.readOnce() 中:

// SortBufferWriteBuffer.java:265
if (rawConsumer != null) {
    rawConsumer.accept(current.getReusedKv());  // 原始记录写入 changelog
}

每条记录从 SortBuffer 读出时,先交给 rawConsumer(写到 changelog 文件),再参与 MergeFunction 计算(最终写到数据文件)。这样:

  • changelog 文件:保留所有原始的 -U / +U / -D / +I 记录
  • 数据文件:只保留 MergeFunction 合并后的最终结果

3.3 文件产物

snapshot-N (CommitKind: APPEND)
├── manifest-list (包含 data files 和 changelog files)
├── data/partition/bucket-x/
│   └── data-xxxxx.orc          ← 合并后的数据文件(Level-0)
└── data/partition/bucket-x/
    └── changelog-xxxxx.orc     ← 原始 changelog 文件(Level-0)

Changelog 文件与数据文件存储在同一目录,文件名以 changelog- 前缀区分。

3.4 下游读取

DataTableStreamScan → ChangelogFollowUpScanner
  └─ snapshot.changelogManifestList() != null → 扫描 changelog 文件

ChangelogFollowUpScanner 直接扫描快照中的 changelogManifestList 对应的 changelog 文件,将 RowKind 信息完整传递给下游。

3.5 源码定位

  • 双写逻辑MergeTreeWriter.java:216
  • rawConsumer 调用SortBufferWriteBuffer.java:265
  • 模式判断MergeTreeWriter.java:217
    changelogProducer == ChangelogProducer.INPUT
        ? writerFactory.createRollingChangelogFileWriter(0)
        : null;
    

3.6 约束与要求

  • 零延迟:changelog 随数据文件同步生成,不需要等待 Compaction
  • 写放大最小:只额外写一份 changelog 文件,无需读取现有数据
  • 上游必须提供完整 RowKind:Paimon 不推断旧值,UPDATE_BEFORE 必须由上游显式产生
  • 不支持与 deletion-vectors.enabled 同用:Deletion Vectors 模式要求 changelog-producer = 'none''lookup'
  • 不支持 upsert 语义的上游:如果上游只发 +I,下游仍然只收到 +I,没有 -U

四、full-compaction 模式

4.1 概念

full-compaction 模式不在写入时生成 changelog,而是等待 Full Compaction(将所有数据合并到最高层)完成后,通过对比合并前后的值来生成 changelog。

其核心逻辑:当一条主键记录参与 Full Compaction(输出到 maxLevel)时,将该主键在最高层已有的旧值topLevelKv)与合并后的新值merged)做比较,生成 INSERT / UPDATE_BEFORE+UPDATE_AFTER / DELETE 事件。

4.2 Full Compaction 的触发时机

Full Compaction 是将所有层级的数据合并输出到 maxLevel(通常是 Level-6)的操作。在 full-compaction 模式下,有两种控制方式:

方式一:full-compaction.delta-commits(推荐)

// GlobalFullCompactionSinkWrite.java:172
if (!writtenBuckets.isEmpty() && isFullCompactedIdentifier(checkpointId, deltaCommits)) {
    waitCompaction = true;  // 强制本次 checkpoint 等待 Full Compaction 完成
}

isFullCompactedIdentifier(checkpointId, deltaCommits)checkpointId % deltaCommits == 0 时返回 true,意味着每隔 N 个 checkpoint 执行一次 Full Compaction

例如 full-compaction.delta-commits = 3,则第 3、6、9... 个 checkpoint 时,所有有写入的 bucket 都会强制触发 Full Compaction,并等待其完成后才提交快照。

方式二:changelog-producer.full-compaction.trigger-interval

按时间间隔折算成 delta-commits,原理相同。

4.3 changelog 生成机制(FullChangelogMergeFunctionWrapper)

Full Compaction 时,每个参与合并的主键记录经过 FullChangelogMergeFunctionWrapper.getResult() 处理:

// FullChangelogMergeFunctionWrapper.java

// topLevelKv: 该主键在 maxLevel(最高层)已有的旧值(若无则为 null)
// merged:     MergeFunction 合并后的最新值

public ChangelogResult getResult() {
    reusedResult.reset();
    if (isInitialized) {  // 有多条记录参与合并(真正的 UPDATE 或 DELETE 场景)
        KeyValue merged = mergeFunction.getResult();
        if (topLevelKv == null) {
            // 旧的最高层没有此主键 → 新增
            if (merged.isAdd()) {
                reusedResult.addChangelog(INSERT);      // +I
            }
        } else {
            if (!merged.isAdd()) {
                // 合并结果是删除 → 输出 DELETE
                reusedResult.addChangelog(DELETE, topLevelKv);   // -D
            } else if (!valueEqualiser.equals(topLevelKv.value(), merged.value())) {
                // 值发生变化 → 输出 UPDATE_BEFORE + UPDATE_AFTER
                reusedResult.addChangelog(UPDATE_BEFORE, topLevelKv);  // -U
                reusedResult.addChangelog(UPDATE_AFTER, merged);       // +U
            }
            // 若值相同(row-deduplicate 开启时会跳过),则不输出 changelog
        }
    } else {
        // 只有一条记录(且来自非最高层)→ 该主键第一次出现在最高层
        if (topLevelKv == null && initialKv.isAdd()) {
            reusedResult.addChangelog(INSERT, initialKv);  // +I
        }
    }
}

判断旧值的关键topLevelKv 是 Full Compaction 参与合并的来自 maxLevel 的那条记录。在 LSM 中,只有 Full Compaction 才会向最高层写入数据,因此 maxLevel 的记录就是"上一次 Full Compaction 写入的最终状态",天然代表旧值。

4.4 仅在 Full Compaction 时生成 changelog

FullChangelogMergeTreeCompactRewriter.rewriteChangelog() 控制何时生成 changelog:

// FullChangelogMergeTreeCompactRewriter.java:73
protected boolean rewriteChangelog(int outputLevel, boolean dropDelete, List<List<SortedRun>> sections) {
    boolean changelog = outputLevel == maxLevel;  // 只在输出到最高层时才生成 changelog
    // ...
    return changelog;
}

普通的 Minor Compaction(L0→L1 或 L1→L2 等)不会生成 changelog,只有最终合并到 maxLevel 的 Full Compaction 才会。

4.5 SinkWrite 的特殊处理

full-compaction 模式使用 GlobalFullCompactionSinkWrite 而非普通的 StoreSinkWriteImpl

StoreSinkWrite.createSinkWrite()
  ├─ changelogProducer == FULL_COMPACTION → GlobalFullCompactionSinkWrite
  ├─ needLookup (LOOKUP) → LookupSinkWrite
  └─ 其他 → StoreSinkWriteImpl

GlobalFullCompactionSinkWrite 的额外职责:

  1. 跟踪有写入的 bucketwrittenBuckets
  2. 在指定 checkpoint 时统一提交 Full Compaction 任务submitFullCompaction
  3. 等待所有 bucket 的 Full Compaction 完成才提交快照

4.6 初始读取的特殊处理

流读启动时(tryFirstPlan()),full-compaction 模式只读最高层数据:

// DataTableStreamScan.java:164
} else if (options.changelogProducer().equals(FULL_COMPACTION)) {
    result = startingScanner.scan(
        snapshotReader.withLevelFilter(level -> level == options.numLevels() - 1)
    );
}

这样保证:读取初始快照时只取每个主键的最终状态(最高层),避免重复处理低层的中间值。

4.7 文件产物

snapshot-N (CommitKind: COMPACT)
├── manifest-list
└── data/partition/bucket-x/
    ├── data-xxxxx.orc            ← Full Compaction 输出的 maxLevel 数据文件
    └── changelog-xxxxx.orc      ← 同步生成的 changelog 文件(记录 INSERT/UPDATE/DELETE)

4.8 配置示例

-- ✅ 建议配置
CREATE TABLE my_table (
    id BIGINT,
    name STRING,
    dt STRING
) TBLPROPERTIES (
    'primary-key' = 'id,dt',
    'changelog-producer' = 'full-compaction',
    'full-compaction.delta-commits' = '3'   -- 每 3 个 checkpoint 触发一次 Full Compaction
);

4.9 特性总结

  • ✅ 上游无需携带 RowKind,Paimon 自行推断变更
  • ✅ 写放大相对较低(changelog 生成复用 Compaction IO)
  • 有延迟:changelog 最多延迟 delta-commits 个 checkpoint 才到达下游
  • 中间状态丢失:两次 Full Compaction 之间的多次 UPDATE,下游只看到最终的一次变更
  • ❌ 要求 Flink 流式写入(batch 写入无法保证 Full Compaction 的周期控制)

五、lookup 模式

5.1 概念

lookup 模式是实时性最强的 changelog 生成方式:在每次涉及 Level-0 文件的 Compaction 时,通过查询(Lookup)高层已有数据来获取旧值,从而生成精确的 changelog

full-compaction 的根本区别:

  • full-compaction:等待数据到达 maxLevel 才生成 changelog,延迟高
  • lookup:任何 涉及 L0 的 Compaction(L0→L1、L0→L0、L0→maxLevel 均可)都实时生成 changelog,延迟低

5.2 什么情况触发 changelog 生成

// ChangelogMergeTreeRewriter.java:85
protected boolean rewriteLookupChangelog(int outputLevel, List<List<SortedRun>> sections) {
    if (outputLevel == 0) {
        return false;  // 输出到 L0 的 compaction 不生成 changelog
    }
    for (List<SortedRun> runs : sections) {
        for (SortedRun run : runs) {
            for (DataFileMeta file : run.files()) {
                if (file.level() == 0) {
                    return true;  // 只要参与合并的文件中有 L0 文件,就生成 changelog
                }
            }
        }
    }
    return false;
}

只要 Compaction 输入中包含 L0 文件,且输出层级 > 0,就会触发 lookup changelog 生成。这保证了每次 L0 数据向上推进时都能捕获变更。

5.3 Lookup 查旧值的机制(LookupLevels)

当一条 L0 记录参与 Compaction,但 Compaction 的其他输入文件(同一 key range 的高层文件)中没有该主键的旧值时,必须通过 LookupLevels 从更高层查找:

// LookupChangelogMergeFunctionWrapper.java:109
if (highLevel == null) {
    // 参与本次 compaction 的高层文件中无此 key,去更高层查
    T lookupResult = lookup.apply(mergeFunction.key());
    // lookup 实际调用:lookupLevels.lookup(key, outputLevel + 1)
    // 即从 (outputLevel + 1) 层开始向上搜索
}

LookupLevels.lookup() 的实现:

  1. 遍历从 startLevel 开始的每一层
  2. 对每个 SST 文件,先检查 Bloom Filter(快速排除不存在的 key)
  3. 若 Bloom Filter 命中,将 SST 文件转换为本地 LookupFile(类似 RocksDB 的 SST 格式,存储在 TaskManager 本地磁盘)
  4. 在 LookupFile 中按 key 查找,返回对应的 KeyValue
// LookupLevels.java:128
private T lookup(InternalRow key, DataFileMeta file) throws IOException {
    LookupFile lookupFile = lookupFileCache.getIfPresent(file.fileName());
    if (lookupFile == null) {
        lookupFile = createLookupFile(file);  // 第一次访问,构建本地索引文件
        // ...
    }
    byte[] keyBytes = keySerializer.serializeToBytes(key);
    byte[] valueBytes = lookupFile.get(keyBytes);  // O(1) 查找
    // ...
}

LookupFile 缓存:通过 Caffeine Cache 缓存 SST 文件对应的本地索引文件,避免重复构建。缓存大小由 lookup.cache-max-disk-size(默认 10GB)和 lookup.cache-file-retention(默认 1h)控制。

5.4 完整 changelog 生成逻辑

// LookupChangelogMergeFunctionWrapper.java

public ChangelogResult getResult() {
    // 1. 从本次 compaction 输入中找最高层(非 L0)的记录作为候选旧值
    KeyValue highLevel = mergeFunction.pickHighLevel();
    boolean containLevel0 = mergeFunction.containLevel0();

    // 2. 若没有高层记录,向更高层 Lookup
    if (highLevel == null) {
        T lookupResult = lookup.apply(mergeFunction.key());
        if (lookupResult != null) {
            highLevel = (KeyValue) lookupResult;
            mergeFunction.insertInto(highLevel, comparator);
        }
    }

    // 3. 计算合并结果
    KeyValue result = mergeFunction.getResult();

    // 4. 只有包含 L0 记录时才生成 changelog
    if (containLevel0 && lookupStrategy.produceChangelog) {
        setChangelog(highLevel, result);
    }

    return reusedResult.setResult(result);
}

private void setChangelog(KeyValue before, KeyValue after) {
    if (before == null || !before.isAdd()) {
        if (after.isAdd()) {
            reusedResult.addChangelog(INSERT);          // +I:该 key 第一次出现
        }
    } else {
        if (!after.isAdd()) {
            reusedResult.addChangelog(DELETE, before);  // -D:该 key 被删除
        } else if (!valueEqualiser.equals(before.value(), after.value())) {
            reusedResult.addChangelog(UPDATE_BEFORE, before);  // -U
            reusedResult.addChangelog(UPDATE_AFTER, after);    // +U
        }
    }
}

5.5 初始读取的特殊处理

流读启动时,lookup 模式只读 Level > 0 的数据(跳过 L0):

// DataTableStreamScan.java:160
} else if (options.changelogProducer().equals(LOOKUP)) {
    // L0 数据将在未来的 compaction 中产生 changelog
    result = startingScanner.scan(snapshotReader.withLevelFilter(level -> level > 0));
    snapshotReader.withLevelFilter(Filter.alwaysTrue());
}

这样保证:

  • 初始全量读只取稳定的高层数据(每个主键的最终状态)
  • L0 数据将在后续 Compaction 时作为 changelog 推送给下游,不会重复

5.6 LookupSinkWrite:故障恢复的连续性保证

lookup 模式使用 LookupSinkWrite,它在 Flink 状态中记录所有活跃 bucket:

// LookupSinkWrite.java

// 故障恢复时:对所有上次活跃的 bucket 触发一次轻量 compaction
List<StoreSinkWriteState.StateValue> activeBucketsStateValues = state.get(...);
if (activeBucketsStateValues != null) {
    for (StoreSinkWriteState.StateValue stateValue : activeBucketsStateValues) {
        write.compact(stateValue.partition(), stateValue.bucket(), false);
    }
}

这确保故障恢复后,L0 中可能遗留的数据会重新触发 Compaction,changelog 不会丢失。

5.7 与 Deletion Vectors 的关系

当同时开启 deletion-vectors.enabled = 'true'changelog-producer = 'lookup' 时,LookupStrategy 会设置 deletionVector = true,此时 Lookup 使用 PositionedKeyValueProcessor 获取记录在数据文件中的物理位置(行号),以便更新删除向量:

// CoreOptions.lookupStrategy():
LookupStrategy.from(
    mergeEngine().equals(MergeEngine.FIRST_ROW),
    changelogProducer().equals(ChangelogProducer.LOOKUP),  // produceChangelog
    deletionVectorsEnabled(),                               // deletionVector
    options.get(FORCE_LOOKUP)
)

5.8 配置示例

-- ✅ 推荐配置(需要低延迟精确 changelog)
CREATE TABLE my_table (
    id BIGINT,
    name STRING,
    update_time TIMESTAMP(3)
) TBLPROPERTIES (
    'primary-key' = 'id',
    'changelog-producer' = 'lookup'
    -- 可选:控制 Lookup 缓存
    -- 'lookup.cache-max-disk-size' = '10 gb',
    -- 'lookup.cache-file-retention' = '1 h'
);

-- ✅ 与 deletion-vectors 组合使用
CREATE TABLE my_table_dv (
    id BIGINT,
    name STRING
) TBLPROPERTIES (
    'primary-key' = 'id',
    'changelog-producer' = 'lookup',
    'deletion-vectors.enabled' = 'true'
);

5.9 特性总结

  • ✅ 低延迟:每次涉及 L0 的 Compaction 就生成 changelog,不需要等到 Full Compaction
  • ✅ 精确:通过 Lookup 查询实际旧值,UPDATE_BEFORE 值准确
  • ✅ 上游无需提供 RowKind,自动推断
  • ✅ 支持与 deletion-vectors.enabled 配合使用
  • ❌ 需要额外的本地磁盘空间存储 LookupFile(约等于 L1+ 数据量的一份副本)
  • ❌ Compaction 时有额外的 IO 开销(构建和查询 LookupFile)

六、四种模式的本质对比

6.1 旧值从哪里来

模式 旧值来源 获取时机
none
input 上游显式提供 写入时
full-compaction LSM maxLevel 的文件(上次 Full Compaction 的结果) Full Compaction 执行时
lookup Lookup 查询 L1+ 层的实时数据 任意涉及 L0 的 Compaction 时

6.2 changelog 文件的位置

三种生成 changelog 的模式(input / full-compaction / lookup)产生的 changelog 文件都存储在与数据文件相同的目录下,以 changelog- 为前缀,在 Manifest 中通过 changelogManifestList 字段索引。

下游统一通过 ChangelogFollowUpScanner 读取,只扫描含有 changelogManifestList 的快照。

6.3 选型建议

上游数据携带完整 RowKind(CDC 场景)
  └─ 使用 input 模式

需要实时精确 changelog,且可接受额外磁盘开销
  └─ 使用 lookup 模式

对延迟不敏感,希望最小化 Compaction 额外开销
  └─ 使用 full-compaction 模式(配合 delta-commits 控制频率)

不需要 changelog,或下游只消费最新状态
  └─ 使用 none 模式(默认)

6.4 关键配置参考

配置项 适用模式 说明
changelog-producer 全部 选择 producer 模式
full-compaction.delta-commits full-compaction 每 N 个 checkpoint 触发一次 Full Compaction
changelog-producer.row-deduplicate full-compaction / lookup 值未变化时是否跳过生成 -U/+U(默认 false)
lookup.cache-max-disk-size lookup LookupFile 本地磁盘缓存上限(默认 10GB)
lookup.cache-file-retention lookup LookupFile 缓存保留时间(默认 1h)
Flink RowKind 详解 2026-04-29

评论区