diff --git a/sdks/java/extensions/sql/iceberg/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergMetastore.java b/sdks/java/extensions/sql/iceberg/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergMetastore.java index 5604b7c0837f..e6518f7192a5 100644 --- a/sdks/java/extensions/sql/iceberg/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergMetastore.java +++ b/sdks/java/extensions/sql/iceberg/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergMetastore.java @@ -94,7 +94,7 @@ public void dropTable(String tableName) { @Override public Map getTables() { for (String id : catalogConfig.listTables(database)) { - String name = TableName.create(id).getTableName(); + String name = tableName(id); @Nullable Table cachedTable = cachedTables.get(name); if (cachedTable == null) { Table table = checkStateNotNull(loadTable(id)); @@ -116,6 +116,10 @@ public Map getTables() { return table; } + private String tableName(String tableId) { + return TableName.create(tableId).getTableName(); + } + private String getIdentifier(String name) { return database + "." + name; } @@ -133,7 +137,7 @@ private String getIdentifier(Table table) { } return Table.builder() .type(getTableType()) - .name(identifier) + .name(tableName(identifier)) .schema(tableInfo.getSchema()) .properties(TableUtils.parseProperties(tableInfo.getProperties())) .build(); diff --git a/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergMetastoreTest.java b/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergMetastoreTest.java index a7baf1191d15..5e76b536134d 100644 --- a/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergMetastoreTest.java +++ b/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergMetastoreTest.java @@ -23,6 +23,7 @@ import java.io.File; import java.io.IOException; +import java.util.Map; import java.util.UUID; import org.apache.beam.sdk.extensions.sql.meta.BeamSqlTable; import org.apache.beam.sdk.extensions.sql.meta.Table; @@ -89,6 +90,18 @@ public void testGetTables() { assertEquals(ImmutableSet.of("my_table_1", "my_table_2"), metastore().getTables().keySet()); } + @Test + public void testGetTablesLoadsUncachedTablesWithSimpleNames() { + String tableName = "uncached_table"; + catalog.catalogConfig.createTable( + catalog.currentDatabase() + "." + tableName, Schema.of(), null, null); + + Map tables = metastore().getTables(); + + assertEquals(ImmutableSet.of(tableName), tables.keySet()); + assertEquals(tableName, tables.get(tableName).getName()); + } + @Test public void testSupportsPartitioning() { Table table = Table.builder().name("my_table_1").schema(Schema.of()).type("iceberg").build(); diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergCatalogConfig.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergCatalogConfig.java index 6fb0a3480a0f..454fe40269d6 100644 --- a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergCatalogConfig.java +++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergCatalogConfig.java @@ -276,7 +276,7 @@ public boolean dropTable(String tableIdentifier) { public Set listTables(String namespace) { return catalog().listTables(Namespace.of(namespace)).stream() - .map(TableIdentifier::name) + .map(TableIdentifier::toString) .collect(Collectors.toSet()); }