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 f512ad7b7816..8ff0180bfb9b 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 @@ -104,5 +104,5 @@ CompactManager create( ExecutorService compactExecutor, List restoreFiles, @Nullable BucketedDvMaintainer dvMaintainer, - boolean lookupEnabled); + boolean ignorePreviousFiles); } 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 8b858e6592ce..1daa51a399e1 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 @@ -148,13 +148,12 @@ public CompactManager create( ExecutorService compactExecutor, List restoreFiles, @Nullable BucketedDvMaintainer dvMaintainer, - boolean lookupEnabled) { + boolean ignorePreviousFiles) { if (options.writeOnly()) { return new NoopCompactManager(); } - CompactStrategy compactStrategy = - createCompactStrategy(options, restoreFiles, lookupEnabled); + CompactStrategy compactStrategy = createCompactStrategy(options, restoreFiles); Comparator keyComparator = keyComparatorSupplier.get(); Levels levels = new Levels(keyComparator, restoreFiles, options.numLevels()); @Nullable FieldsComparator userDefinedSeqComparator = udsComparatorSupplier.get(); @@ -166,7 +165,7 @@ public CompactManager create( userDefinedSeqComparator, levels, dvMaintainer, - lookupEnabled); + ignorePreviousFiles); CompactionMetrics.Reporter metricsReporter = compactionMetrics == null ? null @@ -184,18 +183,18 @@ public CompactManager create( rewriter, metricsReporter, dvMaintainer, - lookupEnabled && options.prepareCommitWaitCompaction(), - lookupEnabled, + options.prepareCommitWaitCompaction(), + options.needLookup(), recordLevelExpire, options.forceRewriteAllFiles(), options.isChainTable()); } private CompactStrategy createCompactStrategy( - CoreOptions options, List restoreFiles, boolean lookupEnabled) { + CoreOptions options, List restoreFiles) { Long initialLastFullCompaction = estimateLastFullCompactionTime(restoreFiles, options.numLevels()); - if (lookupEnabled) { + if (options.needLookup()) { Integer compactMaxInterval = null; switch (options.lookupCompact()) { case GENTLE: @@ -251,7 +250,7 @@ private MergeTreeCompactRewriter createRewriter( @Nullable FieldsComparator userDefinedSeqComparator, Levels levels, @Nullable BucketedDvMaintainer dvMaintainer, - boolean lookupEnabled) { + boolean ignorePreviousFiles) { DeletionVector.Factory dvFactory = DeletionVector.factory(dvMaintainer); KeyValueFileReaderFactory keyReaderFactory = readerFactoryBuilder.build(partition, bucket, dvFactory); @@ -277,7 +276,7 @@ private MergeTreeCompactRewriter createRewriter( mfFactory, mergeSorter, logDedupEqualSupplier.get()); - } else if (lookupEnabled && lookupStrategy.needLookup) { + } else if (lookupStrategy.needLookup) { PersistProcessor.Factory processorFactory; LookupMergeTreeCompactRewriter.MergeFunctionWrapperFactory wrapperFactory; FileReaderFactory lookupReaderFactory = readerFactory; @@ -335,7 +334,7 @@ private MergeTreeCompactRewriter createRewriter( mfFactory, mergeSorter, wrapperFactory, - lookupStrategy.produceChangelog, + lookupStrategy.produceChangelog && !ignorePreviousFiles, dvMaintainer, options, remoteLookupFileManager); diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/clustering/ClusteringCompactManagerFactory.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/clustering/ClusteringCompactManagerFactory.java index e58c172fd679..c63a9812fb4b 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/clustering/ClusteringCompactManagerFactory.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/clustering/ClusteringCompactManagerFactory.java @@ -87,7 +87,7 @@ public CompactManager create( ExecutorService compactExecutor, List restoreFiles, @Nullable BucketedDvMaintainer dvMaintainer, - boolean lookupEnabled) { + boolean ignorePreviousFiles) { if (options.writeOnly()) { return new NoopCompactManager(); } 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 fc46c91c4b76..1bbbfd2c306c 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 @@ -212,7 +212,6 @@ protected MergeTreeWriter createWriter( restoreFiles); } - boolean lookupEnabled = !ignorePreviousFiles && options.needLookup(); KeyValueFileWriterFactory writerFactory = writerFactoryBuilder.build(partition, bucket, options); Comparator keyComparator = keyComparatorSupplier.get(); @@ -223,7 +222,7 @@ protected MergeTreeWriter createWriter( compactExecutor, restoreFiles, dvMaintainer, - lookupEnabled); + ignorePreviousFiles); return new MergeTreeWriter( options.writeBufferSpillable(), diff --git a/paimon-core/src/main/java/org/apache/paimon/table/sink/BatchWriteBuilderImpl.java b/paimon-core/src/main/java/org/apache/paimon/table/sink/BatchWriteBuilderImpl.java index 989db8e65596..c06bd9de42d6 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/sink/BatchWriteBuilderImpl.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/sink/BatchWriteBuilderImpl.java @@ -39,7 +39,6 @@ public class BatchWriteBuilderImpl implements BatchWriteBuilder { private final InnerTable table; private final String commitUser; - private boolean overwrite; private Map staticPartition; private boolean appendCommitCheckConflict = false; private @Nullable Long rowIdCheckFromSnapshot = null; @@ -50,13 +49,9 @@ public BatchWriteBuilderImpl(InnerTable table) { } private BatchWriteBuilderImpl( - InnerTable table, - String commitUser, - boolean overwrite, - @Nullable Map staticPartition) { + InnerTable table, String commitUser, @Nullable Map staticPartition) { this.table = table; this.commitUser = commitUser; - this.overwrite = overwrite; this.staticPartition = staticPartition; } @@ -77,14 +72,13 @@ public Optional newWriteSelector() { @Override public BatchWriteBuilder withOverwrite(@Nullable Map staticPartition) { - this.overwrite = true; this.staticPartition = staticPartition; return this; } @Override public BatchTableWrite newWrite() { - return table.newWrite(commitUser).withIgnorePreviousFiles(overwrite); + return table.newWrite(commitUser).withIgnorePreviousFiles(staticPartition != null); } @Override @@ -102,8 +96,7 @@ public BatchTableCommit newCommit() { } public BatchWriteBuilderImpl copyWithNewTable(Table newTable) { - return new BatchWriteBuilderImpl( - (InnerTable) newTable, commitUser, overwrite, staticPartition); + return new BatchWriteBuilderImpl((InnerTable) newTable, commitUser, staticPartition); } public BatchWriteBuilderImpl appendCommitCheckConflict(boolean appendCommitCheckConflict) { diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/LookupRemoteFileTableTest.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/LookupRemoteFileTableTest.java index b78b5e98d970..95b2b84fbe30 100644 --- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/LookupRemoteFileTableTest.java +++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/LookupRemoteFileTableTest.java @@ -268,4 +268,31 @@ public void testRemoteFileLevelThreshold() throws Exception { assertThat(level5.extraFiles()).hasSize(1); assertThat(level4.extraFiles()).hasSize(0); } + + @Test + public void testOverwriteGeneratesRemoteFile() throws Exception { + Options options = new Options(); + options.set(CoreOptions.BUCKET, 1); + options.set(CoreOptions.DELETION_VECTORS_ENABLED, true); + options.set(CoreOptions.LOOKUP_REMOTE_FILE_ENABLED, true); + Identifier identifier = new Identifier("default", "t"); + Schema schema = + new Schema( + RowType.of(new IntType(), new IntType()).getFields(), + Collections.emptyList(), + Collections.singletonList("f0"), + options.toMap(), + null); + catalog.createTable(identifier, schema, false); + FileStoreTable table = (FileStoreTable) catalog.getTable(identifier); + BatchWriteBuilder writeBuilder = table.newBatchWriteBuilder().withOverwrite(); + try (BatchTableWrite write = writeBuilder.newWrite().withIOManager(ioManager); + BatchTableCommit commit = writeBuilder.newCommit()) { + write.write(GenericRow.of(1, 1)); + write.write(GenericRow.of(2, 1)); + commit.commit(write.prepareCommit()); + } + + allShouldHaveRemoteSst(table.newReadBuilder().newScan().plan().splits()); + } }