This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 7b27a4c25e [core] Fix lookup changelog for overwrite writes (#8287)
7b27a4c25e is described below
commit 7b27a4c25e696b355af8b86e452728c134169f70
Author: WenjunMin <[email protected]>
AuthorDate: Fri Jun 19 15:32:50 2026 +0800
[core] Fix lookup changelog for overwrite writes (#8287)
When a primary-key table uses `changelog-producer = 'lookup'`, overwrite
writes can still go through the lookup changelog path and produce extra
`changelog-*` files.
However, overwrite commits replace the affected data files and do not
reference those lookup changelog files from the latest table data. As a
result, the generated changelog files become orphan files. The lookup
process is unnecessary in this case and can be skipped when the write
already ignores previous files.
The added IT reproduces this on the pre-fix baseline: after one initial
changelog file, `INSERT OVERWRITE` increases the changelog file count
from 1 to 2.
---
.../compact/KvCompactionManagerFactory.java | 3 ++-
.../compact/MergeTreeCompactManagerFactory.java | 22 ++++++++++--------
.../ClusteringCompactManagerFactory.java | 3 ++-
.../paimon/operation/AbstractFileStoreWrite.java | 11 ++++++---
.../paimon/operation/BaseAppendFileStoreWrite.java | 3 ++-
.../paimon/operation/KeyValueFileStoreWrite.java | 11 +++++++--
.../postpone/PostponeBucketFileStoreWrite.java | 3 ++-
.../paimon/table/sink/BatchWriteBuilderImpl.java | 13 ++++++++---
.../org/apache/paimon/spark/SparkWriteITCase.java | 27 ++++++++++++++++++++++
9 files changed, 75 insertions(+), 21 deletions(-)
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 408788ab55..f512ad7b78 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
@@ -103,5 +103,6 @@ public interface KvCompactionManagerFactory extends
Closeable {
int bucket,
ExecutorService compactExecutor,
List<DataFileMeta> restoreFiles,
- @Nullable BucketedDvMaintainer dvMaintainer);
+ @Nullable BucketedDvMaintainer dvMaintainer,
+ boolean lookupEnabled);
}
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 89d1931384..8b858e6592 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
@@ -147,12 +147,14 @@ public class MergeTreeCompactManagerFactory implements
KvCompactionManagerFactor
int bucket,
ExecutorService compactExecutor,
List<DataFileMeta> restoreFiles,
- @Nullable BucketedDvMaintainer dvMaintainer) {
+ @Nullable BucketedDvMaintainer dvMaintainer,
+ boolean lookupEnabled) {
if (options.writeOnly()) {
return new NoopCompactManager();
}
- CompactStrategy compactStrategy = createCompactStrategy(options,
restoreFiles);
+ CompactStrategy compactStrategy =
+ createCompactStrategy(options, restoreFiles, lookupEnabled);
Comparator<InternalRow> keyComparator = keyComparatorSupplier.get();
Levels levels = new Levels(keyComparator, restoreFiles,
options.numLevels());
@Nullable FieldsComparator userDefinedSeqComparator =
udsComparatorSupplier.get();
@@ -163,7 +165,8 @@ public class MergeTreeCompactManagerFactory implements
KvCompactionManagerFactor
keyComparator,
userDefinedSeqComparator,
levels,
- dvMaintainer);
+ dvMaintainer,
+ lookupEnabled);
CompactionMetrics.Reporter metricsReporter =
compactionMetrics == null
? null
@@ -181,18 +184,18 @@ public class MergeTreeCompactManagerFactory implements
KvCompactionManagerFactor
rewriter,
metricsReporter,
dvMaintainer,
- options.prepareCommitWaitCompaction(),
- options.needLookup(),
+ lookupEnabled && options.prepareCommitWaitCompaction(),
+ lookupEnabled,
recordLevelExpire,
options.forceRewriteAllFiles(),
options.isChainTable());
}
private CompactStrategy createCompactStrategy(
- CoreOptions options, List<DataFileMeta> restoreFiles) {
+ CoreOptions options, List<DataFileMeta> restoreFiles, boolean
lookupEnabled) {
Long initialLastFullCompaction =
estimateLastFullCompactionTime(restoreFiles,
options.numLevels());
- if (options.needLookup()) {
+ if (lookupEnabled) {
Integer compactMaxInterval = null;
switch (options.lookupCompact()) {
case GENTLE:
@@ -247,7 +250,8 @@ public class MergeTreeCompactManagerFactory implements
KvCompactionManagerFactor
Comparator<InternalRow> keyComparator,
@Nullable FieldsComparator userDefinedSeqComparator,
Levels levels,
- @Nullable BucketedDvMaintainer dvMaintainer) {
+ @Nullable BucketedDvMaintainer dvMaintainer,
+ boolean lookupEnabled) {
DeletionVector.Factory dvFactory =
DeletionVector.factory(dvMaintainer);
KeyValueFileReaderFactory keyReaderFactory =
readerFactoryBuilder.build(partition, bucket, dvFactory);
@@ -273,7 +277,7 @@ public class MergeTreeCompactManagerFactory implements
KvCompactionManagerFactor
mfFactory,
mergeSorter,
logDedupEqualSupplier.get());
- } else if (lookupStrategy.needLookup) {
+ } else if (lookupEnabled && lookupStrategy.needLookup) {
PersistProcessor.Factory<?> processorFactory;
LookupMergeTreeCompactRewriter.MergeFunctionWrapperFactory<?>
wrapperFactory;
FileReaderFactory<KeyValue> lookupReaderFactory = readerFactory;
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 f0b117df8b..e58c172fd6 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
@@ -86,7 +86,8 @@ public class ClusteringCompactManagerFactory implements
KvCompactionManagerFacto
int bucket,
ExecutorService compactExecutor,
List<DataFileMeta> restoreFiles,
- @Nullable BucketedDvMaintainer dvMaintainer) {
+ @Nullable BucketedDvMaintainer dvMaintainer,
+ boolean lookupEnabled) {
if (options.writeOnly()) {
return new NoopCompactManager();
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreWrite.java
b/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreWrite.java
index eed5259517..953942d738 100644
---
a/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreWrite.java
+++
b/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreWrite.java
@@ -395,7 +395,10 @@ public abstract class AbstractFileStoreWrite<T> implements
FileStoreWrite<T> {
state.maxSequenceNumber,
state.commitIncrement,
compactExecutor(),
- state.deletionVectorsMaintainer);
+ state.deletionVectorsMaintainer,
+ // Restore reconstructs writer state from
checkpointed files, so do
+ // not ignore them.
+ false);
notifyNewWriter(writer);
WriterContainer<T> writerContainer =
new WriterContainer<>(
@@ -483,7 +486,8 @@ public abstract class AbstractFileStoreWrite<T> implements
FileStoreWrite<T> {
getMaxSequenceNumber(restoreFiles),
latestSnapshot),
null,
compactExecutor(),
- dvMaintainer);
+ dvMaintainer,
+ actualIgnorePreviousFiles);
notifyNewWriter(writer);
Snapshot previousSnapshot = restored.snapshot();
@@ -590,7 +594,8 @@ public abstract class AbstractFileStoreWrite<T> implements
FileStoreWrite<T> {
long restoredMaxSeqNumber,
@Nullable CommitIncrement restoreIncrement,
ExecutorService compactExecutor,
- @Nullable BucketedDvMaintainer deletionVectorsMaintainer);
+ @Nullable BucketedDvMaintainer deletionVectorsMaintainer,
+ boolean ignorePreviousFiles);
// force buffer spill to avoid out of memory in batch mode
protected void forceBufferSpill() throws Exception {}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java
b/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java
index 94620205a2..b40e8ee2bc 100644
---
a/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java
+++
b/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java
@@ -125,7 +125,8 @@ public abstract class BaseAppendFileStoreWrite extends
MemoryFileStoreWrite<Inte
long restoredMaxSeqNumber,
@Nullable CommitIncrement restoreIncrement,
ExecutorService compactExecutor,
- @Nullable BucketedDvMaintainer dvMaintainer) {
+ @Nullable BucketedDvMaintainer dvMaintainer,
+ boolean ignorePreviousFiles) {
return new AppendOnlyWriter(
fileIO,
ioManager,
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 51cfc29d1e..fc46c91c4b 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
@@ -202,7 +202,8 @@ public class KeyValueFileStoreWrite extends
MemoryFileStoreWrite<KeyValue> {
long restoredMaxSeqNumber,
@Nullable CommitIncrement restoreIncrement,
ExecutorService compactExecutor,
- @Nullable BucketedDvMaintainer dvMaintainer) {
+ @Nullable BucketedDvMaintainer dvMaintainer,
+ boolean ignorePreviousFiles) {
if (LOG.isDebugEnabled()) {
LOG.debug(
"Creating merge tree writer for partition {} bucket {}
from restored files {}",
@@ -211,12 +212,18 @@ public class KeyValueFileStoreWrite extends
MemoryFileStoreWrite<KeyValue> {
restoreFiles);
}
+ boolean lookupEnabled = !ignorePreviousFiles && options.needLookup();
KeyValueFileWriterFactory writerFactory =
writerFactoryBuilder.build(partition, bucket, options);
Comparator<InternalRow> keyComparator = keyComparatorSupplier.get();
CompactManager compactManager =
compactManagerFactory.create(
- partition, bucket, compactExecutor, restoreFiles,
dvMaintainer);
+ partition,
+ bucket,
+ compactExecutor,
+ restoreFiles,
+ dvMaintainer,
+ lookupEnabled);
return new MergeTreeWriter(
options.writeBufferSpillable(),
diff --git
a/paimon-core/src/main/java/org/apache/paimon/postpone/PostponeBucketFileStoreWrite.java
b/paimon-core/src/main/java/org/apache/paimon/postpone/PostponeBucketFileStoreWrite.java
index a6ba8febe5..9fa8afcd67 100644
---
a/paimon-core/src/main/java/org/apache/paimon/postpone/PostponeBucketFileStoreWrite.java
+++
b/paimon-core/src/main/java/org/apache/paimon/postpone/PostponeBucketFileStoreWrite.java
@@ -191,7 +191,8 @@ public class PostponeBucketFileStoreWrite extends
MemoryFileStoreWrite<KeyValue>
long restoredMaxSeqNumber,
@Nullable CommitIncrement restoreIncrement,
ExecutorService compactExecutor,
- @Nullable BucketedDvMaintainer deletionVectorsMaintainer) {
+ @Nullable BucketedDvMaintainer deletionVectorsMaintainer,
+ boolean ignorePreviousFiles) {
Preconditions.checkArgument(bucket == BucketMode.POSTPONE_BUCKET);
Preconditions.checkArgument(
restoreFiles.isEmpty(),
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 c06bd9de42..989db8e655 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,6 +39,7 @@ public class BatchWriteBuilderImpl implements
BatchWriteBuilder {
private final InnerTable table;
private final String commitUser;
+ private boolean overwrite;
private Map<String, String> staticPartition;
private boolean appendCommitCheckConflict = false;
private @Nullable Long rowIdCheckFromSnapshot = null;
@@ -49,9 +50,13 @@ public class BatchWriteBuilderImpl implements
BatchWriteBuilder {
}
private BatchWriteBuilderImpl(
- InnerTable table, String commitUser, @Nullable Map<String, String>
staticPartition) {
+ InnerTable table,
+ String commitUser,
+ boolean overwrite,
+ @Nullable Map<String, String> staticPartition) {
this.table = table;
this.commitUser = commitUser;
+ this.overwrite = overwrite;
this.staticPartition = staticPartition;
}
@@ -72,13 +77,14 @@ public class BatchWriteBuilderImpl implements
BatchWriteBuilder {
@Override
public BatchWriteBuilder withOverwrite(@Nullable Map<String, String>
staticPartition) {
+ this.overwrite = true;
this.staticPartition = staticPartition;
return this;
}
@Override
public BatchTableWrite newWrite() {
- return
table.newWrite(commitUser).withIgnorePreviousFiles(staticPartition != null);
+ return table.newWrite(commitUser).withIgnorePreviousFiles(overwrite);
}
@Override
@@ -96,7 +102,8 @@ public class BatchWriteBuilderImpl implements
BatchWriteBuilder {
}
public BatchWriteBuilderImpl copyWithNewTable(Table newTable) {
- return new BatchWriteBuilderImpl((InnerTable) newTable, commitUser,
staticPartition);
+ return new BatchWriteBuilderImpl(
+ (InnerTable) newTable, commitUser, overwrite, staticPartition);
}
public BatchWriteBuilderImpl appendCommitCheckConflict(boolean
appendCommitCheckConflict) {
diff --git
a/paimon-spark/paimon-spark-ut/src/test/java/org/apache/paimon/spark/SparkWriteITCase.java
b/paimon-spark/paimon-spark-ut/src/test/java/org/apache/paimon/spark/SparkWriteITCase.java
index c47ea797e7..d040116310 100644
---
a/paimon-spark/paimon-spark-ut/src/test/java/org/apache/paimon/spark/SparkWriteITCase.java
+++
b/paimon-spark/paimon-spark-ut/src/test/java/org/apache/paimon/spark/SparkWriteITCase.java
@@ -22,6 +22,7 @@ import org.apache.paimon.fs.FileIO;
import org.apache.paimon.fs.FileStatus;
import org.apache.paimon.fs.Path;
import org.apache.paimon.fs.local.LocalFileIO;
+import org.apache.paimon.io.DataFileMeta;
import org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions;
import org.apache.paimon.table.FileStoreTable;
import org.apache.paimon.table.FileStoreTableFactory;
@@ -642,6 +643,32 @@ public class SparkWriteITCase {
spark.conf().unset("spark.paimon.changelog-file.prefix");
}
+ @Test
+ public void testInsertOverwriteDoesNotProduceLookupChangelogFiles() throws
Exception {
+ spark.sql(
+ "CREATE TABLE T (a INT, b INT, c STRING) TBLPROPERTIES
('primary-key'='a', 'bucket' = '1', 'changelog-producer' = 'lookup',
'compaction.force-up-level-0' = 'true', 'file.format' = 'avro',
'file.compression' = 'null', 'manifest.compression' = 'null')");
+
+ FileStoreTable table = getTable("T");
+ Path bucketPath = new Path(table.location(), "bucket-0");
+
+ spark.sql("INSERT INTO T VALUES (1, 1, 'aa')");
+ long changelogFiles =
dataFileCount(table.fileIO().listStatus(bucketPath), "changelog-");
+ Assertions.assertEquals(1, changelogFiles);
+
+ spark.sql("INSERT OVERWRITE T VALUES (2, 2, 'bb')");
+
+ Assertions.assertEquals(
+ changelogFiles,
dataFileCount(table.fileIO().listStatus(bucketPath), "changelog-"));
+
Assertions.assertNull(table.latestSnapshot().get().changelogManifestList());
+ List<DataFileMeta> dataFiles =
+ table.newSnapshotReader().read().dataSplits().stream()
+ .flatMap(split -> split.dataFiles().stream())
+ .collect(Collectors.toList());
+ assertThat(dataFiles).isNotEmpty();
+ assertThat(dataFiles).allMatch(file -> file.level() > 0);
+ assertThat(spark.sql("SELECT * FROM
T").collectAsList().toString()).isEqualTo("[[2,2,bb]]");
+ }
+
@Test
public void testMarkDone() throws IOException {
spark.sql(