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 e435a774e2 [core] Support lightweight sequence initialization for 
write-only primary-key tables (#8256)
e435a774e2 is described below

commit e435a774e2d1e4c31cd1fb7f65b69a221ab6af50
Author: WenjunMin <[email protected]>
AuthorDate: Fri Jun 19 09:34:58 2026 +0800

    [core] Support lightweight sequence initialization for write-only 
primary-key tables (#8256)
    
    Primary-key table writers currently need to scan existing metadata to
    initialize the max sequence number. For write-only workloads this can be
    heavier than necessary because the writer only needs a safe starting
    sequence number.
---
 docs/generated/core_configuration.html             |   6 +
 .../main/java/org/apache/paimon/CoreOptions.java   |  40 ++++++
 .../paimon/operation/AbstractFileStoreWrite.java   |  44 +++++-
 .../paimon/operation/FileStoreCommitImpl.java      |  24 ++++
 .../paimon/operation/KeyValueFileStoreWrite.java   |  24 ++++
 .../operation/SequenceSnapshotProperties.java      |  73 ++++++++++
 .../operation/KeyValueFileStoreWriteTest.java      | 147 +++++++++++++++++++--
 .../flink/PrimaryKeyFileStoreTableITCase.java      |  30 +++++
 8 files changed, 372 insertions(+), 16 deletions(-)

diff --git a/docs/generated/core_configuration.html 
b/docs/generated/core_configuration.html
index c17b6a5007..575c40e21f 100644
--- a/docs/generated/core_configuration.html
+++ b/docs/generated/core_configuration.html
@@ -1722,6 +1722,12 @@ If the data size allocated for the sorting task is 
uneven,which may lead to perf
             <td>Boolean</td>
             <td>If set to true, compactions and snapshot expiration will be 
skipped. This option is used along with dedicated compact jobs.</td>
         </tr>
+        <tr>
+            <td><h5>write.sequence-number-init-mode</h5></td>
+            <td style="word-wrap: break-word;">scan</td>
+            <td><p>Enum</p></td>
+            <td>Specify how to initialize the next sequence number for primary 
key table writers.<br /><br />Possible values:<ul><li>"scan": initialize by 
scanning existing file metadata.</li><li>"snapshot": initialize from the 
maximum sequence number recorded in snapshot properties, which can avoid 
scanning existing file metadata in write-only mode.</li></ul></td>
+        </tr>
         <tr>
             <td><h5>write.batch-memory</h5></td>
             <td style="word-wrap: break-word;">128 mb</td>
diff --git a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java 
b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
index 81bf1fe0e0..f2a0d32a4d 100644
--- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
+++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
@@ -838,6 +838,14 @@ public class CoreOptions implements Serializable {
                     .defaultValue(false)
                     .withDescription("Whether to force a compaction before 
commit.");
 
+    public static final ConfigOption<SequenceNumberInitMode> 
WRITE_SEQUENCE_NUMBER_INIT_MODE =
+            key("write.sequence-number-init-mode")
+                    .enumType(SequenceNumberInitMode.class)
+                    .defaultValue(SequenceNumberInitMode.SCAN)
+                    .withDescription(
+                            "Specify how to initialize the next sequence 
number for primary key "
+                                    + "table writers.");
+
     public static final ConfigOption<Duration> COMMIT_TIMEOUT =
             key("commit.timeout")
                     .durationType()
@@ -3246,6 +3254,10 @@ public class CoreOptions implements Serializable {
         return options.get(COMMIT_FORCE_COMPACT);
     }
 
+    public SequenceNumberInitMode writeSequenceNumberInitMode() {
+        return options.get(WRITE_SEQUENCE_NUMBER_INIT_MODE);
+    }
+
     public long commitTimeout() {
         return options.get(COMMIT_TIMEOUT) == null
                 ? Long.MAX_VALUE
@@ -4232,6 +4244,34 @@ public class CoreOptions implements Serializable {
         }
     }
 
+    /** Specifies how to initialize the next sequence number for primary key 
table writers. */
+    public enum SequenceNumberInitMode implements DescribedEnum {
+        SCAN("scan", "initialize by scanning existing file metadata."),
+
+        SNAPSHOT(
+                "snapshot",
+                "initialize from the maximum sequence number recorded in 
snapshot properties, "
+                        + "which can avoid scanning existing file metadata in 
write-only mode.");
+
+        private final String value;
+        private final String description;
+
+        SequenceNumberInitMode(String value, String description) {
+            this.value = value;
+            this.description = description;
+        }
+
+        @Override
+        public String toString() {
+            return value;
+        }
+
+        @Override
+        public InlineElement getDescription() {
+            return text(description);
+        }
+    }
+
     /** Specifies the changelog producer for table. */
     public enum ChangelogProducer implements DescribedEnum {
         NONE("none", "No changelog file."),
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 a49fff06d0..eed5259517 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
@@ -58,6 +58,7 @@ import java.util.HashMap;
 import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
+import java.util.OptionalLong;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.function.Function;
@@ -96,9 +97,11 @@ public abstract class AbstractFileStoreWrite<T> implements 
FileStoreWrite<T> {
     private boolean ignorePreviousFiles = false;
     private boolean ignoreNumBucketCheck = false;
 
+    protected final SnapshotManager snapshotManager;
     protected CompactionMetrics compactionMetrics = null;
     protected final String tableName;
     private final boolean legacyPartitionName;
+    private final CoreOptions options;
 
     protected AbstractFileStoreWrite(
             SnapshotManager snapshotManager,
@@ -114,6 +117,7 @@ public abstract class AbstractFileStoreWrite<T> implements 
FileStoreWrite<T> {
         } else if (dvMaintainerFactory != null) {
             indexFileHandler = dvMaintainerFactory.indexFileHandler();
         }
+        this.snapshotManager = snapshotManager;
         this.restore = new FileSystemWriteRestore(options, snapshotManager, 
scan, indexFileHandler);
         this.dbMaintainerFactory = dbMaintainerFactory;
         this.dvMaintainerFactory = dvMaintainerFactory;
@@ -123,6 +127,7 @@ public abstract class AbstractFileStoreWrite<T> implements 
FileStoreWrite<T> {
         this.tableName = tableName;
         this.writerNumberMax = options.writeMaxWritersToSpill();
         this.legacyPartitionName = options.legacyPartitionName();
+        this.options = options;
         this.partitionTimestampValidator =
                 PartitionTimestampValidator.create(options, partitionType);
     }
@@ -445,8 +450,12 @@ public abstract class AbstractFileStoreWrite<T> implements 
FileStoreWrite<T> {
             }
         }
 
+        Snapshot latestSnapshot = snapshotManager.latestSnapshot();
+        boolean actualIgnorePreviousFiles =
+                ignorePreviousFilesForWriter(
+                        partition, bucket, latestSnapshot, 
ignorePreviousFiles);
         RestoreFiles restored = RestoreFiles.empty();
-        if (!ignorePreviousFiles) {
+        if (!actualIgnorePreviousFiles) {
             restored = scanExistingFileMetas(partition, bucket);
         }
 
@@ -470,7 +479,8 @@ public abstract class AbstractFileStoreWrite<T> implements 
FileStoreWrite<T> {
                         partition.copy(),
                         bucket,
                         restoreFiles,
-                        getMaxSequenceNumber(restoreFiles),
+                        startingMaxSequenceNumber(
+                                getMaxSequenceNumber(restoreFiles), 
latestSnapshot),
                         null,
                         compactExecutor(),
                         dvMaintainer);
@@ -489,6 +499,36 @@ public abstract class AbstractFileStoreWrite<T> implements 
FileStoreWrite<T> {
         return writers.values().stream().mapToLong(Map::size).sum();
     }
 
+    protected boolean ignorePreviousFilesForWriter(
+            BinaryRow partition,
+            int bucket,
+            @Nullable Snapshot latestSnapshot,
+            boolean ignorePreviousFiles) {
+        return ignorePreviousFiles;
+    }
+
+    @VisibleForTesting
+    long startingMaxSequenceNumber(long restoredMaxSeqNumber, @Nullable 
Snapshot latestSnapshot) {
+        if (options.writeSequenceNumberInitMode() != 
CoreOptions.SequenceNumberInitMode.SNAPSHOT) {
+            return restoredMaxSeqNumber;
+        }
+
+        OptionalLong snapshotMaxSeqNumber =
+                SequenceSnapshotProperties.maxSequenceNumber(latestSnapshot);
+        long startingMaxSeqNumber =
+                Math.max(restoredMaxSeqNumber, 
snapshotMaxSeqNumber.orElse(-1L));
+        LOG.info(
+                "Start writer sequence number for table {} from restored max 
sequence number {}, "
+                        + "snapshot max sequence number {}, selected max 
sequence number {}, "
+                        + "snapshot id {}.",
+                tableName,
+                restoredMaxSeqNumber,
+                snapshotMaxSeqNumber.isPresent() ? 
snapshotMaxSeqNumber.getAsLong() : null,
+                startingMaxSeqNumber,
+                latestSnapshot == null ? null : latestSnapshot.id());
+        return startingMaxSeqNumber;
+    }
+
     @Override
     public FileStoreWrite<T> withMetricRegistry(MetricRegistry metricRegistry) 
{
         this.compactionMetrics = new CompactionMetrics(metricRegistry, 
tableName);
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommitImpl.java
 
b/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommitImpl.java
index afb7ad6698..a316c25711 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommitImpl.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommitImpl.java
@@ -93,6 +93,7 @@ import java.util.List;
 import java.util.Map;
 import java.util.Objects;
 import java.util.Optional;
+import java.util.OptionalLong;
 import java.util.Set;
 import java.util.function.Supplier;
 import java.util.stream.Collectors;
@@ -785,6 +786,17 @@ public class FileStoreCommitImpl implements 
FileStoreCommit {
                 && !options.bucketAppendOrdered();
     }
 
+    private OptionalLong maxSequenceNumber(List<ManifestFileMeta> manifests) {
+        return manifests.stream()
+                .flatMap(
+                        manifest ->
+                                manifestFile.read(manifest.fileName(), 
manifest.fileSize())
+                                        .stream())
+                .filter(entry -> entry.kind() == FileKind.ADD)
+                .mapToLong(entry -> entry.file().maxSequenceNumber())
+                .max();
+    }
+
     /**
      * Try to overwrite partition.
      *
@@ -1038,6 +1050,18 @@ public class FileStoreCommitImpl implements 
FileStoreCommit {
                 }
             }
 
+            if (options.writeSequenceNumberInitMode()
+                    == CoreOptions.SequenceNumberInitMode.SNAPSHOT) {
+                OptionalLong latestMaxSequenceNumber =
+                        
SequenceSnapshotProperties.maxSequenceNumber(latestSnapshot);
+                if (!latestMaxSequenceNumber.isPresent() && latestSnapshot != 
null) {
+                    latestMaxSequenceNumber = 
maxSequenceNumber(mergeBeforeManifests);
+                }
+                properties =
+                        SequenceSnapshotProperties.mergeMaxSequenceNumber(
+                                properties, latestMaxSequenceNumber, 
deltaFiles);
+            }
+
             // prepare snapshot file
             newSnapshot =
                     new Snapshot(
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 41ecc6d0aa..51cfc29d1e 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
@@ -21,6 +21,7 @@ package org.apache.paimon.operation;
 import org.apache.paimon.CoreOptions;
 import org.apache.paimon.KeyValue;
 import org.apache.paimon.KeyValueFileStore;
+import org.apache.paimon.Snapshot;
 import org.apache.paimon.codegen.RecordEqualiser;
 import org.apache.paimon.compact.CompactManager;
 import org.apache.paimon.data.BinaryRow;
@@ -151,6 +152,29 @@ public class KeyValueFileStoreWrite extends 
MemoryFileStoreWrite<KeyValue> {
                         schema,
                         recordLevelExpire,
                         cacheManager);
+        if (options.writeOnly()
+                && options.writeSequenceNumberInitMode()
+                        == CoreOptions.SequenceNumberInitMode.SNAPSHOT) {
+            super.withIgnorePreviousFiles(true);
+            LOG.info(
+                    "Enable ignoring previous files for write-only snapshot 
sequence number initialization of table {}.",
+                    tableName);
+        }
+    }
+
+    @Override
+    protected boolean ignorePreviousFilesForWriter(
+            BinaryRow partition,
+            int bucket,
+            @Nullable Snapshot latestSnapshot,
+            boolean ignorePreviousFiles) {
+        if (options.writeOnly()
+                && options.writeSequenceNumberInitMode()
+                        == CoreOptions.SequenceNumberInitMode.SNAPSHOT) {
+            return latestSnapshot == null
+                    || 
SequenceSnapshotProperties.maxSequenceNumber(latestSnapshot).isPresent();
+        }
+        return ignorePreviousFiles;
     }
 
     @Override
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/operation/SequenceSnapshotProperties.java
 
b/paimon-core/src/main/java/org/apache/paimon/operation/SequenceSnapshotProperties.java
new file mode 100644
index 0000000000..a45c7f24fa
--- /dev/null
+++ 
b/paimon-core/src/main/java/org/apache/paimon/operation/SequenceSnapshotProperties.java
@@ -0,0 +1,73 @@
+/*
+ * 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.operation;
+
+import org.apache.paimon.Snapshot;
+import org.apache.paimon.manifest.FileKind;
+import org.apache.paimon.manifest.ManifestEntry;
+
+import javax.annotation.Nullable;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.OptionalLong;
+
+/** Helpers for persisting generated sequence state in snapshot properties. */
+class SequenceSnapshotProperties {
+
+    static final String MAX_SEQUENCE_NUMBER = 
"sequence.generation.max-sequence-number";
+
+    private SequenceSnapshotProperties() {}
+
+    static OptionalLong maxSequenceNumber(@Nullable Snapshot snapshot) {
+        if (snapshot == null || snapshot.properties() == null) {
+            return OptionalLong.empty();
+        }
+
+        String value = snapshot.properties().get(MAX_SEQUENCE_NUMBER);
+        if (value == null) {
+            return OptionalLong.empty();
+        }
+
+        return OptionalLong.of(Long.parseLong(value));
+    }
+
+    static Map<String, String> mergeMaxSequenceNumber(
+            @Nullable Map<String, String> properties,
+            OptionalLong latestMax,
+            List<ManifestEntry> deltaFiles) {
+        OptionalLong currentMax =
+                deltaFiles.stream()
+                        .filter(entry -> entry.kind() == FileKind.ADD)
+                        .mapToLong(entry -> entry.file().maxSequenceNumber())
+                        .max();
+
+        if (!latestMax.isPresent() && !currentMax.isPresent()) {
+            return properties == null ? new HashMap<>() : properties;
+        }
+
+        long maxSequenceNumber =
+                Math.max(latestMax.orElse(Long.MIN_VALUE), 
currentMax.orElse(Long.MIN_VALUE));
+        Map<String, String> merged =
+                properties == null ? new HashMap<>() : new 
HashMap<>(properties);
+        merged.put(MAX_SEQUENCE_NUMBER, String.valueOf(maxSequenceNumber));
+        return merged;
+    }
+}
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/operation/KeyValueFileStoreWriteTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/operation/KeyValueFileStoreWriteTest.java
index c13fcc1f64..81beed117e 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/operation/KeyValueFileStoreWriteTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/operation/KeyValueFileStoreWriteTest.java
@@ -20,6 +20,7 @@ package org.apache.paimon.operation;
 
 import org.apache.paimon.CoreOptions;
 import org.apache.paimon.KeyValue;
+import org.apache.paimon.Snapshot;
 import org.apache.paimon.TestFileStore;
 import org.apache.paimon.TestKeyValueGenerator;
 import org.apache.paimon.disk.IOManager;
@@ -32,6 +33,7 @@ import 
org.apache.paimon.mergetree.compact.DeduplicateMergeFunction;
 import org.apache.paimon.mergetree.compact.ForceUpLevel0Compaction;
 import org.apache.paimon.mergetree.compact.MergeTreeCompactManager;
 import org.apache.paimon.schema.Schema;
+import org.apache.paimon.schema.SchemaChange;
 import org.apache.paimon.schema.SchemaManager;
 import org.apache.paimon.schema.TableSchema;
 
@@ -40,6 +42,7 @@ import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.io.TempDir;
 
 import java.io.IOException;
+import java.util.Collections;
 import java.util.HashMap;
 import java.util.Map;
 
@@ -91,6 +94,117 @@ public class KeyValueFileStoreWriteTest {
         
assertThat(compactStrategy).isInstanceOf(ForceUpLevel0Compaction.class);
     }
 
+    @Test
+    public void testSnapshotSequenceNumberInitRestoresFromSnapshotProperties() 
throws Exception {
+        Map<String, String> options = new HashMap<>();
+        options.put(CoreOptions.WRITE_ONLY.key(), "true");
+        options.put(CoreOptions.WRITE_SEQUENCE_NUMBER_INIT_MODE.key(), 
"snapshot");
+        TestFileStore store = createStoreWithOptions(options);
+        TestKeyValueGenerator gen = new TestKeyValueGenerator();
+
+        KeyValue first = gen.nextInsert("20201110", 10, 1L, new int[] {1, 1}, 
"first");
+        store.commitData(Collections.singletonList(first), gen::getPartition, 
kv -> 1);
+        Snapshot firstSnapshot = store.snapshotManager().latestSnapshot();
+        assertThat(firstSnapshot.properties())
+                .containsEntry(SequenceSnapshotProperties.MAX_SEQUENCE_NUMBER, 
"0");
+
+        KeyValue second = gen.nextInsert("20201110", 10, 2L, new int[] {2, 2}, 
"second");
+        store.commitData(Collections.singletonList(second), gen::getPartition, 
kv -> 1);
+        Snapshot secondSnapshot = store.snapshotManager().latestSnapshot();
+        assertThat(secondSnapshot.properties())
+                .containsEntry(SequenceSnapshotProperties.MAX_SEQUENCE_NUMBER, 
"1");
+    }
+
+    @Test
+    public void 
testSnapshotSequenceNumberInitBootstrapsGlobalMaxSequenceNumber() throws 
Exception {
+        Map<String, String> options = new HashMap<>();
+        options.put(CoreOptions.WRITE_ONLY.key(), "true");
+        TestFileStore store = createStoreWithOptions(options);
+        TestKeyValueGenerator gen = new TestKeyValueGenerator();
+
+        store.commitData(
+                Collections.singletonList(
+                        gen.nextInsert("20201110", 10, 1L, new int[] {1, 1}, 
"first")),
+                gen::getPartition,
+                kv -> 2);
+        store.commitData(
+                Collections.singletonList(
+                        gen.nextInsert("20201110", 10, 2L, new int[] {2, 2}, 
"second")),
+                gen::getPartition,
+                kv -> 2);
+        store.commitData(
+                Collections.singletonList(
+                        gen.nextInsert("20201110", 10, 3L, new int[] {3, 3}, 
"third")),
+                gen::getPartition,
+                kv -> 2);
+        assertThat(
+                        SequenceSnapshotProperties.maxSequenceNumber(
+                                store.snapshotManager().latestSnapshot()))
+                .isEmpty();
+
+        SchemaManager schemaManager =
+                new SchemaManager(LocalFileIO.create(), new 
Path(tempDir.toUri()));
+        TableSchema snapshotSchema =
+                schemaManager.commitChanges(
+                        SchemaChange.setOption(
+                                
CoreOptions.WRITE_SEQUENCE_NUMBER_INIT_MODE.key(), "snapshot"));
+        TestFileStore snapshotStore = createStore(snapshotSchema);
+
+        snapshotStore.commitData(
+                Collections.singletonList(
+                        gen.nextInsert("20201110", 10, 4L, new int[] {4, 4}, 
"fourth")),
+                gen::getPartition,
+                kv -> 1);
+        Snapshot bootstrappedSnapshot = 
snapshotStore.snapshotManager().latestSnapshot();
+        assertThat(bootstrappedSnapshot.properties())
+                .containsEntry(SequenceSnapshotProperties.MAX_SEQUENCE_NUMBER, 
"2");
+
+        snapshotStore.commitData(
+                Collections.singletonList(
+                        gen.nextInsert("20201110", 10, 5L, new int[] {5, 5}, 
"fifth")),
+                gen::getPartition,
+                kv -> 1);
+        Snapshot restoredSnapshot = 
snapshotStore.snapshotManager().latestSnapshot();
+        assertThat(restoredSnapshot.properties())
+                .containsEntry(SequenceSnapshotProperties.MAX_SEQUENCE_NUMBER, 
"3");
+    }
+
+    @Test
+    public void testSnapshotSequenceNumberInitUsesWriterCreationSnapshot() 
throws Exception {
+        Map<String, String> options = new HashMap<>();
+        options.put(CoreOptions.WRITE_ONLY.key(), "true");
+        options.put(CoreOptions.WRITE_SEQUENCE_NUMBER_INIT_MODE.key(), 
"snapshot");
+        TestFileStore snapshotStore = createStoreWithOptions(options);
+        TestKeyValueGenerator gen = new TestKeyValueGenerator();
+
+        KeyValue first = gen.nextInsert("20201110", 10, 1L, new int[] {1, 1}, 
"first");
+        snapshotStore.commitData(Collections.singletonList(first), 
gen::getPartition, kv -> 1);
+        Snapshot snapshotWithProperties = 
snapshotStore.snapshotManager().latestSnapshot();
+        assertThat(snapshotWithProperties.properties())
+                .containsEntry(SequenceSnapshotProperties.MAX_SEQUENCE_NUMBER, 
"0");
+
+        SchemaManager schemaManager =
+                new SchemaManager(LocalFileIO.create(), new 
Path(tempDir.toUri()));
+        TableSchema scanSchema =
+                schemaManager.commitChanges(
+                        SchemaChange.setOption(
+                                
CoreOptions.WRITE_SEQUENCE_NUMBER_INIT_MODE.key(), "scan"));
+        TestFileStore scanStore = createStore(scanSchema);
+
+        scanStore.commitData(
+                Collections.singletonList(
+                        gen.nextInsert("20201110", 10, 2L, new int[] {2, 2}, 
"second")),
+                gen::getPartition,
+                kv -> 2);
+        assertThat(
+                        SequenceSnapshotProperties.maxSequenceNumber(
+                                
snapshotStore.snapshotManager().latestSnapshot()))
+                .isEmpty();
+
+        KeyValueFileStoreWrite write = (KeyValueFileStoreWrite) 
snapshotStore.newWrite();
+        assertThat(write.startingMaxSequenceNumber(-1, 
snapshotWithProperties)).isEqualTo(0);
+    }
+
     private CompactStrategy createCompactStrategy(Map<String, String> options) 
throws Exception {
         KeyValueFileStoreWrite write = createWriteWithOptions(options);
         write.withIOManager(ioManager);
@@ -108,6 +222,10 @@ public class KeyValueFileStoreWriteTest {
 
     private KeyValueFileStoreWrite createWriteWithOptions(Map<String, String> 
options)
             throws Exception {
+        return (KeyValueFileStoreWrite) 
createStoreWithOptions(options).newWrite();
+    }
+
+    private TestFileStore createStoreWithOptions(Map<String, String> options) 
throws Exception {
         SchemaManager schemaManager =
                 new SchemaManager(LocalFileIO.create(), new 
Path(tempDir.toUri()));
 
@@ -120,19 +238,20 @@ public class KeyValueFileStoreWriteTest {
                                         
TestKeyValueGenerator.GeneratorMode.MULTI_PARTITIONED),
                                 options,
                                 null));
-        TestFileStore store =
-                new TestFileStore.Builder(
-                                "avro",
-                                tempDir.toString(),
-                                NUM_BUCKETS,
-                                TestKeyValueGenerator.DEFAULT_PART_TYPE,
-                                TestKeyValueGenerator.KEY_TYPE,
-                                TestKeyValueGenerator.DEFAULT_ROW_TYPE,
-                                
TestKeyValueGenerator.TestKeyValueFieldsExtractor.EXTRACTOR,
-                                DeduplicateMergeFunction.factory(),
-                                schema)
-                        .build();
-
-        return (KeyValueFileStoreWrite) store.newWrite();
+        return createStore(schema);
+    }
+
+    private TestFileStore createStore(TableSchema schema) {
+        return new TestFileStore.Builder(
+                        "avro",
+                        tempDir.toString(),
+                        NUM_BUCKETS,
+                        TestKeyValueGenerator.DEFAULT_PART_TYPE,
+                        TestKeyValueGenerator.KEY_TYPE,
+                        TestKeyValueGenerator.DEFAULT_ROW_TYPE,
+                        
TestKeyValueGenerator.TestKeyValueFieldsExtractor.EXTRACTOR,
+                        DeduplicateMergeFunction.factory(),
+                        schema)
+                .build();
     }
 }
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/PrimaryKeyFileStoreTableITCase.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/PrimaryKeyFileStoreTableITCase.java
index 94bd5a30e7..a39f7bc725 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/PrimaryKeyFileStoreTableITCase.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/PrimaryKeyFileStoreTableITCase.java
@@ -171,6 +171,36 @@ public class PrimaryKeyFileStoreTableITCase extends 
AbstractTestBase {
     //  Constructed Tests
     // ------------------------------------------------------------------------
 
+    @Test
+    @Timeout(TIMEOUT)
+    public void 
testWriteOnlySnapshotSequenceNumberInitOverwritePreviousValue() throws 
Exception {
+        TableEnvironment bEnv = 
tableEnvironmentBuilder().batchMode().parallelism(1).build();
+        bEnv.executeSql(createCatalogSql("testCatalog", path));
+        bEnv.executeSql("USE CATALOG testCatalog");
+        bEnv.executeSql(
+                "CREATE TABLE T ("
+                        + "  k INT,"
+                        + "  v STRING,"
+                        + "  PRIMARY KEY (k) NOT ENFORCED"
+                        + ") WITH ("
+                        + "  'bucket' = '1',"
+                        + "  'write-only' = 'true',"
+                        + "  'write.sequence-number-init-mode' = 'snapshot'"
+                        + ")");
+
+        bEnv.executeSql("INSERT INTO T VALUES (1, 'old'), (2, 
'keep')").await();
+        bEnv.executeSql("INSERT INTO T VALUES (1, 'new')").await();
+
+        List<Row> actual = new ArrayList<>();
+        try (CloseableIterator<Row> it = bEnv.executeSql("SELECT * FROM T 
ORDER BY k").collect()) {
+            while (it.hasNext()) {
+                actual.add(it.next());
+            }
+        }
+
+        assertThat(actual).containsExactly(Row.of(1, "new"), Row.of(2, 
"keep"));
+    }
+
     @Test
     @Timeout(TIMEOUT)
     public void testFullCompactionTriggerInterval() throws Exception {

Reply via email to