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"; }