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 156ac1d211 [core] Allow custom primary-key compaction rewriters (#9762)
156ac1d211 is described below

commit 156ac1d211debf3f33ab73202b74f8fb5a78c075
Author: Jordan Epstein <[email protected]>
AuthorDate: Sun Sep 13 22:37:40 2026 -0400

    [core] Allow custom primary-key compaction rewriters (#9762)
---
 docs/docs/program-api/java-writing.md              |  29 +++
 .../mergetree/compact/CompactRewriterFactory.java  |  46 +++++
 .../compact/KvCompactionManagerFactory.java        |   4 +
 .../compact/MergeTreeCompactManagerFactory.java    |  45 +++-
 .../apache/paimon/operation/FileStoreWrite.java    |   6 +
 .../paimon/operation/KeyValueFileStoreWrite.java   |   7 +
 .../org/apache/paimon/table/sink/TableWrite.java   |  11 +
 .../apache/paimon/table/sink/TableWriteImpl.java   |   7 +
 .../MergeTreeCompactManagerFactoryTest.java        |  45 ++++
 .../table/sink/CompactRewriterFactoryTest.java     | 230 +++++++++++++++++++++
 10 files changed, 428 insertions(+), 2 deletions(-)

diff --git a/docs/docs/program-api/java-writing.md 
b/docs/docs/program-api/java-writing.md
index 3d29235d73..69b83734cf 100644
--- a/docs/docs/program-api/java-writing.md
+++ b/docs/docs/program-api/java-writing.md
@@ -180,3 +180,32 @@ selector API: they require dedicated bucket assignment and 
`write(row, bucket)`
 
 For a Flink job, use [FlinkSinkBuilder](flink-api#write-to-table) to integrate 
routing, checkpoints,
 and commits with the engine.
+
+## Custom Primary-Key Compaction Rewriters
+
+Applications can install a `CompactRewriterFactory` on a table writer to 
replace or wrap the
+file-rewrite work for each primary-key partition and bucket. Paimon continues 
to select compaction
+inputs, schedule work, and collect results for checkpoint commits. The factory 
receives the normal
+rewriter selected for the table's merge engine, changelog producer, and 
deletion-vector options,
+so an implementation can delegate unsupported operations to it.
+
+```java
+write.withCompactRewriterFactory((partition, bucket, defaultRewriter) -> {
+    // Return a custom CompactRewriter here, or retain Paimon's implementation.
+    return defaultRewriter;
+});
+```
+
+Configure the factory before writing, restoring, or compacting any bucket. 
Install it again on each
+recovered writer; the factory and rewriters are not checkpoint state. The 
callback receives an
+independent partition copy and is invoked for each newly opened or restored 
bucket writer.
+
+A custom rewriter implements `rewrite(outputLevel, dropDelete, sections)` and
+`upgrade(outputLevel, file)`, returning `CompactResult` file changes. It must 
preserve Paimon's
+merge, sequence, changelog, deletion-vector, record-expiration, and metadata 
contracts. The returned rewriter owns the
+default rewriter and must close it when closed, even if it handles every 
operation itself. Paimon
+closes the default rewriter if factory creation fails or returns null.
+
+This hook supports primary-key merge-tree writers. Append, postpone, and 
primary-key clustering
+writers reject it. With `write-only = true`, the factory is never invoked. 
Installing a factory
+after a bucket writer has been created is rejected.
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/CompactRewriterFactory.java
 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/CompactRewriterFactory.java
new file mode 100644
index 0000000000..b1a9081de3
--- /dev/null
+++ 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/CompactRewriterFactory.java
@@ -0,0 +1,46 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.mergetree.compact;
+
+import org.apache.paimon.data.BinaryRow;
+
+/**
+ * Creates a compaction rewriter for one primary-key partition and bucket.
+ *
+ * <p>The supplied rewriter is Paimon's implementation selected for the 
table's merge engine,
+ * changelog producer, and deletion-vector options. A factory may return it 
unchanged, or wrap it to
+ * delegate compactions that its implementation does not support. A 
replacement must preserve the
+ * same records, sequence numbers, changelogs, deletion vectors, record 
expiration, and file
+ * metadata contracts.
+ *
+ * <p>After a successful call, the returned rewriter owns the supplied 
rewriter and must close it
+ * when closed. If creation fails or returns null, Paimon closes the supplied 
rewriter. The factory
+ * must release any other resources it allocated before failing. Each call 
must return a rewriter
+ * owned exclusively by that bucket; it is closed by Paimon's compaction 
manager.
+ */
+@FunctionalInterface
+public interface CompactRewriterFactory {
+
+    /**
+     * Creates a rewriter before the bucket starts compacting. The partition 
is an independent copy
+     * that may be retained. Capture table schema, options, and file access in 
the factory as
+     * needed.
+     */
+    CompactRewriter create(BinaryRow partition, int bucket, CompactRewriter 
defaultRewriter);
+}
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/KvCompactionManagerFactory.java
 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/KvCompactionManagerFactory.java
index 8ff0180bfb..c6083ea56e 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/KvCompactionManagerFactory.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/KvCompactionManagerFactory.java
@@ -97,6 +97,10 @@ public interface KvCompactionManagerFactory extends 
Closeable {
 
     void withCompactionMetrics(@Nullable CompactionMetrics compactionMetrics);
 
+    default void withCompactRewriterFactory(CompactRewriterFactory factory) {
+        throw new UnsupportedOperationException("Custom compaction rewriters 
are not supported.");
+    }
+
     /** Create a {@link CompactManager} for the given partition and bucket. */
     CompactManager create(
             BinaryRow partition,
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 0ef421b0cb..87ed4868ac 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
@@ -22,6 +22,7 @@ import org.apache.paimon.CoreOptions;
 import org.apache.paimon.CoreOptions.ChangelogProducer;
 import org.apache.paimon.CoreOptions.MergeEngine;
 import org.apache.paimon.KeyValue;
+import org.apache.paimon.annotation.VisibleForTesting;
 import org.apache.paimon.codegen.RecordEqualiser;
 import org.apache.paimon.compact.CompactManager;
 import org.apache.paimon.compact.NoopCompactManager;
@@ -74,6 +75,8 @@ import static 
org.apache.paimon.CoreOptions.ChangelogProducer.FULL_COMPACTION;
 import static org.apache.paimon.CoreOptions.MergeEngine.DEDUPLICATE;
 import static 
org.apache.paimon.lookup.LookupStoreFactory.bloomFilterBuilderFactory;
 import static org.apache.paimon.mergetree.LookupFile.localFilePrefix;
+import static org.apache.paimon.utils.Preconditions.checkNotNull;
+import static org.apache.paimon.utils.Preconditions.checkState;
 
 /** Factory to create {@link MergeTreeCompactManager}. */
 public class MergeTreeCompactManagerFactory implements 
KvCompactionManagerFactory {
@@ -97,6 +100,8 @@ public class MergeTreeCompactManagerFactory implements 
KvCompactionManagerFactor
     @Nullable private IOManager ioManager;
     @Nullable private CompactionMetrics compactionMetrics;
     @Nullable private Cache<String, LookupFile> lookupFileCache;
+    @Nullable private CompactRewriterFactory compactRewriterFactory;
+    private boolean initialized;
 
     public MergeTreeCompactManagerFactory(
             KeyValueFileReaderFactory.Builder readerFactoryBuilder,
@@ -141,6 +146,14 @@ public class MergeTreeCompactManagerFactory implements 
KvCompactionManagerFactor
         this.compactionMetrics = compactionMetrics;
     }
 
+    @Override
+    public void withCompactRewriterFactory(CompactRewriterFactory factory) {
+        checkState(
+                !initialized,
+                "Configure the compaction rewriter factory before creating 
bucket writers.");
+        this.compactRewriterFactory = checkNotNull(factory);
+    }
+
     @Override
     public CompactManager create(
             BinaryRow partition,
@@ -149,6 +162,7 @@ public class MergeTreeCompactManagerFactory implements 
KvCompactionManagerFactor
             List<DataFileMeta> restoreFiles,
             @Nullable BucketedDvMaintainer dvMaintainer,
             boolean ignorePreviousFiles) {
+        initialized = true;
         if (options.writeOnly()) {
             return new NoopCompactManager();
         }
@@ -157,7 +171,7 @@ public class MergeTreeCompactManagerFactory implements 
KvCompactionManagerFactor
         Comparator<InternalRow> keyComparator = keyComparatorSupplier.get();
         Levels levels = new Levels(keyComparator, restoreFiles, 
options.numLevels());
         @Nullable FieldsComparator userDefinedSeqComparator = 
udsComparatorSupplier.get();
-        MergeTreeCompactRewriter rewriter =
+        MergeTreeCompactRewriter defaultRewriter =
                 createRewriter(
                         partition,
                         bucket,
@@ -166,12 +180,14 @@ public class MergeTreeCompactManagerFactory implements 
KvCompactionManagerFactor
                         levels,
                         dvMaintainer,
                         ignorePreviousFiles);
+        CompactRewriter rewriter =
+                wrapRewriter(compactRewriterFactory, partition, bucket, 
defaultRewriter);
         CompactionMetrics.Reporter metricsReporter =
                 compactionMetrics == null
                         ? null
                         : compactionMetrics.createReporter(partition, bucket);
         if (metricsReporter != null) {
-            rewriter.setMetricsReporter(metricsReporter);
+            defaultRewriter.setMetricsReporter(metricsReporter);
         }
         String bucketInfo = "bucket=" + bucket;
         if (partition.getFieldCount() > 0) {
@@ -250,6 +266,31 @@ public class MergeTreeCompactManagerFactory implements 
KvCompactionManagerFactor
         return max < 0 ? null : max;
     }
 
+    @VisibleForTesting
+    static CompactRewriter wrapRewriter(
+            @Nullable CompactRewriterFactory compactRewriterFactory,
+            BinaryRow partition,
+            int bucket,
+            CompactRewriter defaultRewriter) {
+        if (compactRewriterFactory == null) {
+            return defaultRewriter;
+        }
+        try {
+            return checkNotNull(
+                    compactRewriterFactory.create(partition.copy(), bucket, 
defaultRewriter),
+                    "The compaction rewriter factory must return a rewriter.");
+        } catch (RuntimeException | Error failure) {
+            try {
+                defaultRewriter.close();
+            } catch (Exception closeFailure) {
+                if (closeFailure != failure) {
+                    failure.addSuppressed(closeFailure);
+                }
+            }
+            throw failure;
+        }
+    }
+
     private MergeTreeCompactRewriter createRewriter(
             BinaryRow partition,
             int bucket,
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreWrite.java 
b/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreWrite.java
index 9f85047fb6..6945120d8a 100644
--- a/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreWrite.java
+++ b/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreWrite.java
@@ -27,6 +27,7 @@ import org.apache.paimon.index.DynamicBucketIndexMaintainer;
 import org.apache.paimon.index.pk.BucketedPrimaryKeyIndexMaintainer;
 import org.apache.paimon.io.DataFileMeta;
 import org.apache.paimon.memory.MemoryPoolFactory;
+import org.apache.paimon.mergetree.compact.CompactRewriterFactory;
 import org.apache.paimon.metrics.MetricRegistry;
 import org.apache.paimon.table.sink.CommitMessage;
 import org.apache.paimon.table.sink.SinkRecord;
@@ -86,6 +87,11 @@ public interface FileStoreWrite<T> extends 
Restorable<List<FileStoreWrite.State<
 
     void withCompactExecutor(ExecutorService compactExecutor);
 
+    /** Installs a compaction rewriter factory before any bucket writer is 
created. */
+    default FileStoreWrite<T> 
withCompactRewriterFactory(CompactRewriterFactory factory) {
+        throw new UnsupportedOperationException("Custom compaction rewriters 
are not supported.");
+    }
+
     /**
      * Write the data to the store according to the partition and bucket.
      *
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/operation/KeyValueFileStoreWrite.java
 
b/paimon-core/src/main/java/org/apache/paimon/operation/KeyValueFileStoreWrite.java
index c672853684..a81516b8e7 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/operation/KeyValueFileStoreWrite.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/operation/KeyValueFileStoreWrite.java
@@ -37,6 +37,7 @@ import org.apache.paimon.io.KeyValueFileReaderFactory;
 import org.apache.paimon.io.KeyValueFileWriterFactory;
 import org.apache.paimon.io.RecordLevelExpire;
 import org.apache.paimon.mergetree.MergeTreeWriter;
+import org.apache.paimon.mergetree.compact.CompactRewriterFactory;
 import org.apache.paimon.mergetree.compact.KvCompactionManagerFactory;
 import org.apache.paimon.mergetree.compact.LookupMergeFunction;
 import org.apache.paimon.mergetree.compact.MergeFunctionFactory;
@@ -180,6 +181,12 @@ public class KeyValueFileStoreWrite extends 
MemoryFileStoreWrite<KeyValue> {
         return ignorePreviousFiles;
     }
 
+    @Override
+    public KeyValueFileStoreWrite 
withCompactRewriterFactory(CompactRewriterFactory factory) {
+        compactManagerFactory.withCompactRewriterFactory(factory);
+        return this;
+    }
+
     @Override
     public KeyValueFileStoreWrite withIOManager(IOManager ioManager) {
         super.withIOManager(ioManager);
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/sink/TableWrite.java 
b/paimon-core/src/main/java/org/apache/paimon/table/sink/TableWrite.java
index 536ebf856a..50564d19cf 100644
--- a/paimon-core/src/main/java/org/apache/paimon/table/sink/TableWrite.java
+++ b/paimon-core/src/main/java/org/apache/paimon/table/sink/TableWrite.java
@@ -25,6 +25,7 @@ import org.apache.paimon.data.InternalRow;
 import org.apache.paimon.disk.IOManager;
 import org.apache.paimon.io.BundleRecords;
 import org.apache.paimon.memory.MemoryPoolFactory;
+import org.apache.paimon.mergetree.compact.CompactRewriterFactory;
 import org.apache.paimon.metrics.MetricRegistry;
 import org.apache.paimon.table.Table;
 import org.apache.paimon.types.RowType;
@@ -53,6 +54,16 @@ public interface TableWrite extends AutoCloseable {
      */
     TableWrite withBlobConsumer(BlobConsumer blobConsumer);
 
+    /**
+     * Installs a rewriter factory for primary-key merge-tree compaction. 
Configure this before
+     * writing, restoring, or compacting any bucket, and configure it again on 
each recovered
+     * writer. Paimon retains compaction scheduling and commit coordination. 
Append and clustering
+     * writers do not support this hook; write-only writers never invoke it.
+     */
+    default TableWrite withCompactRewriterFactory(CompactRewriterFactory 
factory) {
+        throw new UnsupportedOperationException("Custom compaction rewriters 
are not supported.");
+    }
+
     /** Calculate which partition {@code row} belongs to. */
     BinaryRow getPartition(InternalRow row);
 
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/sink/TableWriteImpl.java 
b/paimon-core/src/main/java/org/apache/paimon/table/sink/TableWriteImpl.java
index 9c815db10d..cb26d47110 100644
--- a/paimon-core/src/main/java/org/apache/paimon/table/sink/TableWriteImpl.java
+++ b/paimon-core/src/main/java/org/apache/paimon/table/sink/TableWriteImpl.java
@@ -27,6 +27,7 @@ import org.apache.paimon.disk.IOManager;
 import org.apache.paimon.io.BundleRecords;
 import org.apache.paimon.io.DataFileMeta;
 import org.apache.paimon.memory.MemoryPoolFactory;
+import org.apache.paimon.mergetree.compact.CompactRewriterFactory;
 import org.apache.paimon.metrics.MetricRegistry;
 import org.apache.paimon.operation.BundleFileStoreWriter;
 import org.apache.paimon.operation.FileStoreWrite;
@@ -135,6 +136,12 @@ public class TableWriteImpl<T> implements InnerTableWrite, 
Restorable<List<State
         return this;
     }
 
+    @Override
+    public TableWriteImpl<T> withCompactRewriterFactory(CompactRewriterFactory 
factory) {
+        write.withCompactRewriterFactory(factory);
+        return this;
+    }
+
     public TableWriteImpl<T> withCompactExecutor(ExecutorService 
compactExecutor) {
         write.withCompactExecutor(compactExecutor);
         return this;
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactoryTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactoryTest.java
index 5bcff3b171..cc970e5a25 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactoryTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactoryTest.java
@@ -36,15 +36,20 @@ import org.apache.paimon.types.DataTypes;
 import org.apache.paimon.types.RowType;
 
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
 
+import java.io.IOException;
 import java.util.Collections;
 import java.util.Comparator;
 import java.util.concurrent.ExecutorService;
 
 import static org.apache.paimon.CoreOptions.DELETION_VECTORS_ENABLED;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
 import static org.mockito.Answers.RETURNS_SELF;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyInt;
+import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
@@ -61,6 +66,46 @@ public class MergeTreeCompactManagerFactoryTest {
                     DataTypes.FIELD(0, "key", DataTypes.INT()),
                     DataTypes.FIELD(1, "value", DataTypes.INT()));
 
+    @ParameterizedTest
+    @ValueSource(booleans = {false, true})
+    public void testFailedFactoryClosesDefaultRewriter(boolean returnNull) 
throws Exception {
+        CompactRewriter delegate = mock(CompactRewriter.class);
+        CompactRewriterFactory factory =
+                (partition, bucket, rewriter) -> {
+                    if (returnNull) {
+                        return null;
+                    }
+                    throw new IllegalStateException("factory failed");
+                };
+        assertThatThrownBy(
+                        () ->
+                                MergeTreeCompactManagerFactory.wrapRewriter(
+                                        factory, BinaryRow.EMPTY_ROW, 0, 
delegate))
+                .isInstanceOf(returnNull ? NullPointerException.class : 
IllegalStateException.class)
+                .hasMessageContaining(returnNull ? "must return a rewriter" : 
"factory failed");
+        verify(delegate).close();
+    }
+
+    @Test
+    public void testCloseFailureDoesNotReplaceFactoryFailure() throws 
Exception {
+        CompactRewriter delegate = mock(CompactRewriter.class);
+        IOException closeFailure = new IOException("close failed");
+        doThrow(closeFailure).when(delegate).close();
+        IllegalStateException failure = new IllegalStateException("factory 
failed");
+        assertThatThrownBy(
+                        () ->
+                                MergeTreeCompactManagerFactory.wrapRewriter(
+                                        (partition, bucket, rewriter) -> {
+                                            throw failure;
+                                        },
+                                        BinaryRow.EMPTY_ROW,
+                                        0,
+                                        delegate))
+                .isSameAs(failure)
+                .hasSuppressedException(closeFailure);
+        verify(delegate).close();
+    }
+
     @Test
     public void testLookupValueProjection() throws Exception {
         Options options = new Options();
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/sink/CompactRewriterFactoryTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/table/sink/CompactRewriterFactoryTest.java
new file mode 100644
index 0000000000..3fbde9d25f
--- /dev/null
+++ 
b/paimon-core/src/test/java/org/apache/paimon/table/sink/CompactRewriterFactoryTest.java
@@ -0,0 +1,230 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.table.sink;
+
+import org.apache.paimon.compact.CompactResult;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.data.BinaryRowWriter;
+import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.disk.IOManagerImpl;
+import org.apache.paimon.fs.Path;
+import org.apache.paimon.fs.local.LocalFileIO;
+import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.mergetree.SortedRun;
+import org.apache.paimon.mergetree.compact.CompactRewriter;
+import 
org.apache.paimon.mergetree.compact.FullChangelogMergeTreeCompactRewriter;
+import org.apache.paimon.mergetree.compact.LookupMergeTreeCompactRewriter;
+import org.apache.paimon.mergetree.compact.MergeTreeCompactRewriter;
+import org.apache.paimon.reader.RecordReaderIterator;
+import org.apache.paimon.schema.FileSystemSchemaManager;
+import org.apache.paimon.schema.Schema;
+import org.apache.paimon.schema.SchemaUtils;
+import org.apache.paimon.schema.TableSchema;
+import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.FileStoreTableFactory;
+import org.apache.paimon.types.DataType;
+import org.apache.paimon.types.DataTypes;
+import org.apache.paimon.types.RowKind;
+import org.apache.paimon.types.RowType;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.CsvSource;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+/** Tests compaction rewriter installation through the table write API. */
+public class CompactRewriterFactoryTest {
+
+    @TempDir java.nio.file.Path tempDir;
+
+    @ParameterizedTest
+    @CsvSource({
+        "none,false",
+        "input,false",
+        "lookup,false",
+        "full-compaction,false",
+        "none,true",
+        "input,true",
+        "lookup,true"
+    })
+    public void testRewriteAndUpgradeAcrossReopenedWriters(String producer, 
boolean deletionVectors)
+            throws Exception {
+        Map<String, String> options = new HashMap<>();
+        options.put("changelog-producer", producer);
+        options.put("deletion-vectors.enabled", 
Boolean.toString(deletionVectors));
+        options.put("num-sorted-run.compaction-trigger", "100");
+        FileStoreTable table = createTable(options, true);
+        List<TrackingRewriter> rewriters = new ArrayList<>();
+        Class<?> expected =
+                producer.equals("full-compaction")
+                        ? FullChangelogMergeTreeCompactRewriter.class
+                        : producer.equals("lookup") || deletionVectors
+                                ? LookupMergeTreeCompactRewriter.class
+                                : MergeTreeCompactRewriter.class;
+
+        for (int run = 0; run < 2; run++) {
+            try (IOManagerImpl io = new IOManagerImpl(tempDir.toString());
+                    StreamTableWrite write = 
table.newWrite("test").withIOManager(io);
+                    StreamTableCommit commit = table.newCommit("test")) {
+                write.withCompactRewriterFactory(
+                        (partition, bucket, delegate) -> {
+                            assertThat(partition.getInt(0)).isBetween(1, 2);
+                            assertThat(bucket).isZero();
+                            assertThat(delegate).isExactlyInstanceOf(expected);
+                            TrackingRewriter rewriter = new 
TrackingRewriter(delegate);
+                            rewriters.add(rewriter);
+                            return rewriter;
+                        });
+                if (run == 0) {
+                    write.write(GenericRow.of(1, 1, 10));
+                    write.write(GenericRow.of(1, 2, 20));
+                    write.write(GenericRow.of(2, 1, 30));
+                    commit.commit(0, write.prepareCommit(true, 0));
+                } else {
+                    // Restore and compact existing buckets before receiving 
any new records.
+                    write.compact(partition(1), 0, true);
+                    write.compact(partition(2), 0, true);
+                    commit.commit(1, write.prepareCommit(true, 1));
+                    write.write(GenericRow.of(1, 1, 11));
+                    write.write(GenericRow.ofKind(RowKind.DELETE, 1, 2, 20));
+                    write.write(GenericRow.of(2, 1, 31));
+                    write.compact(partition(1), 0, true);
+                    write.compact(partition(2), 0, true);
+                    commit.commit(2, write.prepareCommit(true, 2));
+                }
+            }
+        }
+
+        assertThat(rewriters.size()).isGreaterThanOrEqualTo(4);
+        assertThat(rewriters.stream().mapToInt(r -> 
r.rewrites.get()).sum()).isPositive();
+        if (producer.equals("none") && !deletionVectors) {
+            assertThat(rewriters.stream().mapToInt(r -> 
r.upgrades.get()).sum()).isPositive();
+        }
+        rewriters.forEach(r -> assertThat(r.closes.get()).isEqualTo(1));
+        List<String> rows = new ArrayList<>();
+        try (RecordReaderIterator<InternalRow> reader =
+                new 
RecordReaderIterator<>(table.newRead().createReader(table.newScan().plan()))) {
+            while (reader.hasNext()) {
+                InternalRow row = reader.next();
+                rows.add(row.getInt(0) + "/" + row.getInt(1) + "/" + 
row.getInt(2));
+            }
+        }
+        assertThat(rows).containsExactlyInAnyOrder("1/1/11", "2/1/31");
+    }
+
+    @Test
+    public void 
testWriteOnlyDoesNotCreateRewritersAndLateInstallationIsRejected()
+            throws Exception {
+        FileStoreTable table = 
createTable(Collections.singletonMap("write-only", "true"), true);
+        try (StreamTableWrite write = table.newWrite("test");
+                StreamTableCommit commit = table.newCommit("test")) {
+            write.withCompactRewriterFactory(
+                    (partition, bucket, delegate) -> {
+                        throw new AssertionError("write-only must not create a 
rewriter");
+                    });
+            write.write(GenericRow.of(1, 1, 10));
+            commit.commit(0, write.prepareCommit(true, 0));
+            assertThatThrownBy(() -> write.withCompactRewriterFactory((p, b, 
delegate) -> delegate))
+                    .isInstanceOf(IllegalStateException.class)
+                    .hasMessageContaining("before creating bucket writers");
+        }
+    }
+
+    @Test
+    public void testAppendWriterRejectsFactory() throws Exception {
+        FileStoreTable table = createTable(Collections.emptyMap(), false);
+        try (StreamTableWrite write = table.newWrite("test")) {
+            assertThatThrownBy(() -> write.withCompactRewriterFactory((p, b, 
delegate) -> delegate))
+                    .isInstanceOf(UnsupportedOperationException.class);
+        }
+    }
+
+    private FileStoreTable createTable(Map<String, String> extraOptions, 
boolean primaryKey)
+            throws Exception {
+        Map<String, String> options = new HashMap<>(extraOptions);
+        options.put("bucket", primaryKey ? "1" : "-1");
+        RowType rowType =
+                RowType.of(
+                        new DataType[] {DataTypes.INT(), DataTypes.INT(), 
DataTypes.INT()},
+                        new String[] {"pt", "k", "v"});
+        Path path = new Path(tempDir.toUri());
+        TableSchema schema =
+                SchemaUtils.forceCommit(
+                        new FileSystemSchemaManager(LocalFileIO.create(), 
path),
+                        new Schema(
+                                rowType.getFields(),
+                                Collections.singletonList("pt"),
+                                primaryKey ? Arrays.asList("pt", "k") : 
Collections.emptyList(),
+                                options,
+                                ""));
+        return FileStoreTableFactory.create(LocalFileIO.create(), path, 
schema);
+    }
+
+    private static BinaryRow partition(int value) {
+        BinaryRow row = new BinaryRow(1);
+        BinaryRowWriter writer = new BinaryRowWriter(row);
+        writer.writeInt(0, value);
+        writer.complete();
+        return row;
+    }
+
+    private static class TrackingRewriter implements CompactRewriter {
+        private final CompactRewriter delegate;
+        private final AtomicInteger rewrites = new AtomicInteger();
+        private final AtomicInteger upgrades = new AtomicInteger();
+        private final AtomicInteger closes = new AtomicInteger();
+
+        private TrackingRewriter(CompactRewriter delegate) {
+            this.delegate = delegate;
+        }
+
+        @Override
+        public CompactResult rewrite(
+                int outputLevel, boolean dropDelete, List<List<SortedRun>> 
sections)
+                throws Exception {
+            rewrites.incrementAndGet();
+            return delegate.rewrite(outputLevel, dropDelete, sections);
+        }
+
+        @Override
+        public CompactResult upgrade(int outputLevel, DataFileMeta file) 
throws Exception {
+            upgrades.incrementAndGet();
+            return delegate.upgrade(outputLevel, file);
+        }
+
+        @Override
+        public void close() throws IOException {
+            closes.incrementAndGet();
+            delegate.close();
+        }
+    }
+}

Reply via email to