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();
+ }
}
}