Skip to content

Commit 06aa37b

Browse files
authored
use proper table naming in iceberg sql (#39035)
1 parent 5334dc1 commit 06aa37b

3 files changed

Lines changed: 20 additions & 3 deletions

File tree

sdks/java/extensions/sql/iceberg/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergMetastore.java

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -94,7 +94,7 @@ public void dropTable(String tableName) {
9494
@Override
9595
public Map<String, Table> getTables() {
9696
for (String id : catalogConfig.listTables(database)) {
97-
String name = TableName.create(id).getTableName();
97+
String name = tableName(id);
9898
@Nullable Table cachedTable = cachedTables.get(name);
9999
if (cachedTable == null) {
100100
Table table = checkStateNotNull(loadTable(id));
@@ -116,6 +116,10 @@ public Map<String, Table> getTables() {
116116
return table;
117117
}
118118

119+
private String tableName(String tableId) {
120+
return TableName.create(tableId).getTableName();
121+
}
122+
119123
private String getIdentifier(String name) {
120124
return database + "." + name;
121125
}
@@ -133,7 +137,7 @@ private String getIdentifier(Table table) {
133137
}
134138
return Table.builder()
135139
.type(getTableType())
136-
.name(identifier)
140+
.name(tableName(identifier))
137141
.schema(tableInfo.getSchema())
138142
.properties(TableUtils.parseProperties(tableInfo.getProperties()))
139143
.build();

sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergMetastoreTest.java

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323

2424
import java.io.File;
2525
import java.io.IOException;
26+
import java.util.Map;
2627
import java.util.UUID;
2728
import org.apache.beam.sdk.extensions.sql.meta.BeamSqlTable;
2829
import org.apache.beam.sdk.extensions.sql.meta.Table;
@@ -89,6 +90,18 @@ public void testGetTables() {
8990
assertEquals(ImmutableSet.of("my_table_1", "my_table_2"), metastore().getTables().keySet());
9091
}
9192

93+
@Test
94+
public void testGetTablesLoadsUncachedTablesWithSimpleNames() {
95+
String tableName = "uncached_table";
96+
catalog.catalogConfig.createTable(
97+
catalog.currentDatabase() + "." + tableName, Schema.of(), null, null);
98+
99+
Map<String, Table> tables = metastore().getTables();
100+
101+
assertEquals(ImmutableSet.of(tableName), tables.keySet());
102+
assertEquals(tableName, tables.get(tableName).getName());
103+
}
104+
92105
@Test
93106
public void testSupportsPartitioning() {
94107
Table table = Table.builder().name("my_table_1").schema(Schema.of()).type("iceberg").build();

sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergCatalogConfig.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -276,7 +276,7 @@ public boolean dropTable(String tableIdentifier) {
276276

277277
public Set<String> listTables(String namespace) {
278278
return catalog().listTables(Namespace.of(namespace)).stream()
279-
.map(TableIdentifier::name)
279+
.map(TableIdentifier::toString)
280280
.collect(Collectors.toSet());
281281
}
282282

0 commit comments

Comments
 (0)