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..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 @@ -1216,13 +1216,17 @@ 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()); - 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()))); + // 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 = 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 84a987ba5503..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 @@ -1816,6 +1816,38 @@ 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)); + + // 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); + + 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); + } + /* Create snapshots Create tags 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