From 7450f8795332c5ebe744e378252080493529d832 Mon Sep 17 00:00:00 2001 From: Jordan Epstein Date: Sat, 12 Sep 2026 10:44:09 -0400 Subject: [PATCH 1/2] [core] Allow custom primary-key compaction rewriters Expose a per-bucket rewriter factory through the table write API so external engines can replace file rewriting while retaining Paimon compaction scheduling and commit coordination. Supply the default rewriter for selective fallback and define ownership, recovery, and initialization rules. Validate real rewrites and upgrades across reopened writers, producer and deletion-vector fallback selection, closure, write-only behavior, unsupported append writers, late configuration, and failed factory cleanup. All 24 focused Core tests pass on JDK 11 with the normal build checks. --- docs/docs/program-api/java-writing.md | 29 +++ .../compact/CompactRewriterFactory.java | 46 ++++ .../compact/KvCompactionManagerFactory.java | 4 + .../MergeTreeCompactManagerFactory.java | 45 +++- .../paimon/operation/FileStoreWrite.java | 6 + .../operation/KeyValueFileStoreWrite.java | 7 + .../apache/paimon/table/sink/TableWrite.java | 11 + .../paimon/table/sink/TableWriteImpl.java | 7 + .../MergeTreeCompactManagerFactoryTest.java | 45 ++++ .../sink/CompactRewriterFactoryTest.java | 230 ++++++++++++++++++ 10 files changed, 428 insertions(+), 2 deletions(-) create mode 100644 paimon-core/src/main/java/org/apache/paimon/mergetree/compact/CompactRewriterFactory.java create mode 100644 paimon-core/src/test/java/org/apache/paimon/table/sink/CompactRewriterFactoryTest.java 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(); + } + } +} From c15c0fcfee863d2163ed927cc588dbd40d143724 Mon Sep 17 00:00:00 2001 From: Jordan Epstein Date: Sat, 12 Sep 2026 11:11:14 -0400 Subject: [PATCH 2/2] [test] Pull MinIO test image from Quay --- .../org/apache/paimon/testutils/junit/DockerImageVersions.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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"; }