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 {
     /**

Reply via email to