This is an automated email from the ASF dual-hosted git repository.
ahmedabu98 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 20ed72c83e3 [IcebergIO] Support TableIdentifiers with special
characters (#38876)
20ed72c83e3 is described below
commit 20ed72c83e3ad05281d46781f3da00715f8db193
Author: Ahmed Abualsaud <[email protected]>
AuthorDate: Mon Jul 13 17:41:45 2026 -0500
[IcebergIO] Support TableIdentifiers with special characters (#38876)
---
.../IO_Iceberg_Integration_Tests.json | 2 +-
.../org/apache/beam/sdk/io/iceberg/AddFiles.java | 8 +--
.../beam/sdk/io/iceberg/AppendFilesToTables.java | 5 +-
.../iceberg/AssignDestinationsAndPartitions.java | 3 +-
.../beam/sdk/io/iceberg/FileWriteResult.java | 4 +-
.../beam/sdk/io/iceberg/IcebergCatalogConfig.java | 6 +-
.../IcebergCdcReadSchemaTransformProvider.java | 3 +-
.../IcebergReadSchemaTransformProvider.java | 3 +-
.../beam/sdk/io/iceberg/IcebergScanConfig.java | 6 +-
.../apache/beam/sdk/io/iceberg/IcebergUtils.java | 54 ++++++++++++++++++
.../beam/sdk/io/iceberg/IncrementalScanSource.java | 4 +-
.../io/iceberg/OneTableDynamicDestinations.java | 6 +-
.../io/iceberg/PortableIcebergDestinations.java | 3 +-
.../apache/beam/sdk/io/iceberg/SnapshotInfo.java | 3 +-
.../org/apache/beam/sdk/io/iceberg/TableCache.java | 6 +-
.../beam/sdk/io/iceberg/IcebergIOReadTest.java | 3 +-
.../beam/sdk/io/iceberg/IcebergIOWriteTest.java | 3 +-
.../beam/sdk/io/iceberg/IcebergUtilsTest.java | 64 ++++++++++++++++++++++
18 files changed, 152 insertions(+), 34 deletions(-)
diff --git a/.github/trigger_files/IO_Iceberg_Integration_Tests.json
b/.github/trigger_files/IO_Iceberg_Integration_Tests.json
index b73af5e61a4..7ab7bcd9a9c 100644
--- a/.github/trigger_files/IO_Iceberg_Integration_Tests.json
+++ b/.github/trigger_files/IO_Iceberg_Integration_Tests.json
@@ -1,4 +1,4 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to
run.",
- "modification": 1
+ "modification": 2
}
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AddFiles.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AddFiles.java
index 9fad990a439..f37935f89e8 100644
---
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AddFiles.java
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AddFiles.java
@@ -523,7 +523,7 @@ public class AddFiles extends
PTransform<PCollection<String>, PCollectionRowTupl
}
private Table getOrCreateTable(String filePath, FileFormat format) throws
IOException {
- TableIdentifier tableId = TableIdentifier.parse(identifier);
+ TableIdentifier tableId = IcebergUtils.parseTableIdentifier(identifier);
try {
return catalogConfig.catalog().loadTable(tableId);
} catch (NoSuchTableException e) {
@@ -549,7 +549,7 @@ public class AddFiles extends
PTransform<PCollection<String>, PCollectionRowTupl
.create();
} catch (AlreadyExistsException e2) { // if table already exists, just
load it
- return
catalogConfig.catalog().loadTable(TableIdentifier.parse(identifier));
+ return
catalogConfig.catalog().loadTable(IcebergUtils.parseTableIdentifier(identifier));
}
}
}
@@ -684,7 +684,7 @@ public class AddFiles extends
PTransform<PCollection<String>, PCollectionRowTupl
return;
}
if (table == null) {
- table =
catalogConfig.catalog().loadTable(TableIdentifier.parse(identifier));
+ table =
catalogConfig.catalog().loadTable(IcebergUtils.parseTableIdentifier(identifier));
}
PartitionSpec spec =
checkStateNotNull(table.specs().get(batch.getKey()));
@@ -769,7 +769,7 @@ public class AddFiles extends
PTransform<PCollection<String>, PCollectionRowTupl
}
String commitId = commitHash(manifests);
if (table == null) {
- table =
catalogConfig.catalog().loadTable(TableIdentifier.parse(identifier));
+ table =
catalogConfig.catalog().loadTable(IcebergUtils.parseTableIdentifier(identifier));
}
table.refresh();
ensureNameMappingPresent(table);
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AppendFilesToTables.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AppendFilesToTables.java
index c7981d697de..917b087e39e 100644
---
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AppendFilesToTables.java
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AppendFilesToTables.java
@@ -47,7 +47,6 @@ import org.apache.iceberg.PartitionSpec;
import org.apache.iceberg.Snapshot;
import org.apache.iceberg.Table;
import org.apache.iceberg.catalog.Catalog;
-import org.apache.iceberg.catalog.TableIdentifier;
import org.apache.iceberg.io.FileIO;
import org.checkerframework.checker.nullness.qual.MonotonicNonNull;
import org.slf4j.Logger;
@@ -75,7 +74,7 @@ class AppendFilesToTables
new SerializableFunction<FileWriteResult, String>() {
@Override
public String apply(FileWriteResult input) {
- return input.getTableIdentifier().toString();
+ return
IcebergUtils.tableIdentifierToString(input.getTableIdentifier());
}
}))
.apply("Group metadata updates by table", GroupByKey.create())
@@ -128,7 +127,7 @@ class AppendFilesToTables
BoundedWindow window)
throws IOException {
String tableStringIdentifier = element.getKey();
- Table table =
getCatalog().loadTable(TableIdentifier.parse(element.getKey()));
+ Table table =
getCatalog().loadTable(IcebergUtils.parseTableIdentifier(element.getKey()));
Iterable<FileWriteResult> fileWriteResults = element.getValue();
if (shouldSkip(table, fileWriteResults)) {
return;
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AssignDestinationsAndPartitions.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AssignDestinationsAndPartitions.java
index 99cd07b23c8..364b44c37ee 100644
---
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AssignDestinationsAndPartitions.java
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AssignDestinationsAndPartitions.java
@@ -35,7 +35,6 @@ import org.apache.beam.sdk.values.ValueInSingleWindow;
import org.apache.iceberg.PartitionKey;
import org.apache.iceberg.PartitionSpec;
import org.apache.iceberg.Schema;
-import org.apache.iceberg.catalog.TableIdentifier;
import org.apache.iceberg.exceptions.NoSuchTableException;
import org.checkerframework.checker.nullness.qual.MonotonicNonNull;
import org.checkerframework.checker.nullness.qual.Nullable;
@@ -149,7 +148,7 @@ class AssignDestinationsAndPartitions
// see if table already exists with a spec
spec =
TableCache.getAndRefreshIfStale(
- catalogConfig, TableIdentifier.parse(tableIdentifier))
+ catalogConfig,
IcebergUtils.parseTableIdentifier(tableIdentifier))
.spec();
} catch (NoSuchTableException ignored) {
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/FileWriteResult.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/FileWriteResult.java
index b96b1d42c94..713b8029c15 100644
---
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/FileWriteResult.java
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/FileWriteResult.java
@@ -44,7 +44,7 @@ abstract class FileWriteResult {
@SchemaIgnore
public TableIdentifier getTableIdentifier() {
if (cachedTableIdentifier == null) {
- cachedTableIdentifier =
TableIdentifier.parse(getTableIdentifierString());
+ cachedTableIdentifier =
IcebergUtils.parseTableIdentifier(getTableIdentifierString());
}
return cachedTableIdentifier;
}
@@ -70,7 +70,7 @@ abstract class FileWriteResult {
@SchemaIgnore
public Builder setTableIdentifier(TableIdentifier tableId) {
- return setTableIdentifierString(tableId.toString());
+ return
setTableIdentifierString(IcebergUtils.tableIdentifierToString(tableId));
}
public abstract FileWriteResult 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 454fe40269d..8fe05cb5b21 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
@@ -155,7 +155,7 @@ public abstract class IcebergCatalogConfig implements
Serializable {
Schema tableSchema,
@Nullable List<String> partitionFields,
@Nullable Map<String, String> properties) {
- TableIdentifier icebergIdentifier = TableIdentifier.parse(tableIdentifier);
+ TableIdentifier icebergIdentifier =
IcebergUtils.parseTableIdentifier(tableIdentifier);
org.apache.iceberg.Schema icebergSchema =
IcebergUtils.beamSchemaToIcebergSchema(tableSchema);
PartitionSpec icebergSpec =
PartitionUtils.toPartitionSpec(partitionFields, tableSchema);
try {
@@ -178,7 +178,7 @@ public abstract class IcebergCatalogConfig implements
Serializable {
}
public @Nullable IcebergTableInfo loadTable(String tableIdentifier) {
- TableIdentifier icebergIdentifier = TableIdentifier.parse(tableIdentifier);
+ TableIdentifier icebergIdentifier =
IcebergUtils.parseTableIdentifier(tableIdentifier);
try {
Table table = catalog().loadTable(icebergIdentifier);
return new IcebergTableInfo(tableIdentifier, table);
@@ -270,7 +270,7 @@ public abstract class IcebergCatalogConfig implements
Serializable {
}
public boolean dropTable(String tableIdentifier) {
- TableIdentifier icebergIdentifier = TableIdentifier.parse(tableIdentifier);
+ TableIdentifier icebergIdentifier =
IcebergUtils.parseTableIdentifier(tableIdentifier);
return catalog().dropTable(icebergIdentifier);
}
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergCdcReadSchemaTransformProvider.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergCdcReadSchemaTransformProvider.java
index 31ff57a668b..e029a85a812 100644
---
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergCdcReadSchemaTransformProvider.java
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergCdcReadSchemaTransformProvider.java
@@ -41,7 +41,6 @@ import org.apache.beam.sdk.values.PCollectionRowTuple;
import org.apache.beam.sdk.values.Row;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Enums;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Optional;
-import org.apache.iceberg.catalog.TableIdentifier;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.joda.time.Duration;
@@ -109,7 +108,7 @@ public class IcebergCdcReadSchemaTransformProvider
IcebergIO.ReadRows readRows =
IcebergIO.readRows(configuration.getIcebergCatalog())
.withCdc()
- .from(TableIdentifier.parse(configuration.getTable()))
+
.from(IcebergUtils.parseTableIdentifier(configuration.getTable()))
.fromSnapshot(configuration.getFromSnapshot())
.toSnapshot(configuration.getToSnapshot())
.fromTimestamp(configuration.getFromTimestamp())
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergReadSchemaTransformProvider.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergReadSchemaTransformProvider.java
index 63d6f792e56..67ea6050507 100644
---
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergReadSchemaTransformProvider.java
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergReadSchemaTransformProvider.java
@@ -37,7 +37,6 @@ import
org.apache.beam.sdk.schemas.transforms.TypedSchemaTransformProvider;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.PCollectionRowTuple;
import org.apache.beam.sdk.values.Row;
-import org.apache.iceberg.catalog.TableIdentifier;
import org.checkerframework.checker.nullness.qual.Nullable;
/**
@@ -93,7 +92,7 @@ public class IcebergReadSchemaTransformProvider
.getPipeline()
.apply(
IcebergIO.readRows(configuration.getIcebergCatalog())
- .from(TableIdentifier.parse(configuration.getTable()))
+
.from(IcebergUtils.parseTableIdentifier(configuration.getTable()))
.keeping(configuration.getKeep())
.dropping(configuration.getDrop())
.withFilter(configuration.getFilter()));
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergScanConfig.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergScanConfig.java
index bdcf652a5df..95ea6cf1bd4 100644
---
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergScanConfig.java
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergScanConfig.java
@@ -80,7 +80,9 @@ public abstract class IcebergScanConfig implements
Serializable {
@Pure
public Table getTable() {
if (cachedTable == null) {
- cachedTable = TableCache.get(getCatalogConfig(),
TableIdentifier.parse(getTableIdentifier()));
+ cachedTable =
+ TableCache.get(
+ getCatalogConfig(),
IcebergUtils.parseTableIdentifier(getTableIdentifier()));
}
return cachedTable;
}
@@ -294,7 +296,7 @@ public abstract class IcebergScanConfig implements
Serializable {
public abstract Builder setTableIdentifier(String tableIdentifier);
public Builder setTableIdentifier(TableIdentifier tableIdentifier) {
- return this.setTableIdentifier(tableIdentifier.toString());
+ return
this.setTableIdentifier(IcebergUtils.tableIdentifierToString(tableIdentifier));
}
public Builder setTableIdentifier(String... names) {
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergUtils.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergUtils.java
index d0d24532ff3..309205707a9 100644
---
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergUtils.java
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergUtils.java
@@ -46,6 +46,8 @@ import org.apache.beam.sdk.values.Row;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
+import org.apache.iceberg.catalog.TableIdentifier;
+import org.apache.iceberg.catalog.TableIdentifierParser;
import org.apache.iceberg.data.GenericRecord;
import org.apache.iceberg.data.Record;
import org.apache.iceberg.types.Type;
@@ -631,6 +633,58 @@ public class IcebergUtils {
return icebergValue;
}
+ /** Serializes a table identifier without losing dots or other special
characters in each part. */
+ public static String tableIdentifierToString(TableIdentifier
tableIdentifier) {
+ TableIdentifier identifier = checkArgumentNotNull(tableIdentifier);
+ return requiresJsonTableIdentifier(identifier)
+ ? TableIdentifierParser.toJson(identifier)
+ : identifier.toString();
+ }
+
+ /** Parses either Iceberg's JSON table identifier representation or the
legacy dotted form. */
+ public static TableIdentifier parseTableIdentifier(String table) {
+ if (looksLikeJsonObject(table)) {
+ return TableIdentifierParser.fromJson(table);
+ }
+
+ return TableIdentifier.parse(table);
+ }
+
+ private static boolean looksLikeJsonObject(@Nullable String value) {
+ if (value == null) {
+ return false;
+ }
+
+ int start = 0;
+ while (start < value.length() &&
Character.isWhitespace(value.charAt(start))) {
+ start++;
+ }
+ if (start == value.length() || value.charAt(start) != '{') {
+ return false;
+ }
+
+ int end = value.length() - 1;
+ while (end >= 0 && Character.isWhitespace(value.charAt(end))) {
+ end--;
+ }
+
+ return end >= 0 && value.charAt(end) == '}';
+ }
+
+ private static boolean requiresJsonTableIdentifier(TableIdentifier
identifier) {
+ if (looksLikeJsonObject(identifier.toString())) {
+ return true;
+ }
+
+ for (String level : identifier.namespace().levels()) {
+ if (level.contains(".")) {
+ return true;
+ }
+ }
+
+ return identifier.name().contains(".");
+ }
+
static <T> boolean isUnbounded(PCollection<T> input) {
return input.isBounded().equals(PCollection.IsBounded.UNBOUNDED);
}
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IncrementalScanSource.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IncrementalScanSource.java
index 58cc8f50e0b..324eb817276 100644
---
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IncrementalScanSource.java
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IncrementalScanSource.java
@@ -33,7 +33,6 @@ import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.Row;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects;
import org.apache.iceberg.Table;
-import org.apache.iceberg.catalog.TableIdentifier;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.joda.time.Duration;
@@ -54,7 +53,8 @@ class IncrementalScanSource extends PTransform<PBegin,
PCollection<Row>> {
public PCollection<Row> expand(PBegin input) {
Table table =
TableCache.get(
- scanConfig.getCatalogConfig(),
TableIdentifier.parse(scanConfig.getTableIdentifier()));
+ scanConfig.getCatalogConfig(),
+
IcebergUtils.parseTableIdentifier(scanConfig.getTableIdentifier()));
PCollection<KV<String, List<SnapshotInfo>>> snapshots =
MoreObjects.firstNonNull(scanConfig.getStreaming(), false)
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/OneTableDynamicDestinations.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/OneTableDynamicDestinations.java
index 861a8ad198a..afca43fca15 100644
---
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/OneTableDynamicDestinations.java
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/OneTableDynamicDestinations.java
@@ -41,13 +41,13 @@ class OneTableDynamicDestinations implements
DynamicDestinations, Externalizable
@VisibleForTesting
TableIdentifier getTableIdentifier() {
if (tableId == null) {
- tableId = TableIdentifier.parse(checkStateNotNull(tableIdString));
+ tableId =
IcebergUtils.parseTableIdentifier(checkStateNotNull(tableIdString));
}
return tableId;
}
OneTableDynamicDestinations(TableIdentifier tableId, Schema rowSchema) {
- this.tableIdString = tableId.toString();
+ this.tableIdString = IcebergUtils.tableIdentifierToString(tableId);
this.rowSchema = rowSchema;
}
@@ -86,6 +86,6 @@ class OneTableDynamicDestinations implements
DynamicDestinations, Externalizable
@Override
public void readExternal(ObjectInput in) throws IOException {
tableIdString = in.readUTF();
- tableId = TableIdentifier.parse(tableIdString);
+ tableId = IcebergUtils.parseTableIdentifier(tableIdString);
}
}
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/PortableIcebergDestinations.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/PortableIcebergDestinations.java
index f2cec9bdfec..775020879f6 100644
---
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/PortableIcebergDestinations.java
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/PortableIcebergDestinations.java
@@ -25,7 +25,6 @@ import org.apache.beam.sdk.util.RowStringInterpolator;
import org.apache.beam.sdk.values.Row;
import org.apache.beam.sdk.values.ValueInSingleWindow;
import org.apache.iceberg.FileFormat;
-import org.apache.iceberg.catalog.TableIdentifier;
import org.checkerframework.checker.nullness.qual.Nullable;
class PortableIcebergDestinations implements DynamicDestinations {
@@ -84,7 +83,7 @@ class PortableIcebergDestinations implements
DynamicDestinations {
@Override
public IcebergDestination instantiateDestination(String dest) {
return IcebergDestination.builder()
- .setTableIdentifier(TableIdentifier.parse(dest))
+ .setTableIdentifier(IcebergUtils.parseTableIdentifier(dest))
.setTableCreateConfig(
IcebergTableCreateConfig.builder()
.setSchema(getDataSchema())
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SnapshotInfo.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SnapshotInfo.java
index bab5405cd4a..de8c400c08c 100644
---
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SnapshotInfo.java
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SnapshotInfo.java
@@ -94,7 +94,8 @@ public abstract class SnapshotInfo {
@SchemaIgnore
public TableIdentifier getTableIdentifier() {
if (cachedTableIdentifier == null) {
- cachedTableIdentifier =
TableIdentifier.parse(checkStateNotNull(getTableIdentifierString()));
+ cachedTableIdentifier =
+
IcebergUtils.parseTableIdentifier(checkStateNotNull(getTableIdentifierString()));
}
return cachedTableIdentifier;
}
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/TableCache.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/TableCache.java
index 63523bb995f..bc2fc25de30 100644
---
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/TableCache.java
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/TableCache.java
@@ -55,7 +55,7 @@ public class TableCache {
/** Returns the cached table for a string identifier, loading it on a cache
miss. */
public static Table get(IcebergCatalogConfig catalogConfig, String
identifier) {
- return get(catalogConfig, TableIdentifier.parse(identifier));
+ return get(catalogConfig, IcebergUtils.parseTableIdentifier(identifier));
}
/** Returns the cached table, using the given loader only on a cache miss. */
@@ -75,11 +75,11 @@ public class TableCache {
/** Returns the cached table for a string identifier after refreshing any
pre-existing entry. */
public static Table getRefreshed(IcebergCatalogConfig catalogConfig, String
identifier) {
- return getRefreshed(catalogConfig, TableIdentifier.parse(identifier));
+ return getRefreshed(catalogConfig,
IcebergUtils.parseTableIdentifier(identifier));
}
public static Table getAndRefreshIfStale(IcebergCatalogConfig catalogConfig,
String identifier) {
- return getAndRefreshIfStale(catalogConfig,
TableIdentifier.parse(identifier));
+ return getAndRefreshIfStale(catalogConfig,
IcebergUtils.parseTableIdentifier(identifier));
}
/**
diff --git
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOReadTest.java
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOReadTest.java
index d7c97efa19f..edd26145816 100644
---
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOReadTest.java
+++
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOReadTest.java
@@ -341,7 +341,8 @@ public class IcebergIOReadTest {
@Test
public void testSimpleScan() throws Exception {
TableIdentifier tableId =
- TableIdentifier.of("default", "table" +
Long.toString(UUID.randomUUID().hashCode(), 16));
+ TableIdentifier.of(
+ "default", "table.with.dots" +
Long.toString(UUID.randomUUID().hashCode(), 16));
Table simpleTable = warehouse.createTable(tableId,
schemaForMode(TestFixtures.SCHEMA, 1));
final Schema schema = icebergSchemaToBeamSchema(simpleTable.schema());
diff --git
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOWriteTest.java
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOWriteTest.java
index 52d92911f4e..6580f4eeaf6 100644
---
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOWriteTest.java
+++
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOWriteTest.java
@@ -148,7 +148,8 @@ public class IcebergIOWriteTest implements Serializable {
@Test
public void testSimpleAppend() throws Exception {
TableIdentifier tableId =
- TableIdentifier.of("default", "table" +
Long.toString(UUID.randomUUID().hashCode(), 16));
+ TableIdentifier.of(
+ "default", "table.with.dots" +
Long.toString(UUID.randomUUID().hashCode(), 16));
Map<String, String> catalogProps =
ImmutableMap.<String, String>builder()
diff --git
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergUtilsTest.java
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergUtilsTest.java
index c9026522dba..3da31ecc206 100644
---
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergUtilsTest.java
+++
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergUtilsTest.java
@@ -19,11 +19,15 @@ package org.apache.beam.sdk.io.iceberg;
import static org.apache.beam.sdk.io.iceberg.IcebergUtils.TypeAndMaxId;
import static
org.apache.beam.sdk.io.iceberg.IcebergUtils.beamFieldTypeToIcebergFieldType;
+import static org.apache.beam.sdk.io.iceberg.IcebergUtils.parseTableIdentifier;
+import static
org.apache.beam.sdk.io.iceberg.IcebergUtils.tableIdentifierToString;
import static org.apache.iceberg.types.Types.NestedField.optional;
import static org.apache.iceberg.types.Types.NestedField.required;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.hamcrest.Matchers.equalTo;
+import static org.junit.Assert.assertArrayEquals;
import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertThrows;
import static org.junit.Assert.assertTrue;
import java.math.BigDecimal;
@@ -44,6 +48,7 @@ import
org.apache.beam.sdk.schemas.logicaltypes.VariableString;
import org.apache.beam.sdk.values.Row;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
+import org.apache.iceberg.catalog.TableIdentifier;
import org.apache.iceberg.data.GenericRecord;
import org.apache.iceberg.data.Record;
import org.apache.iceberg.types.Type;
@@ -60,6 +65,65 @@ import org.junit.runners.JUnit4;
@RunWith(Enclosed.class)
public class IcebergUtilsTest {
+ @RunWith(JUnit4.class)
+ public static class TableIdentifierTests {
+
+ @Test
+ public void parseTableIdentifierParsesJson() {
+ TableIdentifier identifier =
+ parseTableIdentifier(
+ " {\"namespace\": [\"dogs\", \"owners.and.handlers\"],"
+ + " \"name\": \"food.with.dots\"}");
+
+ assertArrayEquals(
+ new String[] {"dogs", "owners.and.handlers"},
identifier.namespace().levels());
+ assertEquals("food.with.dots", identifier.name());
+ }
+
+ @Test
+ public void parseTableIdentifierParsesLegacyDottedString() {
+ TableIdentifier identifier =
parseTableIdentifier("dogs.owners.and.handlers.food");
+
+ assertArrayEquals(
+ new String[] {"dogs", "owners", "and", "handlers"},
identifier.namespace().levels());
+ assertEquals("food", identifier.name());
+ }
+
+ @Test
+ public void tableIdentifierToStringRoundTripsSpecialCharacters() {
+ TableIdentifier expected =
+ TableIdentifier.of("dogs", "owners.and.handlers", "food.with.dots");
+
+ assertEquals(expected,
parseTableIdentifier(tableIdentifierToString(expected)));
+ }
+
+ @Test
+ public void tableIdentifierToStringUsesLegacyFormWhenUnambiguous() {
+ assertEquals("dogs.food",
tableIdentifierToString(TableIdentifier.of("dogs", "food")));
+ }
+
+ @Test
+ public void
tableIdentifierToStringUsesJsonForLegacyStringsThatLookLikeJson() {
+ TableIdentifier expected = TableIdentifier.of("{dogs}", "{food}");
+
+ assertEquals(expected,
parseTableIdentifier(tableIdentifierToString(expected)));
+ }
+
+ @Test
+ public void
tableIdentifierToStringDoesNotUseJsonForPartialJsonLikeStrings() {
+ TableIdentifier expected = TableIdentifier.of("{dogs}", "food");
+
+ assertEquals("{dogs}.food", tableIdentifierToString(expected));
+ assertEquals(expected,
parseTableIdentifier(tableIdentifierToString(expected)));
+ }
+
+ @Test
+ public void parseTableIdentifierRejectsInvalidJsonIdentifier() {
+ assertThrows(
+ IllegalArgumentException.class, () ->
parseTableIdentifier("{\"table_name\":\"food\"}"));
+ }
+ }
+
@RunWith(JUnit4.class)
public static class RowToRecordTests {
/**