diff --git a/paimon-core/src/main/java/org/apache/paimon/AbstractFileStore.java b/paimon-core/src/main/java/org/apache/paimon/AbstractFileStore.java index b88798b9246d..ebf64c952068 100644 --- a/paimon-core/src/main/java/org/apache/paimon/AbstractFileStore.java +++ b/paimon-core/src/main/java/org/apache/paimon/AbstractFileStore.java @@ -70,6 +70,7 @@ import org.apache.paimon.tag.TagAutoManager; import org.apache.paimon.tag.TagPreview; import org.apache.paimon.types.RowType; +import org.apache.paimon.utils.BranchManager; import org.apache.paimon.utils.ChainTableUtils; import org.apache.paimon.utils.ChangelogManager; import org.apache.paimon.utils.FileStorePathFactory; @@ -398,8 +399,7 @@ private List createCommitPreCallbacks(FileStoreTable table) { if (options.isChainTable()) { callbacks.add(new ChainTableCommitPreCallback(table)); } - if (options.toConfiguration().get(IcebergOptions.METADATA_ICEBERG_STORAGE) - != IcebergOptions.StorageType.DISABLED) { + if (icebergCompatibilityEnabled(table)) { callbacks.add(new IcebergPreCommitValidation(table)); } return callbacks; @@ -433,8 +433,7 @@ private List createCommitCallbacks(String commitUser, FileStoreT } } - if (options.toConfiguration().get(IcebergOptions.METADATA_ICEBERG_STORAGE) - != IcebergOptions.StorageType.DISABLED) { + if (icebergCompatibilityEnabled(table)) { callbacks.add(new IcebergCommitCallback(table, commitUser)); } @@ -597,13 +596,19 @@ public List createTagCallbacks(FileStoreTable table) { if (options.tagCreateSuccessFile()) { callbacks.add(new SuccessFileTagCallback(fileIO, newTagManager().tagDirectory())); } - if (options.toConfiguration().get(IcebergOptions.METADATA_ICEBERG_STORAGE) - != IcebergOptions.StorageType.DISABLED) { + if (icebergCompatibilityEnabled(table)) { callbacks.add(new IcebergCommitCallback(table, "")); } return callbacks; } + private boolean icebergCompatibilityEnabled(FileStoreTable table) { + return options.toConfiguration().get(IcebergOptions.METADATA_ICEBERG_STORAGE) + != IcebergOptions.StorageType.DISABLED + && BranchManager.isMainBranch( + BranchManager.normalizeBranch(table.coreOptions().branch())); + } + @Override public ServiceManager newServiceManager() { return new ServiceManager(fileIO, options.path()); 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 acaf154edcf5..13d450c1db24 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 @@ -38,7 +38,6 @@ import org.apache.paimon.iceberg.manifest.IcebergManifestFile; import org.apache.paimon.iceberg.manifest.IcebergManifestFileMeta; import org.apache.paimon.iceberg.manifest.IcebergManifestList; -import org.apache.paimon.iceberg.metadata.IcebergDataField; import org.apache.paimon.iceberg.metadata.IcebergMetadata; import org.apache.paimon.iceberg.metadata.IcebergRef; import org.apache.paimon.iceberg.metadata.IcebergSchema; @@ -1766,13 +1765,20 @@ public void testExistingTableWithUnpublishableHistoricalTimestampsRefusesToCommi } @Test - public void testCommitOnBranchMirrorsTheBranchSchema() throws Exception { + public void testCommitAndTagOnBranchDoNotModifyMainMetadata() throws Exception { + RecordingIcebergMetadataCommitter.COMMITS.clear(); RowType rowType = RowType.of( new DataType[] {DataTypes.INT(), DataTypes.INT()}, new String[] {"k", "v"}); FileStoreTable table = createPaimonTable( - rowType, Collections.emptyList(), Collections.singletonList("k"), 1); + rowType, + Collections.emptyList(), + Collections.singletonList("k"), + 1, + Collections.singletonMap( + IcebergOptions.METADATA_ICEBERG_STORAGE.key(), + IcebergOptions.StorageType.HADOOP_CATALOG.toString())); String commitUser = UUID.randomUUID().toString(); try (TableWriteImpl write = table.newWrite(commitUser); @@ -1785,25 +1791,35 @@ public void testCommitOnBranchMirrorsTheBranchSchema() throws Exception { new SchemaManager(table.fileIO(), table.location(), "b1") .commitChanges(SchemaChange.addColumn("branch_only", DataTypes.INT())); + try (TableWriteImpl write = table.newWrite(commitUser); + TableCommitImpl commit = table.newCommit(commitUser)) { + write.write(GenericRow.of(3, 30)); + commit.commit(2, write.prepareCommit(false, 2)); + } + + Path metadataDirectory = IcebergCommitCallback.catalogTableMetadataPath(table); + Path mainMetadataPathV1 = new Path(metadataDirectory, "v1.metadata.json"); + String mainMetadataV1 = table.fileIO().readFileUtf8(mainMetadataPathV1); + Path mainMetadataPathV2 = new Path(metadataDirectory, "v2.metadata.json"); + String mainMetadataV2 = table.fileIO().readFileUtf8(mainMetadataPathV2); + Path mainVersionHint = + new Path( + IcebergCommitCallback.catalogTableMetadataPath(table), "version-hint.text"); + String versionHint = table.fileIO().readFileUtf8(mainVersionHint); + FileStoreTable branchTable = table.switchToBranch("b1"); + RecordingIcebergMetadataCommitter.COMMITS.clear(); try (TableWriteImpl write = branchTable.newWrite(commitUser); TableCommitImpl commit = branchTable.newCommit(commitUser)) { write.write(GenericRow.of(2, 20, 200)); commit.commit(2, write.prepareCommit(false, 2)); } + branchTable.createTag("branch-tag"); - IcebergMetadata metadata = - IcebergMetadata.fromPath( - branchTable.fileIO(), - new Path( - branchTable.location(), - "metadata/v" - + branchTable.snapshotManager().latestSnapshotId() - + ".metadata.json")); - assertThat( - metadata.schemas().get(metadata.currentSchemaId()).fields().stream() - .map(IcebergDataField::name)) - .containsExactly("k", "v", "branch_only"); + assertThat(table.fileIO().readFileUtf8(mainMetadataPathV1)).isEqualTo(mainMetadataV1); + assertThat(table.fileIO().readFileUtf8(mainMetadataPathV2)).isEqualTo(mainMetadataV2); + assertThat(table.fileIO().readFileUtf8(mainVersionHint)).isEqualTo(versionHint); + assertThat(RecordingIcebergMetadataCommitter.COMMITS).isEmpty(); } /* @@ -2826,8 +2842,11 @@ private FileStoreTable createPaimonTable( Options options = new Options(customOptions); options.set(CoreOptions.BUCKET, numBuckets); - options.set( - IcebergOptions.METADATA_ICEBERG_STORAGE, IcebergOptions.StorageType.TABLE_LOCATION); + if (!options.contains(IcebergOptions.METADATA_ICEBERG_STORAGE)) { + options.set( + IcebergOptions.METADATA_ICEBERG_STORAGE, + IcebergOptions.StorageType.TABLE_LOCATION); + } options.set(CoreOptions.FILE_FORMAT, "avro"); options.set(CoreOptions.TARGET_FILE_SIZE, MemorySize.ofKibiBytes(32)); options.set(IcebergOptions.COMPACT_MIN_FILE_NUM, 4);