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 54eab75c42 [core] Fix snapshot expiration for chain table delta
commits (#9611)
54eab75c42 is described below
commit 54eab75c42013acc4b1227f3ea80bd16376e46e9
Author: big face cat <[email protected]>
AuthorDate: Mon Sep 21 11:31:08 2026 +0800
[core] Fix snapshot expiration for chain table delta commits (#9611)
---
.../paimon/table/PrimaryKeyFileStoreTable.java | 21 +-
.../paimon/table/ChainTableSnapshotExpireTest.java | 281 +++++++++++++++++++++
2 files changed, 300 insertions(+), 2 deletions(-)
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/PrimaryKeyFileStoreTable.java
b/paimon-core/src/main/java/org/apache/paimon/table/PrimaryKeyFileStoreTable.java
index a478d17bfc..97c541175b 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/PrimaryKeyFileStoreTable.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/PrimaryKeyFileStoreTable.java
@@ -41,6 +41,7 @@ import org.apache.paimon.table.source.MergeTreeSplitGenerator;
import org.apache.paimon.table.source.PrimaryKeyBatchScan;
import org.apache.paimon.table.source.SplitGenerator;
import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.ChainTableUtils;
import org.apache.paimon.utils.RowKindFilter;
import javax.annotation.Nullable;
@@ -242,8 +243,24 @@ public class PrimaryKeyFileStoreTable extends
AbstractFileStoreTable {
protected Runnable newExpireRunnable() {
if (coreOptions().bucket() == BucketMode.POSTPONE_BUCKET) {
return null;
- } else {
- return super.newExpireRunnable();
}
+
+ Runnable expire = super.newExpireRunnable();
+ CoreOptions options = coreOptions();
+ if (expire == null ||
!ChainTableUtils.isScanFallbackDeltaBranch(options)) {
+ return expire;
+ }
+
+ FileStoreTable snapshotTable =
switchToBranch(options.scanFallbackSnapshotBranch());
+ // Use the Snapshot branch's own retention and changelog lifecycle
settings, not the
+ // Delta writer's runtime overrides.
+ ExpireSnapshots snapshotBranchExpire =
+ snapshotTable
+ .newExpireSnapshots()
+ .config(snapshotTable.coreOptions().expireConfig());
+ return () -> {
+ expire.run();
+ snapshotBranchExpire.expire();
+ };
}
}
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/ChainTableSnapshotExpireTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/ChainTableSnapshotExpireTest.java
new file mode 100644
index 0000000000..7971da1e4e
--- /dev/null
+++
b/paimon-core/src/test/java/org/apache/paimon/table/ChainTableSnapshotExpireTest.java
@@ -0,0 +1,281 @@
+/*
+ * 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;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.fs.Path;
+import org.apache.paimon.fs.local.LocalFileIO;
+import org.apache.paimon.manifest.PartitionEntry;
+import org.apache.paimon.options.Options;
+import org.apache.paimon.schema.FileSystemSchemaManager;
+import org.apache.paimon.schema.Schema;
+import org.apache.paimon.schema.SchemaChange;
+import org.apache.paimon.schema.SchemaManager;
+import org.apache.paimon.schema.TableSchema;
+import org.apache.paimon.table.sink.CommitMessage;
+import org.apache.paimon.table.sink.StreamTableWrite;
+import org.apache.paimon.table.sink.TableCommitImpl;
+import org.apache.paimon.types.DataTypes;
+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.time.LocalDate;
+import java.time.format.DateTimeFormatter;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.UUID;
+import java.util.stream.Collectors;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/** Tests snapshot expiration maintenance for chain table branches. */
+public class ChainTableSnapshotExpireTest {
+
+ private static final DateTimeFormatter PARTITION_FORMATTER =
+ DateTimeFormatter.ofPattern("yyyyMMdd");
+
+ @TempDir java.nio.file.Path tempDir;
+
+ @Test
+ public void
testDeltaCommitExpiresSnapshotBranchSnapshotsAfterPartitionExpiration()
+ throws Exception {
+ Path tablePath = new Path(tempDir.toUri().toString(),
"chain_snapshot_expire");
+ createChainTable(tablePath, Collections.emptyMap());
+
+ FileStoreTable mainTable = loadTable(tablePath);
+ FileStoreTable snapshotTable = mainTable.switchToBranch("snapshot");
+ FileStoreTable deltaTable = mainTable.switchToBranch("delta");
+ String commitUser = UUID.randomUUID().toString();
+
+ String day90 = partitionDate(-90);
+ String day65 = partitionDate(-65);
+ String day40 = partitionDate(-40);
+ String day10 = partitionDate(-10);
+
+ // Build three snapshot anchors while partition expiration is disabled.
+ write(snapshotTable, commitUser, day90, "v1");
+ write(snapshotTable, commitUser, day65, "v2");
+ write(snapshotTable, commitUser, day40, "v3");
+
+ snapshotTable = loadTable(tablePath).switchToBranch("snapshot");
+ assertThat(listPartitions(snapshotTable)).containsExactly(day90,
day65, day40);
+
assertThat(snapshotTable.snapshotManager().snapshotCount()).isEqualTo(3);
+ long latestSnapshotBeforeExpiration =
snapshotTable.snapshotManager().latestSnapshotId();
+
+ // Configure Snapshot retention on that branch itself, independently
of the Delta writer.
+ new FileSystemSchemaManager(LocalFileIO.create(), tablePath,
"snapshot")
+ .commitChanges(
+ Arrays.asList(
+ SchemaChange.setOption(
+
CoreOptions.SNAPSHOT_NUM_RETAINED_MIN.key(), "1"),
+ SchemaChange.setOption(
+
CoreOptions.SNAPSHOT_NUM_RETAINED_MAX.key(), "1"),
+ SchemaChange.setOption(
+
CoreOptions.SNAPSHOT_TIME_RETAINED.key(), "0 ms")));
+
+ Map<String, String> expireOptions = new HashMap<>();
+ expireOptions.put(CoreOptions.WRITE_ONLY.key(), "false");
+ expireOptions.put(CoreOptions.PARTITION_EXPIRATION_TIME.key(), "30 d");
+ expireOptions.put(CoreOptions.END_INPUT_CHECK_PARTITION_EXPIRE.key(),
"true");
+ expireOptions.put(CoreOptions.SNAPSHOT_NUM_RETAINED_MIN.key(), "1");
+ expireOptions.put(CoreOptions.SNAPSHOT_NUM_RETAINED_MAX.key(), "1");
+ expireOptions.put(CoreOptions.SNAPSHOT_TIME_RETAINED.key(), "0 ms");
+ expireOptions.put(CoreOptions.SNAPSHOT_EXPIRE_EXECUTION_MODE.key(),
"sync");
+ deltaTable = deltaTable.copy(expireOptions);
+
+ // A bounded Delta commit deterministically triggers
ChainTablePartitionExpire. The two
+ // oldest snapshot anchors are expired and the latest expired-time
anchor (day40) is kept.
+ // Dropping those Snapshot-branch partitions creates a new
Snapshot-branch metadata
+ // snapshot, which must then be expired according to the Snapshot
branch's own policy.
+ write(deltaTable, commitUser, day10, "v4");
+
+ snapshotTable = loadTable(tablePath).switchToBranch("snapshot");
+ assertThat(listPartitions(snapshotTable)).containsExactly(day40);
+ assertThat(snapshotTable.snapshotManager().latestSnapshotId())
+ .isGreaterThan(latestSnapshotBeforeExpiration);
+
assertThat(snapshotTable.snapshotManager().snapshotCount()).isEqualTo(1);
+ }
+
+ @ParameterizedTest
+ @CsvSource({"3, 3, false, 3", "1, 10, false, 4", "3, 3, true, 3"})
+ public void testDeltaCommitPreservesSnapshotBranchRetentionPolicy(
+ int retainMin, int retainMax, boolean decoupledChangelog, int
expectedSnapshots)
+ throws Exception {
+ Path tablePath = new Path(tempDir.toUri().toString(),
"chain_snapshot_retention");
+ Map<String, String> tableOptions = new HashMap<>();
+ tableOptions.put(CoreOptions.SNAPSHOT_NUM_RETAINED_MIN.key(),
String.valueOf(retainMin));
+ tableOptions.put(CoreOptions.SNAPSHOT_NUM_RETAINED_MAX.key(),
String.valueOf(retainMax));
+ tableOptions.put(CoreOptions.SNAPSHOT_TIME_RETAINED.key(), "365 d");
+ if (decoupledChangelog) {
+ tableOptions.put(CoreOptions.CHANGELOG_NUM_RETAINED_MIN.key(),
"5");
+ tableOptions.put(CoreOptions.CHANGELOG_NUM_RETAINED_MAX.key(),
"5");
+ tableOptions.put(CoreOptions.CHANGELOG_TIME_RETAINED.key(), "365
d");
+ }
+ // Persist the same policy on all branches. Only the Delta writer will
override it.
+ createChainTable(tablePath, tableOptions);
+
+ FileStoreTable mainTable = loadTable(tablePath);
+ FileStoreTable snapshotTable = mainTable.switchToBranch("snapshot");
+ FileStoreTable deltaTable = mainTable.switchToBranch("delta");
+ String commitUser = UUID.randomUUID().toString();
+ String day90 = partitionDate(-90);
+ String day65 = partitionDate(-65);
+ String day40 = partitionDate(-40);
+
+ // Seed both branches without triggering maintenance before the commit
under test.
+ Map<String, String> writeOnly =
+ Collections.singletonMap(CoreOptions.WRITE_ONLY.key(), "true");
+ FileStoreTable snapshotWriter = snapshotTable.copy(writeOnly);
+ write(snapshotWriter, commitUser, day90, "v1");
+ write(snapshotWriter, commitUser, day65, "v2");
+ write(snapshotWriter, commitUser, day40, "v3");
+ FileStoreTable deltaWriter = deltaTable.copy(writeOnly);
+ write(deltaWriter, commitUser, partitionDate(-10), "v4");
+ write(deltaWriter, commitUser, partitionDate(-9), "v5");
+ write(deltaWriter, commitUser, partitionDate(-8), "v6");
+
+
assertThat(snapshotTable.snapshotManager().snapshotCount()).isEqualTo(3);
+ assertThat(deltaTable.snapshotManager().snapshotCount()).isEqualTo(3);
+ long earliestSnapshot =
snapshotTable.snapshotManager().earliestSnapshotId();
+ long latestSnapshot =
snapshotTable.snapshotManager().latestSnapshotId();
+ assertThat(snapshotTable.coreOptions().changelogLifecycleDecoupled())
+ .isEqualTo(decoupledChangelog);
+
+ Map<String, String> expireOptions = new HashMap<>();
+ expireOptions.put(CoreOptions.WRITE_ONLY.key(), "false");
+ expireOptions.put(CoreOptions.PARTITION_EXPIRATION_TIME.key(), "30 d");
+ expireOptions.put(CoreOptions.END_INPUT_CHECK_PARTITION_EXPIRE.key(),
"true");
+ expireOptions.put(CoreOptions.SNAPSHOT_NUM_RETAINED_MIN.key(), "1");
+ expireOptions.put(CoreOptions.SNAPSHOT_NUM_RETAINED_MAX.key(), "1");
+ expireOptions.put(CoreOptions.SNAPSHOT_TIME_RETAINED.key(), "0 ms");
+ expireOptions.put(CoreOptions.CHANGELOG_NUM_RETAINED_MIN.key(), "1");
+ expireOptions.put(CoreOptions.CHANGELOG_NUM_RETAINED_MAX.key(), "1");
+ expireOptions.put(CoreOptions.CHANGELOG_TIME_RETAINED.key(), "0 ms");
+ expireOptions.put(CoreOptions.SNAPSHOT_EXPIRE_EXECUTION_MODE.key(),
"sync");
+ deltaWriter = deltaTable.copy(expireOptions);
+
assertThat(deltaWriter.coreOptions().changelogLifecycleDecoupled()).isFalse();
+ write(deltaWriter, commitUser, partitionDate(-7), "v7");
+
+ snapshotTable = loadTable(tablePath).switchToBranch("snapshot");
+ // Partition expiration still creates a metadata snapshot, but Delta's
shorter retention
+ // must not remove Snapshot history protected by its count or time
retention settings.
+ assertThat(listPartitions(snapshotTable)).containsExactly(day40);
+ assertThat(snapshotTable.snapshotManager().latestSnapshotId())
+ .isGreaterThan(latestSnapshot);
+
assertThat(snapshotTable.snapshotManager().snapshotCount()).isEqualTo(expectedSnapshots);
+ assertThat(deltaTable.snapshotManager().snapshotCount()).isEqualTo(1);
+ // A decoupled Snapshot lifecycle must archive the expired snapshot as
a changelog even
+ // though the Delta writer uses a coupled lifecycle.
+
assertThat(snapshotTable.changelogManager().longLivedChangelogExists(earliestSnapshot))
+ .isEqualTo(decoupledChangelog);
+ }
+
+ private void createChainTable(Path tablePath, Map<String, String>
tableOptions)
+ throws Exception {
+ LocalFileIO fileIO = LocalFileIO.create();
+ SchemaManager schemaManager = new FileSystemSchemaManager(fileIO,
tablePath);
+
+ Map<String, String> options = new HashMap<>(tableOptions);
+ options.put(CoreOptions.BUCKET.key(), "1");
+ options.put(CoreOptions.MERGE_ENGINE.key(), "deduplicate");
+ options.put(CoreOptions.SEQUENCE_FIELD.key(), "v");
+
+ Schema schema =
+ new Schema(
+ RowType.of(
+ new org.apache.paimon.types.DataType[]
{
+ DataTypes.STRING(),
+ DataTypes.STRING(),
+ DataTypes.STRING()
+ },
+ new String[] {"dt", "pk", "v"})
+ .getFields(),
+ Collections.singletonList("dt"),
+ Arrays.asList("pk", "dt"),
+ options,
+ "");
+ schemaManager.createTable(schema);
+
+ FileStoreTable mainTable = loadTable(tablePath);
+ mainTable.createBranch("snapshot");
+ mainTable.createBranch("delta");
+
+ List<SchemaChange> chainOptions =
+ Arrays.asList(
+
SchemaChange.setOption(CoreOptions.CHAIN_TABLE_ENABLED.key(), "true"),
+ SchemaChange.setOption(
+
CoreOptions.SCAN_FALLBACK_SNAPSHOT_BRANCH.key(), "snapshot"),
+ SchemaChange.setOption(
+ CoreOptions.SCAN_FALLBACK_DELTA_BRANCH.key(),
"delta"),
+ SchemaChange.setOption(
+ CoreOptions.PARTITION_TIMESTAMP_PATTERN.key(),
"$dt"),
+ SchemaChange.setOption(
+
CoreOptions.PARTITION_TIMESTAMP_FORMATTER.key(), "yyyyMMdd"));
+ schemaManager.commitChanges(chainOptions);
+ new FileSystemSchemaManager(fileIO, tablePath,
"snapshot").commitChanges(chainOptions);
+ new FileSystemSchemaManager(fileIO, tablePath,
"delta").commitChanges(chainOptions);
+ }
+
+ private FileStoreTable loadTable(Path tablePath) {
+ LocalFileIO fileIO = LocalFileIO.create();
+ Options options = new Options();
+ options.set(CoreOptions.PATH, tablePath.toString());
+ String branchName = CoreOptions.branch(options.toMap());
+ TableSchema tableSchema =
+ new FileSystemSchemaManager(fileIO, tablePath,
branchName).latest().get();
+ return FileStoreTableFactory.create(
+ fileIO, tablePath, tableSchema, CatalogEnvironment.empty());
+ }
+
+ private void write(FileStoreTable table, String commitUser, String dt,
String value)
+ throws Exception {
+ StreamTableWrite write = table.newWrite(commitUser);
+ write.write(
+ GenericRow.of(
+ BinaryString.fromString(dt),
+ BinaryString.fromString(value),
+ BinaryString.fromString(value)));
+ try (TableCommitImpl commit = table.newCommit(commitUser)) {
+ List<CommitMessage> commitMessages = write.prepareCommit(true,
Long.MAX_VALUE);
+ commit.commit(Long.MAX_VALUE, commitMessages);
+ }
+ write.close();
+ }
+
+ private List<String> listPartitions(FileStoreTable table) {
+ return table.newSnapshotReader().partitionEntries().stream()
+ .map(PartitionEntry::partition)
+ .map(partition -> partition.getString(0).toString())
+ .sorted()
+ .collect(Collectors.toList());
+ }
+
+ private String partitionDate(int daysFromToday) {
+ return
LocalDate.now().plusDays(daysFromToday).format(PARTITION_FORMATTER);
+ }
+}