This is an automated email from the ASF dual-hosted git repository.

zhoujinsong pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/amoro.git


The following commit(s) were added to refs/heads/master by this push:
     new 38b2b49e7 [AMORO-4240][AMS] Support staged table creation in REST 
catalog (#4284)
38b2b49e7 is described below

commit 38b2b49e72676cc5623264ce6f5d0dd94d3e414c
Author: Power John <[email protected]>
AuthorDate: Thu Jul 30 17:57:15 2026 +0800

    [AMORO-4240][AMS] Support staged table creation in REST catalog (#4284)
    
    * [AMORO-4240][AMS] Support staged table creation in REST catalog
    
    * [AMORO-4240][AMS] Address staged create review comments
---
 .../apache/amoro/server/RestCatalogService.java    | 122 ++++++++++++++++++---
 .../table/internal/InternalIcebergCreator.java     |  20 +++-
 .../internal/InternalMixedIcebergCreator.java      |  68 +++++++-----
 .../table/internal/InternalTableCreator.java       |  15 ++-
 .../server/TestInternalIcebergCatalogService.java  |  89 ++++++++++++++-
 5 files changed, 260 insertions(+), 54 deletions(-)

diff --git 
a/amoro-ams/src/main/java/org/apache/amoro/server/RestCatalogService.java 
b/amoro-ams/src/main/java/org/apache/amoro/server/RestCatalogService.java
index d197021a8..c13ed4013 100644
--- a/amoro-ams/src/main/java/org/apache/amoro/server/RestCatalogService.java
+++ b/amoro-ams/src/main/java/org/apache/amoro/server/RestCatalogService.java
@@ -55,8 +55,10 @@ import 
org.apache.amoro.shade.guava32.com.google.common.collect.Maps;
 import org.apache.amoro.shade.guava32.com.google.common.collect.Sets;
 import org.apache.amoro.utils.CatalogUtil;
 import org.apache.amoro.utils.TablePropertyUtil;
+import org.apache.iceberg.MetadataUpdate;
 import org.apache.iceberg.TableMetadata;
 import org.apache.iceberg.TableOperations;
+import org.apache.iceberg.UpdateRequirement;
 import org.apache.iceberg.catalog.Namespace;
 import org.apache.iceberg.catalog.TableIdentifier;
 import org.apache.iceberg.exceptions.AlreadyExistsException;
@@ -317,6 +319,7 @@ public class RestCatalogService extends PersistentBase 
implements RestExtension
         ctx,
         (catalog, database) -> {
           CreateTableRequest request = bodyAsClass(ctx, 
CreateTableRequest.class);
+          request.validate();
           String tableName = request.name();
           TableFormat format =
               TablePropertyUtil.isBaseStore(request.properties(), 
TableFormat.MIXED_ICEBERG)
@@ -325,6 +328,11 @@ public class RestCatalogService extends PersistentBase 
implements RestExtension
 
           try (InternalTableCreator creator =
               catalog.newTableCreator(database, tableName, format, request)) {
+            if (request.stageCreate()) {
+              checkAlreadyExists(!catalog.tableExists(database, tableName), 
"Table", tableName);
+              return 
LoadTableResponse.builder().withTableMetadata(creator.stage()).build();
+            }
+
             try {
               org.apache.amoro.server.table.TableMetadata metadata = 
creator.create();
               tableManager.createTable(catalog.name(), metadata);
@@ -358,25 +366,105 @@ public class RestCatalogService extends PersistentBase 
implements RestExtension
 
   /** POST PREFIX/v1/catalogs/{catalog}/namespaces/{namespace}/tables/{table} 
*/
   public void commitTable(Context ctx) {
-    handleTable(
-        ctx,
-        handler -> {
-          UpdateTableRequest request = bodyAsClass(ctx, 
UpdateTableRequest.class);
-          TableOperations ops = handler.newTableOperator();
-          TableMetadata base = ops.current();
-          if (base == null) {
-            throw new CommitFailedException("table metadata lost.");
-          }
+    UpdateTableRequest request = bodyAsClass(ctx, UpdateTableRequest.class);
+    request.validate();
+    if (isCreate(request)) {
+      handleNamespace(
+          ctx,
+          (catalog, database) -> {
+            String tableName = ctx.pathParam("table");
+            Preconditions.checkNotNull(tableName, "table name is null");
+            return commitCreateTable(catalog, database, tableName, request);
+          });
+    } else {
+      handleTable(
+          ctx,
+          handler -> {
+            TableOperations ops = handler.newTableOperator();
+            TableMetadata base = ops.current();
+            if (base == null) {
+              throw new CommitFailedException("table metadata lost.");
+            }
 
-          TableMetadata.Builder builder = TableMetadata.buildFrom(base);
-          request.requirements().forEach(r -> r.validate(base));
-          request.updates().forEach(u -> u.applyTo(builder));
-          TableMetadata newMetadata = builder.build();
+            TableMetadata.Builder builder = TableMetadata.buildFrom(base);
+            request.requirements().forEach(r -> r.validate(base));
+            request.updates().forEach(u -> u.applyTo(builder));
+            TableMetadata newMetadata = builder.build();
 
-          ops.commit(base, newMetadata);
-          TableMetadata current = ops.current();
-          return 
LoadTableResponse.builder().withTableMetadata(current).build();
-        });
+            ops.commit(base, newMetadata);
+            TableMetadata current = ops.current();
+            return 
LoadTableResponse.builder().withTableMetadata(current).build();
+          });
+    }
+  }
+
+  private LoadTableResponse commitCreateTable(
+      InternalCatalog catalog, String database, String tableName, 
UpdateTableRequest request) {
+    request.requirements().forEach(requirement -> 
requirement.validate((TableMetadata) null));
+
+    Optional<Integer> formatVersion =
+        request.updates().stream()
+            .filter(MetadataUpdate.UpgradeFormatVersion.class::isInstance)
+            .map(MetadataUpdate.UpgradeFormatVersion.class::cast)
+            .map(MetadataUpdate.UpgradeFormatVersion::formatVersion)
+            .findFirst();
+    TableMetadata.Builder builder =
+        formatVersion.isPresent()
+            ? TableMetadata.buildFromEmpty(formatVersion.get())
+            : TableMetadata.buildFromEmpty();
+    request.updates().forEach(update -> update.applyTo(builder));
+    TableMetadata icebergMetadata = builder.build();
+
+    // InternalTableCreator currently requires a CreateTableRequest for 
initialization. The
+    // metadata reconstructed above is passed to create(...) and is the 
metadata that is persisted.
+    CreateTableRequest createRequest =
+        CreateTableRequest.builder()
+            .withName(tableName)
+            .withSchema(icebergMetadata.schema())
+            .withPartitionSpec(icebergMetadata.spec())
+            .withWriteOrder(icebergMetadata.sortOrder())
+            .withLocation(icebergMetadata.location())
+            .setProperties(icebergMetadata.properties())
+            .build();
+    TableFormat format =
+        TablePropertyUtil.isBaseStore(createRequest.properties(), 
TableFormat.MIXED_ICEBERG)
+            ? TableFormat.MIXED_ICEBERG
+            : TableFormat.ICEBERG;
+
+    try (InternalTableCreator creator =
+        catalog.newTableCreator(database, tableName, format, createRequest)) {
+      try {
+        org.apache.amoro.server.table.TableMetadata metadata = 
creator.create(icebergMetadata);
+        tableManager.createTable(catalog.name(), metadata);
+      } catch (RuntimeException e) {
+        creator.rollback();
+        throw e;
+      }
+    }
+
+    try (InternalTableHandler<TableOperations> handler =
+        catalog.newTableHandler(database, tableName)) {
+      return LoadTableResponse.builder()
+          .withTableMetadata(handler.newTableOperator().current())
+          .build();
+    }
+  }
+
+  private static boolean isCreate(UpdateTableRequest request) {
+    boolean create =
+        request.requirements().stream()
+            
.anyMatch(UpdateRequirement.AssertTableDoesNotExist.class::isInstance);
+    if (create) {
+      List<UpdateRequirement> invalidRequirements =
+          request.requirements().stream()
+              .filter(
+                  requirement ->
+                      !(requirement instanceof 
UpdateRequirement.AssertTableDoesNotExist))
+              .collect(Collectors.toList());
+      Preconditions.checkArgument(
+          invalidRequirements.isEmpty(), "Invalid create requirements: %s", 
invalidRequirements);
+    }
+    return create;
   }
 
   /** DELETE 
PREFIX/v1/catalogs/{catalog}/namespaces/{namespace}/tables/{table} */
diff --git 
a/amoro-ams/src/main/java/org/apache/amoro/server/table/internal/InternalIcebergCreator.java
 
b/amoro-ams/src/main/java/org/apache/amoro/server/table/internal/InternalIcebergCreator.java
index 31f3c2389..29892e6bd 100644
--- 
a/amoro-ams/src/main/java/org/apache/amoro/server/table/internal/InternalIcebergCreator.java
+++ 
b/amoro-ams/src/main/java/org/apache/amoro/server/table/internal/InternalIcebergCreator.java
@@ -80,15 +80,26 @@ public class InternalIcebergCreator implements 
InternalTableCreator {
             request.properties());
   }
 
+  @Override
+  public org.apache.iceberg.TableMetadata stage() {
+    checkClosed();
+    return icebergMetadata;
+  }
+
   @Override
   public TableMetadata create() {
+    return create(icebergMetadata);
+  }
+
+  @Override
+  public TableMetadata create(org.apache.iceberg.TableMetadata metadata) {
     checkClosed();
 
     String icebergMetadataFileLocation =
-        InternalTableUtil.genNewMetadataFileLocation(null, icebergMetadata);
+        InternalTableUtil.genNewMetadataFileLocation(null, metadata);
     TableMeta meta = new TableMeta();
-    meta.putToLocations(MetaTableProperties.LOCATION_KEY_TABLE, 
icebergMetadata.location());
-    meta.putToLocations(MetaTableProperties.LOCATION_KEY_BASE, 
icebergMetadata.location());
+    meta.putToLocations(MetaTableProperties.LOCATION_KEY_TABLE, 
metadata.location());
+    meta.putToLocations(MetaTableProperties.LOCATION_KEY_BASE, 
metadata.location());
     meta.setFormat(format().name());
     meta.putToProperties(
         InternalTableConstants.PROPERTIES_METADATA_LOCATION, 
icebergMetadataFileLocation);
@@ -96,10 +107,9 @@ public class InternalIcebergCreator implements 
InternalTableCreator {
     ServerTableIdentifier serverTableIdentifier =
         ServerTableIdentifier.of(catalogMeta.getCatalogName(), database, 
tableName, format());
     
meta.setTableIdentifier(serverTableIdentifier.getIdentifier().buildTableIdentifier());
-    // write metadata file.
     OutputFile outputFile = io.newOutputFile(icebergMetadataFileLocation);
     this.metadataFileLocation = icebergMetadataFileLocation;
-    TableMetadataParser.overwrite(icebergMetadata, outputFile);
+    TableMetadataParser.overwrite(metadata, outputFile);
     return new TableMetadata(serverTableIdentifier, meta, catalogMeta);
   }
 
diff --git 
a/amoro-ams/src/main/java/org/apache/amoro/server/table/internal/InternalMixedIcebergCreator.java
 
b/amoro-ams/src/main/java/org/apache/amoro/server/table/internal/InternalMixedIcebergCreator.java
index 34a58ad16..ca192d46f 100644
--- 
a/amoro-ams/src/main/java/org/apache/amoro/server/table/internal/InternalMixedIcebergCreator.java
+++ 
b/amoro-ams/src/main/java/org/apache/amoro/server/table/internal/InternalMixedIcebergCreator.java
@@ -59,46 +59,40 @@ public class InternalMixedIcebergCreator extends 
InternalIcebergCreator {
   }
 
   @Override
-  public TableMetadata create() {
-    Map<String, String> properties = request.properties();
-    Preconditions.checkArgument(
-        TablePropertyUtil.isBaseStore(properties, TableFormat.MIXED_ICEBERG),
-        "The table creation request must be base store of mixed-iceberg");
+  public org.apache.iceberg.TableMetadata stage() {
+    validate(icebergMetadata);
+    return super.stage();
+  }
 
-    PrimaryKeySpec keySpec =
-        TablePropertyUtil.parsePrimaryKeySpec(request.schema(), 
request.properties());
+  @Override
+  public TableMetadata create() {
+    return create(icebergMetadata);
+  }
 
-    if (keySpec.primaryKeyExisted()) {
-      TableIdentifier identifier = TableIdentifier.of(database, tableName);
-      TableIdentifier changeIdentifier = 
TablePropertyUtil.parseChangeIdentifier(properties);
-      String expectChangeStoreName =
-          identifier.name() + 
InternalTableConstants.CHANGE_STORE_TABLE_NAME_SUFFIX;
-      TableIdentifier expectChangeIdentifier =
-          TableIdentifier.of(identifier.namespace(), expectChangeStoreName);
-      Preconditions.checkArgument(
-          expectChangeIdentifier.equals(changeIdentifier),
-          "the change store identifier is not expected. expected: %s, but 
found %s",
-          expectChangeIdentifier.toString(),
-          changeIdentifier.toString());
-    }
+  @Override
+  public TableMetadata create(org.apache.iceberg.TableMetadata baseMetadata) {
+    validate(baseMetadata);
 
-    TableMetadata metadata = super.create();
+    TableMetadata metadata = super.create(baseMetadata);
     metadata
         .getProperties()
         .put(InternalTableConstants.MIXED_ICEBERG_BASED_REST, 
Boolean.toString(true));
+
+    PrimaryKeySpec keySpec =
+        TablePropertyUtil.parsePrimaryKeySpec(baseMetadata.schema(), 
baseMetadata.properties());
     if (!keySpec.primaryKeyExisted()) {
       return metadata;
     }
 
-    Map<String, String> changeProperties = 
Maps.newHashMap(request.properties());
+    Map<String, String> changeProperties = 
Maps.newHashMap(baseMetadata.properties());
     changeProperties.putAll(
         TablePropertyUtil.changeStoreProperties(keySpec, 
TableFormat.MIXED_ICEBERG));
     String changeTableLocation = metadata.getTableLocation() + "/change";
     org.apache.iceberg.TableMetadata changeMetadata =
         org.apache.iceberg.TableMetadata.newTableMetadata(
-            icebergMetadata.schema(),
-            icebergMetadata.spec(),
-            icebergMetadata.sortOrder(),
+            baseMetadata.schema(),
+            baseMetadata.spec(),
+            baseMetadata.sortOrder(),
             changeTableLocation,
             changeProperties);
     String changeMetadataFileLocation =
@@ -118,6 +112,30 @@ public class InternalMixedIcebergCreator extends 
InternalIcebergCreator {
     return metadata;
   }
 
+  private void validate(org.apache.iceberg.TableMetadata metadata) {
+    Map<String, String> properties = metadata.properties();
+    Preconditions.checkArgument(
+        TablePropertyUtil.isBaseStore(properties, TableFormat.MIXED_ICEBERG),
+        "The table creation request must be base store of mixed-iceberg");
+
+    PrimaryKeySpec keySpec =
+        TablePropertyUtil.parsePrimaryKeySpec(metadata.schema(), 
metadata.properties());
+
+    if (keySpec.primaryKeyExisted()) {
+      TableIdentifier identifier = TableIdentifier.of(database, tableName);
+      TableIdentifier changeIdentifier = 
TablePropertyUtil.parseChangeIdentifier(properties);
+      String expectChangeStoreName =
+          identifier.name() + 
InternalTableConstants.CHANGE_STORE_TABLE_NAME_SUFFIX;
+      TableIdentifier expectChangeIdentifier =
+          TableIdentifier.of(identifier.namespace(), expectChangeStoreName);
+      Preconditions.checkArgument(
+          expectChangeIdentifier.equals(changeIdentifier),
+          "the change store identifier is not expected. expected: %s, but 
found %s",
+          expectChangeIdentifier.toString(),
+          changeIdentifier.toString());
+    }
+  }
+
   @Override
   public void rollback() {
     super.rollback();
diff --git 
a/amoro-ams/src/main/java/org/apache/amoro/server/table/internal/InternalTableCreator.java
 
b/amoro-ams/src/main/java/org/apache/amoro/server/table/internal/InternalTableCreator.java
index a91b02d55..c22386a0b 100644
--- 
a/amoro-ams/src/main/java/org/apache/amoro/server/table/internal/InternalTableCreator.java
+++ 
b/amoro-ams/src/main/java/org/apache/amoro/server/table/internal/InternalTableCreator.java
@@ -25,16 +25,19 @@ import java.io.Closeable;
 /** Interface to create an internal table. */
 public interface InternalTableCreator extends Closeable {
 
-  /**
-   * Do all things about create an internal table, and prepare the {@link 
TableMetadata} for {@link
-   * InternalTableManager#createTable(java.lang.String, TableMetadata)}
-   */
+  /** Build Iceberg metadata without writing files or registering the table. */
+  org.apache.iceberg.TableMetadata stage();
+
+  /** Write Iceberg metadata and prepare the {@link TableMetadata} for AMS. */
   TableMetadata create();
 
-  /** Rollback all resource created during {@link #create()} */
+  /** Persist Iceberg metadata and prepare the table metadata for AMS. */
+  TableMetadata create(org.apache.iceberg.TableMetadata icebergMetadata);
+
+  /** Remove files written during {@link #create()}. */
   void rollback();
 
-  /** Release resource like {@link org.apache.iceberg.io.FileIO} */
+  /** Release resources such as {@link org.apache.iceberg.io.FileIO}. */
   @Override
   default void close() {}
 }
diff --git 
a/amoro-ams/src/test/java/org/apache/amoro/server/TestInternalIcebergCatalogService.java
 
b/amoro-ams/src/test/java/org/apache/amoro/server/TestInternalIcebergCatalogService.java
index 5fa41d72b..bf810fc11 100644
--- 
a/amoro-ams/src/test/java/org/apache/amoro/server/TestInternalIcebergCatalogService.java
+++ 
b/amoro-ams/src/test/java/org/apache/amoro/server/TestInternalIcebergCatalogService.java
@@ -36,7 +36,9 @@ import org.apache.ibatis.session.SqlSession;
 import org.apache.iceberg.AppendFiles;
 import org.apache.iceberg.DataFile;
 import org.apache.iceberg.FileScanTask;
+import org.apache.iceberg.HasTableOperations;
 import org.apache.iceberg.Table;
+import org.apache.iceberg.TableProperties;
 import org.apache.iceberg.Transaction;
 import org.apache.iceberg.UpdateProperties;
 import org.apache.iceberg.catalog.Namespace;
@@ -51,6 +53,7 @@ import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Nested;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -59,6 +62,8 @@ import java.net.URI;
 import java.net.http.HttpClient;
 import java.net.http.HttpRequest;
 import java.net.http.HttpResponse;
+import java.nio.file.Files;
+import java.nio.file.Path;
 import java.sql.PreparedStatement;
 import java.sql.SQLException;
 import java.util.Arrays;
@@ -237,7 +242,9 @@ public class TestInternalIcebergCatalogService extends 
RestCatalogServiceTestBas
 
     @AfterEach
     public void clean() {
-      nsCatalog.dropTable(identifier);
+      if (nsCatalog.tableExists(identifier)) {
+        nsCatalog.dropTable(identifier);
+      }
       if (serverCatalog.tableExists(database, table)) {
         serverCatalog.dropTable(database, table);
       }
@@ -284,6 +291,82 @@ public class TestInternalIcebergCatalogService extends 
RestCatalogServiceTestBas
       Assertions.assertEquals(namespaceLocation + "/" + table, 
created.location());
     }
 
+    @Test
+    public void testStagedCreate(@TempDir Path tempDir) {
+      Path tablePath = tempDir.resolve("staged-table");
+      String tableLocation = tablePath.toUri().toString();
+      Transaction transaction =
+          nsCatalog
+              .buildTable(identifier, schema)
+              .withLocation(tableLocation)
+              .withProperty("owner", "analytics")
+              .createTransaction();
+
+      Assertions.assertFalse(serverCatalog.tableExists(database, table));
+      Assertions.assertFalse(Files.exists(tablePath));
+
+      transaction.commitTransaction();
+
+      Assertions.assertTrue(serverCatalog.tableExists(database, table));
+      Assertions.assertTrue(Files.exists(tablePath.resolve("metadata")));
+      Table loaded = nsCatalog.loadTable(identifier);
+      Assertions.assertEquals(2, formatVersion(loaded));
+      Assertions.assertEquals(tableLocation, loaded.location());
+      Assertions.assertEquals("analytics", loaded.properties().get("owner"));
+    }
+
+    @Test
+    public void testStagedCreateWithFormatVersionOne(@TempDir Path tempDir) {
+      Path tablePath = tempDir.resolve("staged-v1-table");
+      Transaction transaction =
+          nsCatalog
+              .buildTable(identifier, schema)
+              .withLocation(tablePath.toUri().toString())
+              .withProperty(TableProperties.FORMAT_VERSION, "1")
+              .createTransaction();
+
+      Assertions.assertEquals(1, formatVersion(transaction.table()));
+
+      transaction.commitTransaction();
+
+      Assertions.assertEquals(1, 
formatVersion(nsCatalog.loadTable(identifier)));
+    }
+
+    @Test
+    public void testStagedCreateWithAppendFiles(@TempDir Path tempDir) throws 
IOException {
+      Path tablePath = tempDir.resolve("staged-ctas-table");
+      Transaction transaction =
+          nsCatalog
+              .buildTable(identifier, schema)
+              .withLocation(tablePath.toUri().toString())
+              .createTransaction();
+
+      Assertions.assertFalse(serverCatalog.tableExists(database, table));
+      Assertions.assertFalse(Files.exists(tablePath));
+
+      DataFile[] files = IcebergDataTestHelpers.insert(transaction.table(), 
newRecords).dataFiles();
+      AppendFiles appendFiles = transaction.newAppend();
+      Arrays.stream(files).forEach(appendFiles::appendFile);
+      appendFiles.commit();
+
+      Assertions.assertFalse(serverCatalog.tableExists(database, table));
+
+      transaction.commitTransaction();
+
+      Assertions.assertTrue(serverCatalog.tableExists(database, table));
+      Table loaded = nsCatalog.loadTable(identifier);
+      Assertions.assertNotNull(loaded.currentSnapshot());
+      Set<String> expectedPaths =
+          Arrays.stream(files)
+              .map(dataFile -> dataFile.path().toString())
+              .collect(Collectors.toSet());
+      Set<String> actualPaths =
+          Streams.stream(loaded.newScan().planFiles())
+              .map(scanTask -> scanTask.file().path().toString())
+              .collect(Collectors.toSet());
+      Assertions.assertEquals(expectedPaths, actualPaths);
+    }
+
     @Test
     public void testTableWriteAndCommit() throws IOException {
       Table tbl = nsCatalog.createTable(identifier, schema);
@@ -353,5 +436,9 @@ public class TestInternalIcebergCatalogService extends 
RestCatalogServiceTestBas
           MixedDataTestHelpers.readBaseStore(mixedTable, reader, 
Expressions.alwaysTrue());
       Assertions.assertEquals(newRecords.size(), records.size());
     }
+
+    private int formatVersion(Table icebergTable) {
+      return ((HasTableOperations) 
icebergTable).operations().current().formatVersion();
+    }
   }
 }

Reply via email to