From e2dfe700791c804dc5efc4676577b93d6f8be037 Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Sat, 26 Sep 2026 03:51:44 +0800 Subject: [PATCH 1/3] [core] Keep all sibling tags of a snapshot as Iceberg refs MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The ref rebuild kept only the first tag name of each snapshot, so with several tags on one snapshot the siblings silently disappeared from Iceberg on the next commit — VERSION AS OF for them failed, and tags that notifyCreation had just added vanished again. Flatten every tag name of a snapshot into its own ref. Assisted-by: GLM-5.3 --- .../paimon/iceberg/IcebergCommitCallback.java | 18 ++++++++--- .../iceberg/IcebergCompatibilityTest.java | 30 +++++++++++++++++++ 2 files changed, 44 insertions(+), 4 deletions(-) diff --git a/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java b/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java index ab731753969e..46e8772115ee 100644 --- a/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java +++ b/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java @@ -85,6 +85,7 @@ import java.io.FileNotFoundException; import java.io.IOException; import java.io.UncheckedIOException; +import java.util.AbstractMap; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; @@ -1216,13 +1217,22 @@ private void createMetadataWithBase( // exists. Otherwise an Iceberg client fails to parse the metadata and all reads fail. Set snapshotIds = snapshots.stream().map(IcebergSnapshot::snapshotId).collect(Collectors.toSet()); + // every tag of a snapshot becomes its own ref: dropping sibling tags would make + // them invisible in Iceberg and break VERSION AS OF for them Map refs = table.tagManager().tags().entrySet().stream() .filter(entry -> snapshotIds.contains(entry.getKey().id())) - .collect( - Collectors.toMap( - entry -> entry.getValue().get(0), - entry -> new IcebergRef(entry.getKey().id()))); + .flatMap( + entry -> + entry.getValue().stream() + .map( + name -> + new AbstractMap.SimpleEntry<>( + name, + new IcebergRef( + entry.getKey() + .id())))) + .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); IcebergMetadata metadata = new IcebergMetadata( diff --git a/paimon-core/src/test/java/org/apache/paimon/iceberg/IcebergCompatibilityTest.java b/paimon-core/src/test/java/org/apache/paimon/iceberg/IcebergCompatibilityTest.java index 84a987ba5503..c3e32bc89afc 100644 --- a/paimon-core/src/test/java/org/apache/paimon/iceberg/IcebergCompatibilityTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/iceberg/IcebergCompatibilityTest.java @@ -1816,6 +1816,36 @@ public void testCommitOnBranchMirrorsTheBranchSchema() throws Exception { .containsExactly("k", "v", "branch_only"); } + @Test + public void testSiblingTagsOnOneSnapshotAllBecomeRefs() throws Exception { + RowType rowType = + RowType.of( + new DataType[] {DataTypes.INT(), DataTypes.INT()}, new String[] {"k", "v"}); + FileStoreTable table = + createPaimonTable( + rowType, Collections.emptyList(), Collections.singletonList("k"), 1); + + String commitUser = UUID.randomUUID().toString(); + TableWriteImpl write = table.newWrite(commitUser); + TableCommitImpl commit = table.newCommit(commitUser); + + write.write(GenericRow.of(1, 10)); + commit.commit(1, write.prepareCommit(false, 1)); + + table.createTag("first", 1); + table.createTag("second", 1); + + // a later write commit rebuilds the refs from all tags; both sibling tags must + // survive, not only the first one of the snapshot + write.write(GenericRow.of(2, 20)); + commit.commit(2, write.prepareCommit(false, 2)); + + long latestSnapshotId = table.snapshotManager().latestSnapshotId(); + Map refs = getIcebergRefsFromSnapshot(table, latestSnapshotId); + assertThat(refs.get("first").snapshotId()).isEqualTo(1); + assertThat(refs.get("second").snapshotId()).isEqualTo(1); + } + /* Create snapshots Create tags From 64180648f5d91a1196bfb0d0eebbab1065c1f78d Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Fri, 2 Oct 2026 12:35:49 +0800 Subject: [PATCH 2/3] [core] Pin the sibling-tag test to the tag-list rebuild --- .../paimon/iceberg/IcebergCommitCallback.java | 24 +++++++------------ .../iceberg/IcebergCompatibilityTest.java | 10 ++++---- 2 files changed, 15 insertions(+), 19 deletions(-) diff --git a/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java b/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java index 46e8772115ee..af2f3c3dbdc4 100644 --- a/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java +++ b/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java @@ -85,7 +85,6 @@ import java.io.FileNotFoundException; import java.io.IOException; import java.io.UncheckedIOException; -import java.util.AbstractMap; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; @@ -1219,20 +1218,15 @@ private void createMetadataWithBase( snapshots.stream().map(IcebergSnapshot::snapshotId).collect(Collectors.toSet()); // every tag of a snapshot becomes its own ref: dropping sibling tags would make // them invisible in Iceberg and break VERSION AS OF for them - Map refs = - table.tagManager().tags().entrySet().stream() - .filter(entry -> snapshotIds.contains(entry.getKey().id())) - .flatMap( - entry -> - entry.getValue().stream() - .map( - name -> - new AbstractMap.SimpleEntry<>( - name, - new IcebergRef( - entry.getKey() - .id())))) - .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); + Map refs = new HashMap<>(); + for (Map.Entry> entry : table.tagManager().tags().entrySet()) { + long taggedSnapshotId = entry.getKey().id(); + if (snapshotIds.contains(taggedSnapshotId)) { + for (String tagName : entry.getValue()) { + refs.put(tagName, new IcebergRef(taggedSnapshotId)); + } + } + } IcebergMetadata metadata = new IcebergMetadata( diff --git a/paimon-core/src/test/java/org/apache/paimon/iceberg/IcebergCompatibilityTest.java b/paimon-core/src/test/java/org/apache/paimon/iceberg/IcebergCompatibilityTest.java index c3e32bc89afc..dddeffb3606b 100644 --- a/paimon-core/src/test/java/org/apache/paimon/iceberg/IcebergCompatibilityTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/iceberg/IcebergCompatibilityTest.java @@ -1832,16 +1832,18 @@ public void testSiblingTagsOnOneSnapshotAllBecomeRefs() throws Exception { write.write(GenericRow.of(1, 10)); commit.commit(1, write.prepareCommit(false, 1)); - table.createTag("first", 1); - table.createTag("second", 1); + // create the tags without the Iceberg tag callback, so v1 metadata carries no refs and + // only the next commit's rebuild from the tag list can produce them + Snapshot snapshot = table.snapshotManager().snapshot(1); + table.tagManager().createTag(snapshot, "first", null, Collections.emptyList(), false); + table.tagManager().createTag(snapshot, "second", null, Collections.emptyList(), false); - // a later write commit rebuilds the refs from all tags; both sibling tags must - // survive, not only the first one of the snapshot write.write(GenericRow.of(2, 20)); commit.commit(2, write.prepareCommit(false, 2)); long latestSnapshotId = table.snapshotManager().latestSnapshotId(); Map refs = getIcebergRefsFromSnapshot(table, latestSnapshotId); + assertThat(refs).containsOnlyKeys("first", "second"); assertThat(refs.get("first").snapshotId()).isEqualTo(1); assertThat(refs.get("second").snapshotId()).isEqualTo(1); } From 9a69d01e96178ecac413be209a551c41e0a1c9b2 Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Sat, 3 Oct 2026 01:44:37 +0800 Subject: [PATCH 3/3] [flink] Cover sibling tags in the Iceberg tags ITCase --- .../flink/iceberg/FlinkIcebergITCaseBase.java | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/iceberg/FlinkIcebergITCaseBase.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/iceberg/FlinkIcebergITCaseBase.java index 80279a268861..8210443cbddd 100644 --- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/iceberg/FlinkIcebergITCaseBase.java +++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/iceberg/FlinkIcebergITCaseBase.java @@ -554,6 +554,19 @@ public void testCreateTags(String format) throws Exception { tEnv.executeSql( "SELECT name, type, snapshot_id FROM iceberg.`default`.T$refs"))) .containsExactlyInAnyOrder(Row.of("tag1", "TAG", 1L), Row.of("tag2", "TAG", 4L)); + + // a second tag on snapshot 4 must survive the next commit, which rebuilds the refs + tEnv.executeSql("CALL paimon.sys.create_tag('default.T', 'tag3', 4)"); + tEnv.executeSql("INSERT INTO paimon.`default`.T VALUES (1, 13, 131, 'black')").await(); + + assertThat( + collect( + tEnv.executeSql( + "SELECT name, type, snapshot_id FROM iceberg.`default`.T$refs"))) + .containsExactlyInAnyOrder( + Row.of("tag1", "TAG", 1L), + Row.of("tag2", "TAG", 4L), + Row.of("tag3", "TAG", 4L)); } @ParameterizedTest