This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 684c2e5ecf [core] Support changelog-producer.ignore-update-before and 
changelog-producer.ignore-delete options (#9531)
684c2e5ecf is described below

commit 684c2e5ecfb309c26a3d16b5007e62eecd64ca99
Author: junmuz <[email protected]>
AuthorDate: Wed Sep 2 07:39:17 2026 +0100

    [core] Support changelog-producer.ignore-update-before and 
changelog-producer.ignore-delete options (#9531)
---
 docs/docs/primary-key-table/changelog-producer.md  |  7 +-
 docs/generated/core_configuration.html             | 12 ++++
 .../main/java/org/apache/paimon/CoreOptions.java   | 26 +++++++
 .../compact/ChangelogMergeTreeRewriter.java        | 14 ++++
 .../FullChangelogMergeTreeCompactRewriter.java     |  6 +-
 .../compact/LookupMergeTreeCompactRewriter.java    |  4 ++
 .../compact/MergeTreeCompactManagerFactory.java    |  8 ++-
 .../mergetree/ChangelogMergeTreeRewriterTest.java  |  2 +
 .../LookupMergeTreeCompactRewriterTest.java        |  2 +
 .../flink/FullCompactionFileStoreITCase.java       | 84 ++++++++++++++++++++++
 10 files changed, 161 insertions(+), 4 deletions(-)

diff --git a/docs/docs/primary-key-table/changelog-producer.md 
b/docs/docs/primary-key-table/changelog-producer.md
index f86eefdd35..11260cbf5d 100644
--- a/docs/docs/primary-key-table/changelog-producer.md
+++ b/docs/docs/primary-key-table/changelog-producer.md
@@ -99,7 +99,9 @@ Lookup will cache data on the memory and local disk, you can 
use the following o
 </table>
 
 Lookup changelog-producer supports `changelog-producer.row-deduplicate` to 
avoid generating -U, +U
-changelog for the same record.
+changelog for the same record. It also supports 
`changelog-producer.ignore-update-before` to exclude UPDATE_BEFORE (-U)
+records and `changelog-producer.ignore-delete` to exclude DELETE (-D) records 
from changelog files. These options are
+useful when downstream consumers only need the latest state (e.g. upsert 
sinks) and do not require retraction.
 
 (Note: Please increase `'execution.checkpointing.max-concurrent-checkpoints'` 
Flink configuration, this is very
 important for performance).
@@ -128,7 +130,8 @@ efficient as the input changelog producer and the latency 
to produce changelog m
 :::
 
 Full-compaction changelog-producer supports 
`changelog-producer.row-deduplicate` to avoid generating -U, +U
-changelog for the same record.
+changelog for the same record. It also supports 
`changelog-producer.ignore-update-before` and
+`changelog-producer.ignore-delete` to filter out -U and -D records from 
changelog files respectively.
 
 ## Changelog Merging
 
diff --git a/docs/generated/core_configuration.html 
b/docs/generated/core_configuration.html
index d25ad4a0c0..31a7f0e2ed 100644
--- a/docs/generated/core_configuration.html
+++ b/docs/generated/core_configuration.html
@@ -212,6 +212,18 @@ under the License.
             <td><p>Enum</p></td>
             <td>Whether to double write to a changelog file. This changelog 
file keeps the details of data changes, it can be read directly during stream 
reads. This can be applied to tables with primary keys. <br /><br />Possible 
values:<ul><li>"none": No changelog file.</li><li>"input": Double write to a 
changelog file when flushing memory table, the changelog is from 
input.</li><li>"full-compaction": Generate changelog files with each full 
compaction.</li><li>"lookup": Generate change [...]
         </tr>
+        <tr>
+            <td><h5>changelog-producer.ignore-delete</h5></td>
+            <td style="word-wrap: break-word;">false</td>
+            <td>Boolean</td>
+            <td>Whether to ignore delete records in the changelog. When set to 
true, DELETE (-D) records will not be written to changelog files. This 
configuration is only valid for the changelog-producer is lookup or 
full-compaction.</td>
+        </tr>
+        <tr>
+            <td><h5>changelog-producer.ignore-update-before</h5></td>
+            <td style="word-wrap: break-word;">false</td>
+            <td>Boolean</td>
+            <td>Whether to ignore update-before records in the changelog. When 
set to true, UPDATE_BEFORE (-U) records will not be written to changelog files. 
This configuration is only valid for the changelog-producer is lookup or 
full-compaction.</td>
+        </tr>
         <tr>
             <td><h5>changelog-producer.row-deduplicate</h5></td>
             <td style="word-wrap: break-word;">false</td>
diff --git a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java 
b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
index c9c2c32a93..114caa4553 100644
--- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
+++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
@@ -1084,6 +1084,24 @@ public class CoreOptions implements Serializable {
                                     + "This changelog file keeps the details 
of data changes, "
                                     + "it can be read directly during stream 
reads. This can be applied to tables with primary keys. ");
 
+    public static final ConfigOption<Boolean> 
CHANGELOG_PRODUCER_IGNORE_UPDATE_BEFORE =
+            key("changelog-producer.ignore-update-before")
+                    .booleanType()
+                    .defaultValue(false)
+                    .withDescription(
+                            "Whether to ignore update-before records in the 
changelog. "
+                                    + "When set to true, UPDATE_BEFORE (-U) 
records will not be written to changelog files. "
+                                    + "This configuration is only valid for 
the changelog-producer is lookup or full-compaction.");
+
+    public static final ConfigOption<Boolean> CHANGELOG_PRODUCER_IGNORE_DELETE 
=
+            key("changelog-producer.ignore-delete")
+                    .booleanType()
+                    .defaultValue(false)
+                    .withDescription(
+                            "Whether to ignore delete records in the 
changelog. "
+                                    + "When set to true, DELETE (-D) records 
will not be written to changelog files. "
+                                    + "This configuration is only valid for 
the changelog-producer is lookup or full-compaction.");
+
     public static final ConfigOption<Boolean> 
CHANGELOG_PRODUCER_ROW_DEDUPLICATE =
             key("changelog-producer.row-deduplicate")
                     .booleanType()
@@ -3833,6 +3851,14 @@ public class CoreOptions implements Serializable {
                 options.get(FORCE_LOOKUP));
     }
 
+    public boolean changelogProducerIgnoreUpdateBefore() {
+        return options.get(CHANGELOG_PRODUCER_IGNORE_UPDATE_BEFORE);
+    }
+
+    public boolean changelogProducerIgnoreDelete() {
+        return options.get(CHANGELOG_PRODUCER_IGNORE_DELETE);
+    }
+
     public boolean changelogRowDeduplicate() {
         return options.get(CHANGELOG_PRODUCER_ROW_DEDUPLICATE);
     }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/ChangelogMergeTreeRewriter.java
 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/ChangelogMergeTreeRewriter.java
index 00b9628036..a76c590d11 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/ChangelogMergeTreeRewriter.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/ChangelogMergeTreeRewriter.java
@@ -29,6 +29,7 @@ import org.apache.paimon.io.RollingFileWriter;
 import org.apache.paimon.manifest.FileSource;
 import org.apache.paimon.mergetree.MergeSorter;
 import org.apache.paimon.mergetree.SortedRun;
+import org.apache.paimon.types.RowKind;
 import org.apache.paimon.utils.CloseableIterator;
 import org.apache.paimon.utils.ExceptionUtils;
 import org.apache.paimon.utils.FieldsComparator;
@@ -49,6 +50,8 @@ public abstract class ChangelogMergeTreeRewriter extends 
MergeTreeCompactRewrite
     protected final int maxLevel;
     protected final MergeEngine mergeEngine;
     private final boolean produceChangelog;
+    private final boolean changelogIgnoreUpdateBefore;
+    private final boolean changelogIgnoreDelete;
     private final boolean forceDropDelete;
 
     public ChangelogMergeTreeRewriter(
@@ -61,6 +64,8 @@ public abstract class ChangelogMergeTreeRewriter extends 
MergeTreeCompactRewrite
             MergeFunctionFactory<KeyValue> mfFactory,
             MergeSorter mergeSorter,
             boolean produceChangelog,
+            boolean changelogIgnoreUpdateBefore,
+            boolean changelogIgnoreDelete,
             boolean forceDropDelete) {
         super(
                 readerFactory,
@@ -72,6 +77,8 @@ public abstract class ChangelogMergeTreeRewriter extends 
MergeTreeCompactRewrite
         this.maxLevel = maxLevel;
         this.mergeEngine = mergeEngine;
         this.produceChangelog = produceChangelog;
+        this.changelogIgnoreUpdateBefore = changelogIgnoreUpdateBefore;
+        this.changelogIgnoreDelete = changelogIgnoreDelete;
         this.forceDropDelete = forceDropDelete;
     }
 
@@ -150,6 +157,13 @@ public abstract class ChangelogMergeTreeRewriter extends 
MergeTreeCompactRewrite
                 }
                 if (produceChangelog) {
                     for (KeyValue kv : result.changelogs()) {
+                        if (changelogIgnoreUpdateBefore
+                                && kv.valueKind() == RowKind.UPDATE_BEFORE) {
+                            continue;
+                        }
+                        if (changelogIgnoreDelete && kv.valueKind() == 
RowKind.DELETE) {
+                            continue;
+                        }
                         changelogFileWriter.write(kv);
                     }
                 }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FullChangelogMergeTreeCompactRewriter.java
 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FullChangelogMergeTreeCompactRewriter.java
index 0304819e92..de9a4d9107 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FullChangelogMergeTreeCompactRewriter.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FullChangelogMergeTreeCompactRewriter.java
@@ -53,7 +53,9 @@ public class FullChangelogMergeTreeCompactRewriter extends 
ChangelogMergeTreeRew
             @Nullable FieldsComparator userDefinedSeqComparator,
             MergeFunctionFactory<KeyValue> mfFactory,
             MergeSorter mergeSorter,
-            @Nullable RecordEqualiser valueEqualiser) {
+            @Nullable RecordEqualiser valueEqualiser,
+            boolean changelogIgnoreUpdateBefore,
+            boolean changelogIgnoreDelete) {
         super(
                 maxLevel,
                 mergeEngine,
@@ -64,6 +66,8 @@ public class FullChangelogMergeTreeCompactRewriter extends 
ChangelogMergeTreeRew
                 mfFactory,
                 mergeSorter,
                 true,
+                changelogIgnoreUpdateBefore,
+                changelogIgnoreDelete,
                 false);
         this.valueEqualiser = valueEqualiser;
     }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriter.java
 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriter.java
index 6976d214b3..36cc4cd6a4 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriter.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriter.java
@@ -76,6 +76,8 @@ public class LookupMergeTreeCompactRewriter<T> extends 
ChangelogMergeTreeRewrite
             MergeSorter mergeSorter,
             MergeFunctionWrapperFactory<T> wrapperFactory,
             boolean produceChangelog,
+            boolean changelogIgnoreUpdateBefore,
+            boolean changelogIgnoreDelete,
             @Nullable BucketedDvMaintainer dvMaintainer,
             CoreOptions options,
             @Nullable RemoteLookupFileManager<T> remoteLookupFileManager) {
@@ -89,6 +91,8 @@ public class LookupMergeTreeCompactRewriter<T> extends 
ChangelogMergeTreeRewrite
                 mfFactory,
                 mergeSorter,
                 produceChangelog,
+                changelogIgnoreUpdateBefore,
+                changelogIgnoreDelete,
                 dvMaintainer != null);
         this.dvMaintainer = dvMaintainer;
         this.lookupLevels = lookupLevels;
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java
 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java
index 56e5f516b4..0ef421b0cb 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java
@@ -272,6 +272,8 @@ public class MergeTreeCompactManagerFactory implements 
KvCompactionManagerFactor
         MergeEngine mergeEngine = options.mergeEngine();
         ChangelogProducer changelogProducer = options.changelogProducer();
         LookupStrategy lookupStrategy = options.lookupStrategy();
+        boolean changelogIgnoreUpdateBefore = 
options.changelogProducerIgnoreUpdateBefore();
+        boolean changelogIgnoreDelete = 
options.changelogProducerIgnoreDelete();
         if (changelogProducer.equals(FULL_COMPACTION)) {
             return new FullChangelogMergeTreeCompactRewriter(
                     maxLevel,
@@ -282,7 +284,9 @@ public class MergeTreeCompactManagerFactory implements 
KvCompactionManagerFactor
                     userDefinedSeqComparator,
                     mfFactory,
                     mergeSorter,
-                    logDedupEqualSupplier.get());
+                    logDedupEqualSupplier.get(),
+                    changelogIgnoreUpdateBefore,
+                    changelogIgnoreDelete);
         } else if (lookupStrategy.needLookup) {
             PersistProcessor.Factory<?> processorFactory;
             LookupMergeTreeCompactRewriter.MergeFunctionWrapperFactory<?> 
wrapperFactory;
@@ -350,6 +354,8 @@ public class MergeTreeCompactManagerFactory implements 
KvCompactionManagerFactor
                     mergeSorter,
                     wrapperFactory,
                     lookupStrategy.produceChangelog && !ignorePreviousFiles,
+                    changelogIgnoreUpdateBefore,
+                    changelogIgnoreDelete,
                     dvMaintainer,
                     options,
                     remoteLookupFileManager);
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/mergetree/ChangelogMergeTreeRewriterTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/mergetree/ChangelogMergeTreeRewriterTest.java
index 02ea7309a7..aebb2b1040 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/mergetree/ChangelogMergeTreeRewriterTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/mergetree/ChangelogMergeTreeRewriterTest.java
@@ -314,6 +314,8 @@ public class ChangelogMergeTreeRewriterTest {
                     DeduplicateMergeFunction.factory(),
                     mergeSorter,
                     true,
+                    false,
+                    false,
                     true);
             this.rewriteChangelog = rewriteChangelog;
             this.closeWithException = closeWithException;
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriterTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriterTest.java
index 5c85b37df3..2ef97f54f0 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriterTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriterTest.java
@@ -63,6 +63,8 @@ class LookupMergeTreeCompactRewriterTest {
                         mock(MergeSorter.class),
                         
mock(LookupMergeTreeCompactRewriter.MergeFunctionWrapperFactory.class),
                         false,
+                        false,
+                        false,
                         null,
                         options,
                         null);
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FullCompactionFileStoreITCase.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FullCompactionFileStoreITCase.java
index a40587782b..32cebfd3c6 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FullCompactionFileStoreITCase.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FullCompactionFileStoreITCase.java
@@ -165,6 +165,90 @@ public class FullCompactionFileStoreITCase extends 
CatalogITCaseBase {
                 .containsExactlyInAnyOrder(Row.ofKind(RowKind.INSERT, "+I", 
"1", "4", "5"));
     }
 
+    @Test
+    public void testIgnoreUpdateBeforeAuditLog() throws Exception {
+        sql("ALTER TABLE %s SET ('changelog-producer.ignore-update-before' = 
'true')", table);
+
+        BlockingIterator<Row, Row> iterator =
+                BlockingIterator.of(streamSqlIter("SELECT * FROM 
%s$audit_log", table));
+
+        sql("INSERT INTO %s VALUES ('1', '2', '3')", table);
+        assertThat(iterator.collect(1))
+                .containsExactlyInAnyOrder(Row.ofKind(RowKind.INSERT, "+I", 
"1", "2", "3"));
+
+        // update: should only produce +U, no -U
+        sql("INSERT INTO %s VALUES ('1', '4', '5')", table);
+        assertThat(iterator.collect(1))
+                .containsExactlyInAnyOrder(Row.ofKind(RowKind.INSERT, "+U", 
"1", "4", "5"));
+
+        // delete: should still produce -D
+        sql("DELETE FROM %s WHERE a = '1'", table);
+        assertThat(iterator.collect(1))
+                .containsExactlyInAnyOrder(Row.ofKind(RowKind.INSERT, "-D", 
"1", "4", "5"));
+
+        iterator.close();
+    }
+
+    @Test
+    public void testIgnoreDeleteAuditLog() throws Exception {
+        sql("ALTER TABLE %s SET ('changelog-producer.ignore-delete' = 
'true')", table);
+
+        BlockingIterator<Row, Row> iterator =
+                BlockingIterator.of(streamSqlIter("SELECT * FROM 
%s$audit_log", table));
+
+        sql("INSERT INTO %s VALUES ('1', '2', '3')", table);
+        assertThat(iterator.collect(1))
+                .containsExactlyInAnyOrder(Row.ofKind(RowKind.INSERT, "+I", 
"1", "2", "3"));
+
+        // update: should produce both -U and +U
+        sql("INSERT INTO %s VALUES ('1', '4', '5')", table);
+        assertThat(iterator.collect(2))
+                .containsExactlyInAnyOrder(
+                        Row.ofKind(RowKind.INSERT, "-U", "1", "2", "3"),
+                        Row.ofKind(RowKind.INSERT, "+U", "1", "4", "5"));
+
+        // delete: should not produce -D
+        sql("DELETE FROM %s WHERE a = '1'", table);
+        // insert a new record to trigger compaction and verify no -D appeared
+        sql("INSERT INTO %s VALUES ('2', '3', '4')", table);
+        assertThat(iterator.collect(1))
+                .containsExactlyInAnyOrder(Row.ofKind(RowKind.INSERT, "+I", 
"2", "3", "4"));
+
+        iterator.close();
+    }
+
+    @Test
+    public void testIgnoreUpdateBeforeAndDeleteAuditLog() throws Exception {
+        sql(
+                "ALTER TABLE %s SET ("
+                        + "'changelog-producer.ignore-update-before' = 'true', 
"
+                        + "'changelog-producer.ignore-delete' = 'true')",
+                table);
+
+        BlockingIterator<Row, Row> iterator =
+                BlockingIterator.of(streamSqlIter("SELECT * FROM 
%s$audit_log", table));
+
+        sql("INSERT INTO %s VALUES ('1', '2', '3'), ('2', '5', '6')", table);
+        assertThat(iterator.collect(2))
+                .containsExactlyInAnyOrder(
+                        Row.ofKind(RowKind.INSERT, "+I", "1", "2", "3"),
+                        Row.ofKind(RowKind.INSERT, "+I", "2", "5", "6"));
+
+        // update: should only produce +U, no -U
+        sql("INSERT INTO %s VALUES ('1', '4', '5')", table);
+        assertThat(iterator.collect(1))
+                .containsExactlyInAnyOrder(Row.ofKind(RowKind.INSERT, "+U", 
"1", "4", "5"));
+
+        // delete key '2': should not produce -D
+        // update key '1' again to trigger compaction and verify no -D appeared
+        sql("DELETE FROM %s WHERE a = '2'", table);
+        sql("INSERT INTO %s VALUES ('1', '7', '8')", table);
+        assertThat(iterator.collect(1))
+                .containsExactlyInAnyOrder(Row.ofKind(RowKind.INSERT, "+U", 
"1", "7", "8"));
+
+        iterator.close();
+    }
+
     @Test
     public void testRowDeduplicateWithArrayRow() throws Exception {
         String table = "T_ARRAY_ROW";

Reply via email to