diff --git a/docs/docs/program-api/java-writing.md b/docs/docs/program-api/java-writing.md
index 3d29235d7359..69b83734cf7f 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 000000000000..b1a9081de346
--- /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.
+ *
+ *
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.
+ *
+ *
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 8ff0180bfb9b..c6083ea56e4f 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 @@ static KvCompactionManagerFactory create(
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 0ef421b0cbbf..87ed4868accb 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.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.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 lookupFileCache;
+ @Nullable private CompactRewriterFactory compactRewriterFactory;
+ private boolean initialized;
public MergeTreeCompactManagerFactory(
KeyValueFileReaderFactory.Builder readerFactoryBuilder,
@@ -141,6 +146,14 @@ public void withCompactionMetrics(@Nullable CompactionMetrics compactionMetrics)
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 CompactManager create(
List restoreFiles,
@Nullable BucketedDvMaintainer dvMaintainer,
boolean ignorePreviousFiles) {
+ initialized = true;
if (options.writeOnly()) {
return new NoopCompactManager();
}
@@ -157,7 +171,7 @@ public CompactManager create(
Comparator 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 CompactManager create(
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 @@ private static Long estimateLastFullCompactionTime(
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 9f85047fb685..6945120d8ae5 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.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 @@ default void withWriteType(RowType writeType) {
void withCompactExecutor(ExecutorService compactExecutor);
+ /** Installs a compaction rewriter factory before any bucket writer is created. */
+ default FileStoreWrite 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 c67285368434..a81516b8e7b9 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.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 @@ protected boolean ignorePreviousFilesForWriter(
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 536ebf856a4b..50564d19cfe1 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.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 9c815db10d27..cb26d471106e 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.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 TableWrite withBlobConsumer(BlobConsumer blobConsumer) {
return this;
}
+ @Override
+ public TableWriteImpl withCompactRewriterFactory(CompactRewriterFactory factory) {
+ write.withCompactRewriterFactory(factory);
+ return this;
+ }
+
public TableWriteImpl 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 5bcff3b17184..cc970e5a25ec 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.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 000000000000..3fbde9d25f63
--- /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 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 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 rows = new ArrayList<>();
+ try (RecordReaderIterator 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 extraOptions, boolean primaryKey)
+ throws Exception {
+ Map 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> 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();
+ }
+ }
+}
diff --git a/paimon-test-utils/src/main/java/org/apache/paimon/testutils/junit/DockerImageVersions.java b/paimon-test-utils/src/main/java/org/apache/paimon/testutils/junit/DockerImageVersions.java
index 56be88d2cf62..b405686f3c77 100644
--- a/paimon-test-utils/src/main/java/org/apache/paimon/testutils/junit/DockerImageVersions.java
+++ b/paimon-test-utils/src/main/java/org/apache/paimon/testutils/junit/DockerImageVersions.java
@@ -24,5 +24,5 @@
*/
public class DockerImageVersions {
- public static final String MINIO = "minio/minio:RELEASE.2022-02-07T08-17-33Z";
+ public static final String MINIO = "quay.io/minio/minio:RELEASE.2022-02-07T08-17-33Z";
}