This is an automated email from the ASF dual-hosted git repository.

jrmccluskey pushed a commit to branch release-2.77
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/release-2.77 by this push:
     new 12d09df300e [Cherry-pick][#39723] Iceberg Side Input Cache Public API 
and Wrap-Up  (#40248)
12d09df300e is described below

commit 12d09df300e6e0941ba1ad6f7bdbf0172add08b0
Author: Jack McCluskey <[email protected]>
AuthorDate: Wed Sep 23 16:58:52 2026 -0400

    [Cherry-pick][#39723] Iceberg Side Input Cache Public API and Wrap-Up  
(#40248)
    
    * Expose Public API for Iceberg Side Input Cache (#40140)
    
    * Expose Public API for Iceberg Side Input Cache
    
    * Update 
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergIO.java
    
    Co-authored-by: Ahmed Abualsaud 
<[email protected]>
    
    * Streamline API, remove redundant checks
    
    * Check polled table metric
    
    * Evolve Iceberg schema mid-execution in streaming side input cache test
    
    * Throw if sub-options are enabled without side input cache enabled
    
    * test with partition spec evolution
    
    ---------
    
    Co-authored-by: Ahmed Abualsaud 
<[email protected]>
    
    * Address wrap-up items for Iceberg side input table cache (#40222)
    
    * Address wrap-up items for Iceberg side input table cache
    
    - Bind historical sort orders unchecked via 
SortOrderParser.fromJson(schema, node, getOrderId()) in SerializableTableSpec.
    - Expose Hadoop Configuration via 
IcebergCatalogConfig.getHadoopConfiguration() and propagate to 
SerializableTableSpec.getFileIO(conf) through SideInputTable in write 
transforms.
    - Add worker-local transient Guava cache to ExtractTableIdsDoFn to 
eliminate hot-key shuffle bottlenecks on streaming writes.
    - Add comprehensive unit and integration tests.
    
    * Propagate Clock to ExtractTableIdsDoFn for deterministic cache expiration 
in tests
    
    - Ensure TableMetadataDriver passes its Clock to ExtractTableIdsDoFn so 
virtual-time advancement in tests (e.g. 
testUnusedTablesEvictedFromStreamingCache) correctly expires worker-local table 
ID caches instead of using wall-clock time.
    - Fixes testUnusedTablesEvictedFromStreamingCache flakiness on CI 
environments with limited concurrency.
    
    * Make streaming refresh and spec evolution tests robust against worker 
cache TTL
    
    - Advance ControllableTestClock on trigger rows in 
testMetadataRefreshedAcrossIntervals, 
testMetadataRefreshedAcrossIntervalsAsSideInput, and 
testMetadataRefreshedAcrossIntervalsAsSideInputWithMultipleTables so 
worker-local table ID caches expire as intended across virtual intervals.
    - In testStreamingSpecEvolutionWithoutPipelineRestart, wait for the 
worker-local cache TTL (500ms) to elapse after spec evolution before emitting 
the second element.
    
    * Declare jackson-databind dependency in Iceberg IO module
    
    - Add library.java.jackson_databind to dependencies in 
sdks/java/io/iceberg/build.gradle to resolve usedUndeclaredArtifacts dependency 
analysis warning caused by JsonNode usage in SerializableTableSpec.
    
    * Avoid potential flakes in EvolveSpec DoFn, address nits
    
    * Simplify worker-local cache in TableMetadataDriver using LinkedHashMap LRU
    
    - Replace Guava Cache, CacheBuilder, and Ticker in ExtractTableIdsDoFn with 
a standard LinkedHashMap-based LRU cache capped at 10,000 entries.
    - Track lastEmitted timestamp per table identifier and emit when (now - 
lastEmitted) >= refreshInterval / 2.
    - Remove getClock() and setClock() from TableMetadataDriver and its Builder 
to keep the AutoValue API clean and free of testing methods.
    - Provide testing hooks via static globalTestClock and instance setClock on 
ExtractTableIdsDoFn.
    
    * simplify dofn + fix unit test
    
    ---------
    
    Co-authored-by: Ahmed Abualsaud 
<[email protected]>
---
 sdks/java/io/iceberg/build.gradle                  |   1 +
 .../beam/sdk/io/iceberg/IcebergCatalogConfig.java  |  24 +-
 .../org/apache/beam/sdk/io/iceberg/IcebergIO.java  | 125 ++++-
 .../IcebergWriteSchemaTransformProvider.java       |  57 ++
 .../beam/sdk/io/iceberg/RecordWriterManager.java   |   5 +-
 .../beam/sdk/io/iceberg/SerializableTableSpec.java |  33 +-
 .../apache/beam/sdk/io/iceberg/SideInputTable.java |  33 +-
 .../beam/sdk/io/iceberg/TableMetadataDriver.java   |  85 ++-
 .../io/iceberg/WritePartitionedRowsToFiles.java    |   5 +-
 .../sdk/io/iceberg/IcebergCatalogConfigTest.java   |  56 ++
 .../iceberg/IcebergIOSideInputTableCacheTest.java  | 621 +++++++++++++++++++++
 .../IcebergWriteSchemaTransformProviderTest.java   | 185 ++++++
 .../sdk/io/iceberg/SerializableTableSpecTest.java  |  77 +++
 .../beam/sdk/io/iceberg/SideInputTableTest.java    |  48 ++
 .../sdk/io/iceberg/TableMetadataDriverTest.java    | 227 ++++++--
 15 files changed, 1515 insertions(+), 67 deletions(-)

diff --git a/sdks/java/io/iceberg/build.gradle 
b/sdks/java/io/iceberg/build.gradle
index b3228004ed3..7278ca56159 100644
--- a/sdks/java/io/iceberg/build.gradle
+++ b/sdks/java/io/iceberg/build.gradle
@@ -50,6 +50,7 @@ dependencies {
     implementation library.java.avro
     implementation library.java.slf4j_api
     implementation library.java.joda_time
+    implementation library.java.jackson_databind
     implementation "org.apache.parquet:parquet-column:$parquet_version"
     implementation "org.apache.parquet:parquet-hadoop:$parquet_version"
     implementation "org.apache.parquet:parquet-common:$parquet_version"
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 8fe05cb5b21..fd836dc30dd 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
@@ -80,6 +80,21 @@ public abstract class IcebergCatalogConfig implements 
Serializable {
     return CATALOG_CACHE.computeIfAbsent(this, 
IcebergCatalogConfig::buildCatalog);
   }
 
+  /**
+   * Constructs and returns a new {@link Configuration} populated with 
properties from {@link
+   * #getConfigProperties()}.
+   */
+  public Configuration getHadoopConfiguration() {
+    Configuration config = new Configuration();
+    Map<String, String> confProps = getConfigProperties();
+    if (confProps != null) {
+      for (Map.Entry<String, String> prop : confProps.entrySet()) {
+        config.set(prop.getKey(), prop.getValue());
+      }
+    }
+    return config;
+  }
+
   private static Catalog buildCatalog(IcebergCatalogConfig catalogConfig) {
     String catalogName = catalogConfig.getCatalogName();
     if (catalogName == null) {
@@ -89,14 +104,7 @@ public abstract class IcebergCatalogConfig implements 
Serializable {
     if (catalogProps == null) {
       catalogProps = Maps.newHashMap();
     }
-    Map<String, String> confProps = catalogConfig.getConfigProperties();
-    if (confProps == null) {
-      confProps = Maps.newHashMap();
-    }
-    Configuration config = new Configuration();
-    for (Map.Entry<String, String> prop : confProps.entrySet()) {
-      config.set(prop.getKey(), prop.getValue());
-    }
+    Configuration config = catalogConfig.getHadoopConfiguration();
     return CatalogUtil.buildIcebergCatalog(catalogName, catalogProps, config);
   }
 
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergIO.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergIO.java
index d2119a241e3..bbcaef7a0a6 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergIO.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergIO.java
@@ -30,8 +30,10 @@ import org.apache.beam.sdk.io.iceberg.cdc.sink.WriteCdcRows;
 import org.apache.beam.sdk.options.StreamingOptions;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.transforms.PTransform;
+import org.apache.beam.sdk.transforms.display.DisplayData;
 import org.apache.beam.sdk.values.PBegin;
 import org.apache.beam.sdk.values.PCollection;
+import org.apache.beam.sdk.values.PCollectionView;
 import org.apache.beam.sdk.values.Row;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Predicates;
@@ -391,6 +393,7 @@ public class IcebergIO {
         .setCatalogConfig(catalog)
         .setDistributionMode(DistributionMode.NONE)
         .setAutoSharding(false)
+        .setUsingSideInputTableCache(false)
         .build();
   }
 
@@ -434,6 +437,14 @@ public class IcebergIO {
 
     abstract @Nullable List<String> getSortFields();
 
+    abstract boolean getUsingSideInputTableCache();
+
+    abstract @Nullable Integer getMaximumCacheSize();
+
+    abstract @Nullable Duration getTableRefreshInterval();
+
+    abstract @Nullable Integer getPollingBuckets();
+
     abstract Builder toBuilder();
 
     @AutoValue.Builder
@@ -458,6 +469,14 @@ public class IcebergIO {
 
       abstract Builder setSortFields(List<String> sortFields);
 
+      abstract Builder setUsingSideInputTableCache(boolean 
usingSideInputTableCache);
+
+      abstract Builder setMaximumCacheSize(@Nullable Integer maximumCacheSize);
+
+      abstract Builder setTableRefreshInterval(@Nullable Duration 
refreshInterval);
+
+      abstract Builder setPollingBuckets(@Nullable Integer pollingBuckets);
+
       abstract WriteRows build();
     }
 
@@ -494,11 +513,11 @@ public class IcebergIO {
      * Defines distribution of write data. Supported distributions:
      *
      * <ol>
-     *   <li>{@link DistributionMode.NONE}: don't shuffle rows (default)
-     *   <li>{@link DistributionMode.HASH}: shuffle rows by partition key 
before writing data
+     *   <li>{@link DistributionMode#NONE}: don't shuffle rows (default)
+     *   <li>{@link DistributionMode#HASH}: shuffle rows by partition key 
before writing data
      * </ol>
      *
-     * {@link DistributionMode.RANGE} is not supported yet
+     * {@link DistributionMode#RANGE} is not supported yet
      */
     public WriteRows withDistributionMode(DistributionMode mode) {
       return toBuilder().setDistributionMode(mode).build();
@@ -542,6 +561,68 @@ public class IcebergIO {
       return toBuilder().setSortFields(sortFields).build();
     }
 
+    /**
+     * Enables expirable side-input caching of Iceberg table metadata across 
workers.
+     *
+     * <p>When enabled, a driver transform periodically polls the Iceberg 
catalog and broadcasts
+     * lightweight table specifications as a side input. Workers construct 
in-memory {@link Table}
+     * representations without issuing remote catalog RPCs, drastically 
reducing catalog load.
+     */
+    public WriteRows withSideInputTableCache() {
+      return toBuilder().setUsingSideInputTableCache(true).build();
+    }
+
+    /**
+     * Sets the maximum number of distinct table metadata specifications to 
broadcast in the
+     * side-input cache. Any tables exceeding this limit fall back to 
worker-local catalog loading.
+     *
+     * <p><b>Note:</b> This option is only supported for bounded (batch) 
pipelines. Calling this on
+     * an unbounded streaming pipeline will throw an exception at pipeline 
construction.
+     */
+    public WriteRows withMaximumCacheSize(int maximumCacheSize) {
+      Preconditions.checkArgument(maximumCacheSize > 0, "maximumCacheSize must 
be greater than 0");
+      return toBuilder().setMaximumCacheSize(maximumCacheSize).build();
+    }
+
+    /**
+     * Sets the interval at which table metadata is refreshed from the Iceberg 
catalog.
+     *
+     * <p>Applicable for unbounded streaming pipelines. Defaults to 5 minutes.
+     */
+    public WriteRows withTableRefreshInterval(Duration refreshInterval) {
+      Preconditions.checkNotNull(refreshInterval, "refreshInterval must not be 
null");
+      Preconditions.checkArgument(
+          refreshInterval.isLongerThan(Duration.ZERO), "refreshInterval must 
be greater than 0");
+      return toBuilder().setTableRefreshInterval(refreshInterval).build();
+    }
+
+    /**
+     * Sets the number of parallel buckets/workers used to query the Iceberg 
catalog during
+     * refreshes. Defaults to 1 to serialize catalog queries and protect 
catalogs from connection
+     * spikes.
+     */
+    public WriteRows withPollingBuckets(int pollingBuckets) {
+      Preconditions.checkArgument(pollingBuckets > 0, "pollingBuckets must be 
greater than 0");
+      return toBuilder().setPollingBuckets(pollingBuckets).build();
+    }
+
+    @Override
+    public void populateDisplayData(DisplayData.Builder builder) {
+      super.populateDisplayData(builder);
+      builder.add(
+          DisplayData.item("usingSideInputTableCache", 
getUsingSideInputTableCache())
+              .withLabel("Using Side-Input Table Cache"));
+      builder.addIfNotNull(
+          DisplayData.item("maximumCacheSize", getMaximumCacheSize())
+              .withLabel("Maximum Cache Size"));
+      builder.addIfNotNull(
+          DisplayData.item("tableRefreshInterval", getTableRefreshInterval())
+              .withLabel("Table Refresh Interval"));
+      builder.addIfNotNull(
+          DisplayData.item("pollingBuckets", getPollingBuckets())
+              .withLabel("Catalog Polling Buckets"));
+    }
+
     @Override
     public IcebergWriteResult expand(PCollection<Row> input) {
       List<?> allToArgs = Arrays.asList(getTableIdentifier(), 
getDynamicDestinations());
@@ -567,6 +648,35 @@ public class IcebergIO {
             "Must only provide direct write limit for unbounded pipelines.");
       }
 
+      boolean hasSideInputOptions =
+          getMaximumCacheSize() != null
+              || getTableRefreshInterval() != null
+              || getPollingBuckets() != null;
+      Preconditions.checkArgument(
+          getUsingSideInputTableCache() || !hasSideInputOptions,
+          "Cannot specify side-input cache sub-options (maximumCacheSize, "
+              + "tableRefreshInterval, pollingBuckets) without enabling 
side-input table cache via withSideInputTableCache().");
+
+      PCollectionView<Map<String, SerializableTableSpec>> metadataView = null;
+      if (getUsingSideInputTableCache()) {
+        TableMetadataDriver.Builder driverBuilder =
+            TableMetadataDriver.builder()
+                .setCatalogConfig(getCatalogConfig())
+                .setDynamicDestinations(destinations);
+
+        if (getMaximumCacheSize() != null) {
+          driverBuilder.setMaximumCacheSize(getMaximumCacheSize());
+        }
+        if (getTableRefreshInterval() != null) {
+          driverBuilder.setRefreshInterval(getTableRefreshInterval());
+        }
+        if (getPollingBuckets() != null) {
+          driverBuilder.setPollingBuckets(getPollingBuckets());
+        }
+
+        metadataView = input.apply("GenerateTableMetadataView", 
driverBuilder.build().asView());
+      }
+
       switch (getDistributionMode()) {
         case NONE:
           Preconditions.checkArgument(
@@ -581,12 +691,14 @@ public class IcebergIO {
                       destinations,
                       getTriggeringFrequency(),
                       getDirectWriteByteLimit(),
-                      getWriteProperties()));
+                      getWriteProperties(),
+                      metadataView));
         case HASH:
           return input
               .apply(
                   "AssignDestinationAndPartition",
-                  new AssignDestinationsAndPartitions(destinations, 
getCatalogConfig()))
+                  new AssignDestinationsAndPartitions(
+                      destinations, getCatalogConfig(), metadataView))
               .apply(
                   "Write Rows to Partitions",
                   new WriteToPartitions(
@@ -594,7 +706,8 @@ public class IcebergIO {
                       destinations,
                       getTriggeringFrequency(),
                       getAutoSharding(),
-                      getWriteProperties()));
+                      getWriteProperties(),
+                      metadataView));
         default:
           throw new UnsupportedOperationException(
               "Unsupported distribution mode: " + getDistributionMode());
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java
index 0ae9d5fb0ec..e5a81e41a8b 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java
@@ -163,6 +163,23 @@ public class IcebergWriteSchemaTransformProvider
             + "'write.parquet.bloom-filter-enabled.column.<col>').")
     public abstract @Nullable Map<String, String> getWriteProperties();
 
+    @SchemaFieldDescription(
+        "Enables expirable side-input caching of Iceberg table metadata across 
workers to reduce catalog load.")
+    public abstract @Nullable Boolean getUsingSideInputTableCache();
+
+    @SchemaFieldDescription(
+        "For a streaming pipeline, sets the interval in seconds at which table 
metadata is refreshed from the catalog.")
+    public abstract @Nullable Integer getTableRefreshIntervalSeconds();
+
+    @SchemaFieldDescription(
+        "For a batch pipeline, sets the maximum number of table metadata specs 
to cache in memory. "
+            + "Tables exceeding this limit fall back to worker-local catalog 
loading.")
+    public abstract @Nullable Integer getMaximumCacheSize();
+
+    @SchemaFieldDescription(
+        "Sets the number of parallel buckets/workers used to query the Iceberg 
catalog during refreshes. Defaults to 1.")
+    public abstract @Nullable Integer getPollingBuckets();
+
     @AutoValue.Builder
     public abstract static class Builder {
       public abstract Builder setTable(String table);
@@ -195,6 +212,14 @@ public class IcebergWriteSchemaTransformProvider
 
       public abstract Builder setWriteProperties(Map<String, String> 
writeProperties);
 
+      public abstract Builder setUsingSideInputTableCache(Boolean 
usingSideInputTableCache);
+
+      public abstract Builder setTableRefreshIntervalSeconds(Integer 
tableRefreshIntervalSeconds);
+
+      public abstract Builder setMaximumCacheSize(Integer maximumCacheSize);
+
+      public abstract Builder setPollingBuckets(Integer pollingBuckets);
+
       public abstract Configuration build();
     }
 
@@ -291,6 +316,38 @@ public class IcebergWriteSchemaTransformProvider
         writeTransform = writeTransform.withWriteProperties(writeProperties);
       }
 
+      boolean hasSideInputOptions =
+          configuration.getTableRefreshIntervalSeconds() != null
+              || configuration.getMaximumCacheSize() != null
+              || configuration.getPollingBuckets() != null;
+
+      if (!Boolean.TRUE.equals(configuration.getUsingSideInputTableCache())
+          && hasSideInputOptions) {
+        throw new IllegalArgumentException(
+            "Cannot specify side-input cache sub-options 
(table_refresh_interval_seconds, "
+                + "maximum_cache_size, polling_buckets) without explicitly 
setting using_side_input_table_cache to true.");
+      }
+
+      boolean enableSideInputCache =
+          Boolean.TRUE.equals(configuration.getUsingSideInputTableCache());
+
+      if (enableSideInputCache) {
+        writeTransform = writeTransform.withSideInputTableCache();
+        @Nullable Integer refreshSec = 
configuration.getTableRefreshIntervalSeconds();
+        if (refreshSec != null) {
+          writeTransform =
+              
writeTransform.withTableRefreshInterval(Duration.standardSeconds(refreshSec));
+        }
+        @Nullable Integer maxCacheSize = configuration.getMaximumCacheSize();
+        if (maxCacheSize != null) {
+          writeTransform = writeTransform.withMaximumCacheSize(maxCacheSize);
+        }
+        @Nullable Integer pollingBuckets = configuration.getPollingBuckets();
+        if (pollingBuckets != null) {
+          writeTransform = writeTransform.withPollingBuckets(pollingBuckets);
+        }
+      }
+
       // TODO: support dynamic destinations
       IcebergWriteResult result = rows.apply(writeTransform);
 
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriterManager.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriterManager.java
index 25e5a13da43..e7e79c49553 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriterManager.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriterManager.java
@@ -27,7 +27,6 @@ import java.time.ZoneOffset;
 import java.time.format.DateTimeFormatter;
 import java.time.temporal.ChronoUnit;
 import java.util.ArrayList;
-import java.util.Collections;
 import java.util.List;
 import java.util.Map;
 import java.util.UUID;
@@ -318,9 +317,7 @@ class RecordWriterManager implements AutoCloseable {
     if (sideInputTableSpecs != null && 
sideInputTableSpecs.containsKey(tableIdString)) {
       SerializableTableSpec spec = sideInputTableSpecs.get(tableIdString);
       if (spec != null) {
-        Map<String, String> catalogProperties = 
catalogConfig.getCatalogProperties();
-        return new SideInputTable(
-            spec, catalogProperties != null ? catalogProperties : 
Collections.emptyMap());
+        return new SideInputTable(spec, catalogConfig);
       }
     }
     return TableCache.getAndRefreshIfStale(
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SerializableTableSpec.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SerializableTableSpec.java
index e89a341f55c..ea61b496cc1 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SerializableTableSpec.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SerializableTableSpec.java
@@ -19,7 +19,9 @@ package org.apache.beam.sdk.io.iceberg;
 
 import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
 
+import com.fasterxml.jackson.databind.JsonNode;
 import com.google.auto.value.AutoValue;
+import java.io.IOException;
 import java.io.Serializable;
 import java.util.Collections;
 import java.util.List;
@@ -33,6 +35,8 @@ import org.apache.beam.sdk.schemas.annotations.DefaultSchema;
 import org.apache.beam.sdk.schemas.annotations.SchemaFieldNumber;
 import org.apache.beam.sdk.schemas.annotations.SchemaIgnore;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
+import org.apache.hadoop.conf.Configurable;
+import org.apache.hadoop.conf.Configuration;
 import org.apache.iceberg.EncryptedKeyParser;
 import org.apache.iceberg.HasTableOperations;
 import org.apache.iceberg.PartitionSpec;
@@ -47,6 +51,7 @@ import org.apache.iceberg.catalog.TableIdentifier;
 import org.apache.iceberg.encryption.EncryptedKey;
 import org.apache.iceberg.io.FileIO;
 import org.apache.iceberg.io.FileIOParser;
+import org.apache.iceberg.util.JsonUtil;
 import org.checkerframework.checker.nullness.qual.MonotonicNonNull;
 import org.checkerframework.checker.nullness.qual.Nullable;
 
@@ -185,7 +190,14 @@ public abstract class SerializableTableSpec implements 
Serializable {
         if (local == null) {
           ImmutableMap.Builder<Integer, SortOrder> builder = 
ImmutableMap.builder();
           for (Map.Entry<Integer, String> entry : 
getSortOrdersJson().entrySet()) {
-            builder.put(entry.getKey(), SortOrderParser.fromJson(getSchema(), 
entry.getValue()));
+            try {
+              JsonNode node = JsonUtil.mapper().readTree(entry.getValue());
+              builder.put(
+                  entry.getKey(), SortOrderParser.fromJson(getSchema(), node, 
getOrderId()));
+            } catch (IOException e) {
+              throw new IllegalArgumentException(
+                  "Failed to parse sort order JSON for orderId " + 
entry.getKey(), e);
+            }
           }
           cachedSortOrders = local = builder.build();
         }
@@ -224,17 +236,34 @@ public abstract class SerializableTableSpec implements 
Serializable {
     return local;
   }
 
+  /** Returns a cached {@link FileIO} instance for this table using default 
configuration. */
   @SchemaIgnore
   public FileIO getFileIO() {
+    return getFileIO(null);
+  }
+
+  /**
+   * Returns a cached {@link FileIO} instance for this table, configured with 
the provided Hadoop
+   * {@link Configuration} if supported.
+   */
+  @SchemaIgnore
+  public FileIO getFileIO(@Nullable Configuration conf) {
     FileIO local = cachedFileIO;
     if (local == null) {
       synchronized (this) {
         local = cachedFileIO;
         if (local == null) {
-          cachedFileIO = local = FileIOParser.fromJson(getFileIoJson());
+          cachedFileIO =
+              local =
+                  conf != null
+                      ? FileIOParser.fromJson(getFileIoJson(), conf)
+                      : FileIOParser.fromJson(getFileIoJson());
         }
       }
     }
+    if (conf != null && local instanceof Configurable) {
+      ((Configurable) local).setConf(conf);
+    }
     return local;
   }
 
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SideInputTable.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SideInputTable.java
index aa51571c4cb..c5a4ddc3d82 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SideInputTable.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SideInputTable.java
@@ -25,6 +25,7 @@ import java.util.Map;
 import java.util.Objects;
 import org.apache.beam.sdk.annotations.Internal;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects;
+import org.apache.hadoop.conf.Configuration;
 import org.apache.iceberg.AppendFiles;
 import org.apache.iceberg.DeleteFiles;
 import org.apache.iceberg.ExpireSnapshots;
@@ -61,6 +62,7 @@ import org.apache.iceberg.encryption.KeyManagementClient;
 import org.apache.iceberg.encryption.PlaintextEncryptionManager;
 import org.apache.iceberg.io.FileIO;
 import org.apache.iceberg.io.LocationProvider;
+import org.checkerframework.checker.nullness.qual.Nullable;
 
 /**
  * A lightweight adapter that implements {@link Table} backed by a {@link 
SerializableTableSpec}.
@@ -80,14 +82,23 @@ public class SideInputTable implements Table {
   private final SerializableTableSpec spec;
   private final EncryptionManager encryptionManager;
   private final LocationProvider locationProvider;
+  private final @Nullable Configuration hadoopConf;
 
   public SideInputTable(SerializableTableSpec spec) {
-    this(spec, Collections.emptyMap());
+    this(spec, Collections.emptyMap(), null);
   }
 
   public SideInputTable(SerializableTableSpec spec, Map<String, String> 
catalogProperties) {
+    this(spec, catalogProperties, null);
+  }
+
+  public SideInputTable(
+      SerializableTableSpec spec,
+      Map<String, String> catalogProperties,
+      @Nullable Configuration hadoopConf) {
     this.spec = checkNotNull(spec, "spec must not be null");
     checkNotNull(catalogProperties, "catalogProperties must not be null");
+    this.hadoopConf = hadoopConf;
     this.locationProvider =
         LocationProviders.locationsFor(spec.getLocation(), 
spec.getProperties());
 
@@ -102,10 +113,27 @@ public class SideInputTable implements Table {
   }
 
   public SideInputTable(SerializableTableSpec spec, EncryptionManager 
encryptionManager) {
+    this(spec, encryptionManager, null);
+  }
+
+  public SideInputTable(
+      SerializableTableSpec spec,
+      EncryptionManager encryptionManager,
+      @Nullable Configuration hadoopConf) {
     this.spec = checkNotNull(spec, "spec must not be null");
     this.encryptionManager = checkNotNull(encryptionManager, 
"encryptionManager must not be null");
     this.locationProvider =
         LocationProviders.locationsFor(spec.getLocation(), 
spec.getProperties());
+    this.hadoopConf = hadoopConf;
+  }
+
+  public SideInputTable(SerializableTableSpec spec, IcebergCatalogConfig 
catalogConfig) {
+    this(
+        spec,
+        checkNotNull(catalogConfig, "catalogConfig must not be 
null").getCatalogProperties() != null
+            ? catalogConfig.getCatalogProperties()
+            : Collections.emptyMap(),
+        catalogConfig.getHadoopConfiguration());
   }
 
   public SerializableTableSpec getTableSpec() {
@@ -164,7 +192,7 @@ public class SideInputTable implements Table {
 
   @Override
   public FileIO io() {
-    return spec.getFileIO();
+    return spec.getFileIO(hadoopConf);
   }
 
   @Override
@@ -363,6 +391,7 @@ public class SideInputTable implements Table {
     return MoreObjects.toStringHelper(this)
         .add("spec", spec)
         .add("encryptionManager", encryptionManager)
+        .add("hadoopConf", hadoopConf)
         .toString();
   }
 }
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/TableMetadataDriver.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/TableMetadataDriver.java
index cf5d6310a1d..d7442e3d23b 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/TableMetadataDriver.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/TableMetadataDriver.java
@@ -24,6 +24,7 @@ import java.io.Serializable;
 import java.util.ArrayList;
 import java.util.Collections;
 import java.util.HashMap;
+import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Map;
 import org.apache.beam.sdk.annotations.Internal;
@@ -252,19 +253,24 @@ public abstract class TableMetadataDriver
 
   @Override
   public PCollection<KV<String, @Nullable SerializableTableSpec>> 
expand(PCollection<Row> input) {
+    boolean isStreaming = input.isBounded() == PCollection.IsBounded.UNBOUNDED;
+
+    Duration customInterval = getRefreshInterval();
+    Duration interval =
+        checkNotNull(customInterval != null ? customInterval : 
DEFAULT_REFRESH_INTERVAL);
+
     PCollection<String> tableIds =
         input
-            .apply("ExtractTableIds", ParDo.of(new 
ExtractTableIdsDoFn(getDynamicDestinations())))
+            .apply(
+                "ExtractTableIds",
+                ParDo.of(
+                    new ExtractTableIdsDoFn(
+                        getDynamicDestinations(), isStreaming ? interval : 
null)))
             .setCoder(StringUtf8Coder.of())
             .apply("MetadataGlobalWindow", Window.into(new GlobalWindows()));
 
-    boolean isStreaming = input.isBounded() == PCollection.IsBounded.UNBOUNDED;
-
     PCollection<String> distinctTableIds;
     if (isStreaming) {
-      Duration customInterval = getRefreshInterval();
-      Duration interval =
-          checkNotNull(customInterval != null ? customInterval : 
DEFAULT_REFRESH_INTERVAL);
       distinctTableIds =
           tableIds.apply(
               "DeduplicateTableIds", 
Deduplicate.<String>values().withDuration(interval));
@@ -276,7 +282,7 @@ public abstract class TableMetadataDriver
     Integer maxCacheSize = getMaximumCacheSize();
     if (maxCacheSize != null) {
       if (isStreaming) {
-        throw new UnsupportedOperationException(
+        throw new IllegalArgumentException(
             "maximumCacheSize is currently not supported for unbounded 
streaming pipelines.");
       }
       cachedTableIds = distinctTableIds.apply("CapCacheSize", 
Sample.any(maxCacheSize));
@@ -323,10 +329,61 @@ public abstract class TableMetadataDriver
   }
 
   static class ExtractTableIdsDoFn extends DoFn<Row, String> {
+    private static final int DEFAULT_LOCAL_CACHE_MAX_SIZE = 10_000;
+
+    private static volatile @Nullable Clock globalTestClock;
+
     private final DynamicDestinations dynamicDestinations;
+    private final @Nullable Duration refreshInterval;
+    private transient @Nullable Clock clock;
+    private transient @Nullable LinkedHashMap<String, Long> lastEmittedCache;
 
     ExtractTableIdsDoFn(DynamicDestinations dynamicDestinations) {
+      this(dynamicDestinations, null);
+    }
+
+    ExtractTableIdsDoFn(
+        DynamicDestinations dynamicDestinations, @Nullable Duration 
refreshInterval) {
       this.dynamicDestinations = dynamicDestinations;
+      this.refreshInterval = refreshInterval;
+    }
+
+    @VisibleForTesting
+    void setClock(@Nullable Clock clock) {
+      this.clock = clock;
+    }
+
+    @VisibleForTesting
+    static void setGlobalTestClock(@Nullable Clock clock) {
+      globalTestClock = clock;
+    }
+
+    @Setup
+    public void setup() {
+      initCache();
+    }
+
+    private void initCache() {
+      if (lastEmittedCache == null && refreshInterval != null) {
+        this.lastEmittedCache =
+            new LinkedHashMap<String, Long>(16, 0.75f, true) {
+              @Override
+              protected boolean removeEldestEntry(Map.Entry<String, Long> 
eldest) {
+                return size() > DEFAULT_LOCAL_CACHE_MAX_SIZE;
+              }
+            };
+      }
+    }
+
+    private long getNow() {
+      if (clock != null) {
+        return clock.currentTimeMillis();
+      }
+      Clock global = globalTestClock;
+      if (global != null) {
+        return global.currentTimeMillis();
+      }
+      return System.currentTimeMillis();
     }
 
     @ProcessElement
@@ -340,7 +397,19 @@ public abstract class TableMetadataDriver
           dynamicDestinations.getTableStringIdentifier(
               ValueInSingleWindow.of(element, timestamp, window, paneInfo));
       if (tableIdentifier != null && !tableIdentifier.trim().isEmpty()) {
-        out.output(tableIdentifier.trim());
+        Map<String, Long> cache = lastEmittedCache;
+        Duration interval = refreshInterval;
+        if (cache != null && interval != null) {
+          long now = getNow();
+          Long lastEmitted = cache.get(tableIdentifier);
+          long minInterval = Math.max(1L, interval.getMillis() / 2);
+          if (lastEmitted == null || (now - lastEmitted) >= minInterval) {
+            cache.put(tableIdentifier, now);
+            out.output(tableIdentifier);
+          }
+        } else {
+          out.output(tableIdentifier);
+        }
       }
     }
   }
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WritePartitionedRowsToFiles.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WritePartitionedRowsToFiles.java
index 881d2577fad..fbb6cb08aca 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WritePartitionedRowsToFiles.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WritePartitionedRowsToFiles.java
@@ -22,7 +22,6 @@ import static 
org.apache.beam.sdk.io.iceberg.AssignDestinationsAndPartitions.PAR
 import static 
org.apache.beam.sdk.io.iceberg.RecordWriterManager.getPartitionDataPath;
 import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
 
-import java.util.Collections;
 import java.util.Map;
 import java.util.UUID;
 import org.apache.beam.sdk.coders.IterableCoder;
@@ -223,9 +222,7 @@ class WritePartitionedRowsToFiles
       if (sideInputTableSpecs != null && 
sideInputTableSpecs.containsKey(tableIdString)) {
         SerializableTableSpec spec = sideInputTableSpecs.get(tableIdString);
         if (spec != null) {
-          Map<String, String> catalogProperties = 
catalogConfig.getCatalogProperties();
-          return new SideInputTable(
-              spec, catalogProperties != null ? catalogProperties : 
Collections.emptyMap());
+          return new SideInputTable(spec, catalogConfig);
         }
       }
       return TableCache.getAndRefreshIfStale(
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergCatalogConfigTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergCatalogConfigTest.java
new file mode 100644
index 00000000000..4e0e9c95b50
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergCatalogConfigTest.java
@@ -0,0 +1,56 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.iceberg;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertNull;
+
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
+import org.apache.hadoop.conf.Configuration;
+import org.junit.Test;
+
+public class IcebergCatalogConfigTest {
+
+  @Test
+  public void testGetHadoopConfigurationWhenPropertiesNull() {
+    IcebergCatalogConfig config =
+        IcebergCatalogConfig.builder().setCatalogName("test_catalog").build();
+
+    Configuration hadoopConf = config.getHadoopConfiguration();
+    assertNotNull(hadoopConf);
+    assertNull(hadoopConf.get("non.existent.key"));
+  }
+
+  @Test
+  public void testGetHadoopConfigurationPopulatesProperties() {
+    IcebergCatalogConfig config =
+        IcebergCatalogConfig.builder()
+            .setCatalogName("test_catalog")
+            .setConfigProperties(
+                ImmutableMap.of(
+                    "fs.defaultFS", "file:///test/path",
+                    "custom.hadoop.key", "custom-hadoop-val"))
+            .build();
+
+    Configuration hadoopConf = config.getHadoopConfiguration();
+    assertNotNull(hadoopConf);
+    assertEquals("file:///test/path", hadoopConf.get("fs.defaultFS"));
+    assertEquals("custom-hadoop-val", hadoopConf.get("custom.hadoop.key"));
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOSideInputTableCacheTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOSideInputTableCacheTest.java
new file mode 100644
index 00000000000..289fabc6991
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOSideInputTableCacheTest.java
@@ -0,0 +1,621 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.iceberg;
+
+import static java.util.Arrays.asList;
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertThrows;
+import static org.junit.Assert.assertTrue;
+
+import java.io.Serializable;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.UUID;
+import org.apache.beam.sdk.Pipeline;
+import org.apache.beam.sdk.PipelineResult;
+import org.apache.beam.sdk.metrics.MetricNameFilter;
+import org.apache.beam.sdk.metrics.MetricQueryResults;
+import org.apache.beam.sdk.metrics.MetricResult;
+import org.apache.beam.sdk.metrics.MetricsFilter;
+import org.apache.beam.sdk.schemas.Schema;
+import org.apache.beam.sdk.testing.TestPipeline;
+import org.apache.beam.sdk.testing.TestStream;
+import org.apache.beam.sdk.transforms.Create;
+import org.apache.beam.sdk.transforms.DoFn;
+import org.apache.beam.sdk.transforms.ParDo;
+import org.apache.beam.sdk.transforms.display.DisplayData;
+import org.apache.beam.sdk.values.PCollection;
+import org.apache.beam.sdk.values.Row;
+import org.apache.beam.sdk.values.ValueInSingleWindow;
+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.CatalogUtil;
+import org.apache.iceberg.DataFile;
+import org.apache.iceberg.DistributionMode;
+import org.apache.iceberg.FileFormat;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Snapshot;
+import org.apache.iceberg.SnapshotChanges;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.catalog.TableIdentifier;
+import org.apache.iceberg.data.IcebergGenerics;
+import org.apache.iceberg.data.Record;
+import org.apache.iceberg.types.Types;
+import org.hamcrest.Matchers;
+import org.joda.time.Duration;
+import org.junit.Before;
+import org.junit.ClassRule;
+import org.junit.Rule;
+import org.junit.Test;
+import org.junit.rules.TemporaryFolder;
+import org.junit.runner.RunWith;
+import org.junit.runners.Parameterized;
+
+/** Tests for {@link IcebergIO.WriteRows} with side-input table caching 
enabled. */
+@RunWith(Parameterized.class)
+public class IcebergIOSideInputTableCacheTest implements Serializable {
+
+  private static final String NONE = "none";
+  private static final String HASH = "hash";
+  private static final String HASH_WITH_AUTOSHARDING = "hashWithAutoSharding";
+
+  @Parameterized.Parameters(name = "distributionMode={0}")
+  public static Iterable<Object[]> data() {
+    return asList(new Object[][] {{NONE}, {HASH}, {HASH_WITH_AUTOSHARDING}});
+  }
+
+  @Parameterized.Parameter(0)
+  public String distributionMode;
+
+  @ClassRule public static final TemporaryFolder TEMPORARY_FOLDER = new 
TemporaryFolder();
+
+  @Rule
+  public transient TestDataWarehouse warehouse =
+      new TestDataWarehouse(TEMPORARY_FOLDER, "default_side_input");
+
+  @Rule public transient TestPipeline testPipeline = TestPipeline.create();
+
+  private IcebergCatalogConfig catalogConfig;
+
+  @Before
+  public void setUp() {
+    Map<String, String> catalogProps =
+        ImmutableMap.<String, String>builder()
+            .put("type", CatalogUtil.ICEBERG_CATALOG_TYPE_HADOOP)
+            .put("warehouse", warehouse.location)
+            .build();
+
+    catalogConfig =
+        IcebergCatalogConfig.builder()
+            .setCatalogName("hadoop")
+            .setCatalogProperties(catalogProps)
+            .build();
+    TableCache.invalidateAll();
+  }
+
+  private IcebergIO.WriteRows applyDistribution(IcebergIO.WriteRows write) {
+    if (distributionMode.contains(HASH)) {
+      write = write.withDistributionMode(DistributionMode.HASH);
+    }
+    if (distributionMode.equals(HASH_WITH_AUTOSHARDING)) {
+      write = write.withAutosharding();
+    }
+    return write;
+  }
+
+  private long getTablesPolledCount(PipelineResult result) {
+    MetricQueryResults metrics =
+        result
+            .metrics()
+            .queryMetrics(
+                MetricsFilter.builder()
+                    .addNameFilter(
+                        MetricNameFilter.named(TableMetadataDriver.class, 
"tablesPolled"))
+                    .build());
+    long total = 0;
+    for (MetricResult<Long> counter : metrics.getCounters()) {
+      Long val = counter.getCommitted() != null ? counter.getCommitted() : 
counter.getAttempted();
+      if (val != null) {
+        total += val;
+      }
+    }
+    return total;
+  }
+
+  @Test
+  public void testBatchSingleTableWithSideInputCache() throws Exception {
+    TableIdentifier tableId =
+        TableIdentifier.of(
+            "default_side_input",
+            "single_table_" + Long.toString(UUID.randomUUID().hashCode(), 16));
+
+    warehouse.createTable(tableId, TestFixtures.SCHEMA);
+
+    PCollection<Row> input =
+        testPipeline
+            .apply("CreateRecords", 
Create.of(TestFixtures.asRows(TestFixtures.FILE1SNAPSHOT1)))
+            
.setRowSchema(IcebergUtils.icebergSchemaToBeamSchema(TestFixtures.SCHEMA));
+
+    IcebergIO.WriteRows write =
+        IcebergIO.writeRows(catalogConfig)
+            .to(tableId)
+            .withSideInputTableCache()
+            .withPollingBuckets(1);
+
+    input.apply("WriteToTable", applyDistribution(write));
+    PipelineResult result = testPipeline.run();
+    result.waitUntilFinish();
+
+    assertEquals(1L, getTablesPolledCount(result));
+
+    Table table = warehouse.loadTable(tableId);
+    List<Record> writtenRecords = 
ImmutableList.copyOf(IcebergGenerics.read(table).build());
+    assertThat(writtenRecords, 
Matchers.containsInAnyOrder(TestFixtures.FILE1SNAPSHOT1.toArray()));
+    assertNotNull(table.currentSnapshot());
+  }
+
+  @Test
+  public void testBatchDynamicDestinationsWithSideInputCache() throws 
Exception {
+    final String salt = Long.toString(UUID.randomUUID().hashCode(), 16);
+    final TableIdentifier table1Id = TableIdentifier.of("default_side_input", 
"dyn_table1_" + salt);
+    final TableIdentifier table2Id = TableIdentifier.of("default_side_input", 
"dyn_table2_" + salt);
+    final TableIdentifier table3Id = TableIdentifier.of("default_side_input", 
"dyn_table3_" + salt);
+
+    warehouse.createTable(table1Id, TestFixtures.SCHEMA);
+    warehouse.createTable(table2Id, TestFixtures.SCHEMA);
+    warehouse.createTable(table3Id, TestFixtures.SCHEMA);
+
+    Schema beamSchema = 
IcebergUtils.icebergSchemaToBeamSchema(TestFixtures.SCHEMA);
+    Schema inputSchema =
+        
Schema.builder().addStringField("dest").addFields(beamSchema.getFields()).build();
+
+    List<Row> rows = new ArrayList<>();
+    for (Record record : TestFixtures.FILE1SNAPSHOT1) {
+      Row beamRow = IcebergUtils.icebergRecordToBeamRow(beamSchema, record);
+      rows.add(
+          Row.withSchema(inputSchema)
+              .addValue(IcebergUtils.tableIdentifierToString(table1Id))
+              .addValues(beamRow.getValues())
+              .build());
+    }
+    for (Record record : TestFixtures.FILE2SNAPSHOT1) {
+      Row beamRow = IcebergUtils.icebergRecordToBeamRow(beamSchema, record);
+      rows.add(
+          Row.withSchema(inputSchema)
+              .addValue(IcebergUtils.tableIdentifierToString(table2Id))
+              .addValues(beamRow.getValues())
+              .build());
+    }
+    for (Record record : TestFixtures.FILE3SNAPSHOT1) {
+      Row beamRow = IcebergUtils.icebergRecordToBeamRow(beamSchema, record);
+      rows.add(
+          Row.withSchema(inputSchema)
+              .addValue(IcebergUtils.tableIdentifierToString(table3Id))
+              .addValues(beamRow.getValues())
+              .build());
+    }
+
+    DynamicDestinations dynamicDestinations =
+        new DynamicDestinations() {
+          @Override
+          public Schema getDataSchema() {
+            return beamSchema;
+          }
+
+          @Override
+          public Row getData(Row element) {
+            Row.Builder builder = Row.withSchema(beamSchema);
+            for (Schema.Field field : beamSchema.getFields()) {
+              builder.addValue(element.getValue(field.getName()));
+            }
+            return builder.build();
+          }
+
+          @Override
+          public IcebergDestination instantiateDestination(String destination) 
{
+            return IcebergDestination.builder()
+                
.setTableIdentifier(IcebergUtils.parseTableIdentifier(destination))
+                .setFileFormat(FileFormat.PARQUET)
+                .build();
+          }
+
+          @Override
+          public String getTableStringIdentifier(ValueInSingleWindow<Row> 
element) {
+            return element.getValue().getString("dest");
+          }
+        };
+
+    PCollection<Row> input =
+        testPipeline.apply("CreateRows", 
Create.of(rows)).setRowSchema(inputSchema);
+
+    IcebergIO.WriteRows write =
+        IcebergIO.writeRows(catalogConfig)
+            .to(dynamicDestinations)
+            .withSideInputTableCache()
+            .withPollingBuckets(1);
+
+    input.apply("WriteDynamic", applyDistribution(write));
+    PipelineResult result = testPipeline.run();
+    result.waitUntilFinish();
+
+    assertEquals(3L, getTablesPolledCount(result));
+
+    Table table1 = warehouse.loadTable(table1Id);
+    Table table2 = warehouse.loadTable(table2Id);
+    Table table3 = warehouse.loadTable(table3Id);
+
+    List<Record> records1 = 
ImmutableList.copyOf(IcebergGenerics.read(table1).build());
+    List<Record> records2 = 
ImmutableList.copyOf(IcebergGenerics.read(table2).build());
+    List<Record> records3 = 
ImmutableList.copyOf(IcebergGenerics.read(table3).build());
+
+    assertThat(records1, 
Matchers.containsInAnyOrder(TestFixtures.FILE1SNAPSHOT1.toArray()));
+    assertThat(records2, 
Matchers.containsInAnyOrder(TestFixtures.FILE2SNAPSHOT1.toArray()));
+    assertThat(records3, 
Matchers.containsInAnyOrder(TestFixtures.FILE3SNAPSHOT1.toArray()));
+  }
+
+  @Test
+  public void testBatchMaximumCacheSizeSamplingAndFallback() throws Exception {
+    final String salt = Long.toString(UUID.randomUUID().hashCode(), 16);
+    final TableIdentifier table1Id = TableIdentifier.of("default_side_input", 
"sample_t1_" + salt);
+    final TableIdentifier table2Id = TableIdentifier.of("default_side_input", 
"sample_t2_" + salt);
+    final TableIdentifier table3Id = TableIdentifier.of("default_side_input", 
"sample_t3_" + salt);
+    final TableIdentifier table4Id = TableIdentifier.of("default_side_input", 
"sample_t4_" + salt);
+
+    warehouse.createTable(table1Id, TestFixtures.SCHEMA);
+    warehouse.createTable(table2Id, TestFixtures.SCHEMA);
+    warehouse.createTable(table3Id, TestFixtures.SCHEMA);
+    warehouse.createTable(table4Id, TestFixtures.SCHEMA);
+
+    Schema beamSchema = 
IcebergUtils.icebergSchemaToBeamSchema(TestFixtures.SCHEMA);
+    Schema inputSchema =
+        
Schema.builder().addStringField("dest").addFields(beamSchema.getFields()).build();
+
+    List<Row> rows = new ArrayList<>();
+    for (Record record : TestFixtures.FILE1SNAPSHOT1) {
+      Row beamRow = IcebergUtils.icebergRecordToBeamRow(beamSchema, record);
+      rows.add(
+          Row.withSchema(inputSchema)
+              .addValue(IcebergUtils.tableIdentifierToString(table1Id))
+              .addValues(beamRow.getValues())
+              .build());
+      rows.add(
+          Row.withSchema(inputSchema)
+              .addValue(IcebergUtils.tableIdentifierToString(table2Id))
+              .addValues(beamRow.getValues())
+              .build());
+      rows.add(
+          Row.withSchema(inputSchema)
+              .addValue(IcebergUtils.tableIdentifierToString(table3Id))
+              .addValues(beamRow.getValues())
+              .build());
+      rows.add(
+          Row.withSchema(inputSchema)
+              .addValue(IcebergUtils.tableIdentifierToString(table4Id))
+              .addValues(beamRow.getValues())
+              .build());
+    }
+
+    DynamicDestinations dynamicDestinations =
+        new DynamicDestinations() {
+          @Override
+          public Schema getDataSchema() {
+            return beamSchema;
+          }
+
+          @Override
+          public Row getData(Row element) {
+            Row.Builder builder = Row.withSchema(beamSchema);
+            for (Schema.Field field : beamSchema.getFields()) {
+              builder.addValue(element.getValue(field.getName()));
+            }
+            return builder.build();
+          }
+
+          @Override
+          public IcebergDestination instantiateDestination(String destination) 
{
+            return IcebergDestination.builder()
+                
.setTableIdentifier(IcebergUtils.parseTableIdentifier(destination))
+                .setFileFormat(FileFormat.PARQUET)
+                .build();
+          }
+
+          @Override
+          public String getTableStringIdentifier(ValueInSingleWindow<Row> 
element) {
+            return element.getValue().getString("dest");
+          }
+        };
+
+    PCollection<Row> input =
+        testPipeline.apply("CreateSampleRows", 
Create.of(rows)).setRowSchema(inputSchema);
+
+    // Maximum cache size of 2 forces 2 tables to be cached and 2 to fall back 
to TableCache
+    IcebergIO.WriteRows write =
+        IcebergIO.writeRows(catalogConfig)
+            .to(dynamicDestinations)
+            .withSideInputTableCache()
+            .withMaximumCacheSize(2)
+            .withPollingBuckets(1);
+
+    input.apply("WriteWithSampleCap", applyDistribution(write));
+    PipelineResult result = testPipeline.run();
+    result.waitUntilFinish();
+
+    assertEquals(2L, getTablesPolledCount(result));
+
+    // Verify all 4 tables received data successfully
+    for (TableIdentifier tId : asList(table1Id, table2Id, table3Id, table4Id)) 
{
+      Table table = warehouse.loadTable(tId);
+      List<Record> written = 
ImmutableList.copyOf(IcebergGenerics.read(table).build());
+      assertEquals(TestFixtures.FILE1SNAPSHOT1.size(), written.size());
+    }
+  }
+
+  @Test
+  public void testStreamingWithSideInputCache() throws Exception {
+    TableIdentifier tableId =
+        TableIdentifier.of(
+            "default_side_input",
+            "streaming_table_" + Long.toString(UUID.randomUUID().hashCode(), 
16));
+
+    warehouse.createTable(tableId, TestFixtures.SCHEMA);
+
+    Schema beamSchema = 
IcebergUtils.icebergSchemaToBeamSchema(TestFixtures.SCHEMA);
+    List<Row> rows1 = TestFixtures.asRows(TestFixtures.FILE1SNAPSHOT1);
+    List<Row> rows2 = TestFixtures.asRows(TestFixtures.FILE2SNAPSHOT1);
+
+    TestStream<Row> testStream =
+        TestStream.create(beamSchema)
+            .addElements(rows1.get(0), rows1.subList(1, 
rows1.size()).toArray(new Row[0]))
+            .advanceProcessingTime(Duration.standardSeconds(2))
+            .addElements(rows2.get(0), rows2.subList(1, 
rows2.size()).toArray(new Row[0]))
+            .advanceProcessingTime(Duration.standardSeconds(2))
+            .advanceWatermarkToInfinity();
+
+    PCollection<Row> input = testPipeline.apply("StreamingInput", testStream);
+
+    IcebergIO.WriteRows write =
+        IcebergIO.writeRows(catalogConfig)
+            .to(tableId)
+            .withSideInputTableCache()
+            .withTriggeringFrequency(Duration.standardSeconds(1))
+            .withTableRefreshInterval(Duration.standardSeconds(2))
+            .withPollingBuckets(1);
+
+    input.apply("StreamingWrite", applyDistribution(write));
+    PipelineResult result = testPipeline.run();
+    result.waitUntilFinish();
+
+    assertThat(getTablesPolledCount(result), 
Matchers.greaterThanOrEqualTo(1L));
+
+    Table table = warehouse.loadTable(tableId);
+    List<Record> written = 
ImmutableList.copyOf(IcebergGenerics.read(table).build());
+    List<Record> expected = new ArrayList<>();
+    expected.addAll(TestFixtures.FILE1SNAPSHOT1);
+    expected.addAll(TestFixtures.FILE2SNAPSHOT1);
+    assertThat(written, Matchers.containsInAnyOrder(expected.toArray()));
+    assertNotNull(table.currentSnapshot());
+  }
+
+  @Test
+  public void testStreamingSpecEvolutionWithoutPipelineRestart() throws 
Exception {
+    TableIdentifier tableId =
+        TableIdentifier.of(
+            "default_side_input", "spec_evolve_" + 
Long.toString(UUID.randomUUID().hashCode(), 16));
+
+    // Initial table schema: id (long), name (string), city (string)
+    org.apache.iceberg.Schema icebergSchema =
+        new org.apache.iceberg.Schema(
+            Types.NestedField.required(1, "id", Types.LongType.get()),
+            Types.NestedField.optional(2, "name", Types.StringType.get()),
+            Types.NestedField.optional(3, "city", Types.StringType.get()));
+
+    // Create table unpartitioned with format-version 2 to support spec 
evolution
+    Table realTable =
+        warehouse.createTable(
+            tableId,
+            icebergSchema,
+            PartitionSpec.unpartitioned(),
+            ImmutableMap.of("format-version", "2"));
+
+    Schema beamSchema =
+        Schema.builder()
+            .addInt64Field("id")
+            .addNullableStringField("name")
+            .addNullableStringField("city")
+            .build();
+
+    Row row1 = Row.withSchema(beamSchema).addValues(1L, "alice", "New 
York").build();
+    Row row2 = Row.withSchema(beamSchema).addValues(2L, "bob", "San 
Francisco").build();
+
+    TestStream<Row> testStream =
+        TestStream.create(beamSchema)
+            .addElements(row1)
+            .advanceProcessingTime(Duration.standardSeconds(2))
+            .addElements(row2)
+            .advanceProcessingTime(Duration.standardSeconds(2))
+            .advanceWatermarkToInfinity();
+
+    PCollection<Row> input =
+        testPipeline
+            .apply("StreamingEvolvedInput", testStream)
+            .apply(
+                "EvolveSpecMidExecution",
+                ParDo.of(new EvolveSpecMidExecutionDoFn(catalogConfig, 
tableId.toString())))
+            .setRowSchema(beamSchema);
+
+    IcebergIO.WriteRows write =
+        IcebergIO.writeRows(catalogConfig)
+            .to(tableId)
+            .withSideInputTableCache()
+            .withTriggeringFrequency(Duration.standardSeconds(1))
+            .withTableRefreshInterval(Duration.standardSeconds(1))
+            .withPollingBuckets(1);
+
+    input.apply("StreamingWriteEvolved", write);
+    PipelineResult result = testPipeline.run();
+    result.waitUntilFinish();
+
+    assertThat(getTablesPolledCount(result), 
Matchers.greaterThanOrEqualTo(1L));
+
+    realTable.refresh();
+    List<DataFile> addedFiles =
+        
ImmutableList.copyOf(SnapshotChanges.builderFor(realTable).build().addedDataFiles());
+    if (addedFiles.size() < 2) {
+      List<DataFile> allAddedFiles = new ArrayList<>();
+      for (Snapshot s : realTable.snapshots()) {
+        for (DataFile df :
+            
SnapshotChanges.builderFor(realTable).snapshot(s).build().addedDataFiles()) {
+          allAddedFiles.add(df);
+        }
+      }
+      addedFiles = allAddedFiles;
+    }
+
+    assertEquals(2, addedFiles.size());
+
+    DataFile firstFile = addedFiles.get(0);
+    DataFile secondFile = addedFiles.get(1);
+    DataFile unpartitionedFile;
+    DataFile partitionedFile;
+    if (realTable.specs().get(firstFile.specId()).isUnpartitioned()) {
+      unpartitionedFile = firstFile;
+      partitionedFile = secondFile;
+    } else {
+      unpartitionedFile = secondFile;
+      partitionedFile = firstFile;
+    }
+
+    PartitionSpec spec1 = realTable.specs().get(unpartitionedFile.specId());
+    assertNotNull(spec1);
+    assertTrue("First DataFile must have unpartitioned spec", 
spec1.isUnpartitioned());
+
+    PartitionSpec spec2 = realTable.specs().get(partitionedFile.specId());
+    assertNotNull(spec2);
+    assertEquals(1, spec2.fields().size());
+    assertEquals("city", spec2.fields().get(0).name());
+
+    List<Record> records = 
ImmutableList.copyOf(IcebergGenerics.read(realTable).build());
+    assertEquals(2, records.size());
+  }
+
+  @Test
+  public void testPreconditionsAndValidation() {
+    TableIdentifier tableId = TableIdentifier.of("default_side_input", 
"validation_table");
+    IcebergIO.WriteRows write = IcebergIO.writeRows(catalogConfig).to(tableId);
+
+    assertThrows(IllegalArgumentException.class, () -> 
write.withMaximumCacheSize(0));
+    assertThrows(IllegalArgumentException.class, () -> 
write.withMaximumCacheSize(-1));
+
+    assertThrows(
+        IllegalArgumentException.class, () -> 
write.withTableRefreshInterval(Duration.ZERO));
+
+    assertThrows(IllegalArgumentException.class, () -> 
write.withPollingBuckets(0));
+    assertThrows(IllegalArgumentException.class, () -> 
write.withPollingBuckets(-1));
+
+    // Unbounded streaming pipeline with maximumCacheSize must fail at expand
+    Pipeline p = Pipeline.create();
+    Schema schema = Schema.builder().addInt64Field("id").build();
+    TestStream<Row> testStream =
+        TestStream.create(schema)
+            .addElements(Row.withSchema(schema).addValues(1L).build())
+            .advanceWatermarkToInfinity();
+
+    PCollection<Row> streamInput = p.apply("StreamForValidation", testStream);
+
+    // Sub-options specified without withSideInputTableCache() must fail at 
expand
+    IcebergIO.WriteRows writeWithMaxCacheOnly =
+        IcebergIO.writeRows(catalogConfig).to(tableId).withMaximumCacheSize(5);
+    assertThrows(IllegalArgumentException.class, () -> 
streamInput.apply(writeWithMaxCacheOnly));
+
+    IcebergIO.WriteRows writeWithRefreshIntervalOnly =
+        IcebergIO.writeRows(catalogConfig)
+            .to(tableId)
+            .withTableRefreshInterval(Duration.standardMinutes(1));
+    assertThrows(
+        IllegalArgumentException.class, () -> 
streamInput.apply(writeWithRefreshIntervalOnly));
+
+    IcebergIO.WriteRows writeWithPollingBucketsOnly =
+        IcebergIO.writeRows(catalogConfig).to(tableId).withPollingBuckets(2);
+    assertThrows(
+        IllegalArgumentException.class, () -> 
streamInput.apply(writeWithPollingBucketsOnly));
+
+    IcebergIO.WriteRows streamWrite =
+        IcebergIO.writeRows(catalogConfig)
+            .to(tableId)
+            .withSideInputTableCache()
+            .withMaximumCacheSize(5);
+
+    assertThrows(IllegalArgumentException.class, () -> 
streamInput.apply(streamWrite));
+  }
+
+  @Test
+  public void testDisplayData() {
+    TableIdentifier tableId = TableIdentifier.of("default_side_input", 
"display_data_table");
+    IcebergIO.WriteRows write =
+        IcebergIO.writeRows(catalogConfig)
+            .to(tableId)
+            .withSideInputTableCache()
+            .withMaximumCacheSize(100)
+            .withTableRefreshInterval(Duration.standardMinutes(10))
+            .withPollingBuckets(3);
+
+    DisplayData displayData = DisplayData.from(write);
+    Map<String, String> items = new HashMap<>();
+    for (DisplayData.Item item : displayData.items()) {
+      items.put(item.getKey(), item.getValue() != null ? 
item.getValue().toString() : "");
+    }
+
+    assertEquals("true", items.get("usingSideInputTableCache"));
+    assertEquals("100", items.get("maximumCacheSize"));
+    assertEquals("600000", items.get("tableRefreshInterval"));
+    assertEquals("3", items.get("pollingBuckets"));
+  }
+
+  private static class EvolveSpecMidExecutionDoFn extends DoFn<Row, Row> {
+    private final IcebergCatalogConfig catalogConfig;
+    private final String tableIdString;
+
+    EvolveSpecMidExecutionDoFn(IcebergCatalogConfig catalogConfig, String 
tableIdString) {
+      this.catalogConfig = catalogConfig;
+      this.tableIdString = tableIdString;
+    }
+
+    @ProcessElement
+    public void processElement(@Element Row row, OutputReceiver<Row> out) {
+      Long id = row.getInt64("id");
+      if (id != null && id == 2L) {
+        Table table =
+            
catalogConfig.catalog().loadTable(IcebergUtils.parseTableIdentifier(tableIdString));
+        if (table.spec().isUnpartitioned()) {
+          table.updateSpec().addField("city").commit();
+          // Ensure worker-local table ID cache TTL (interval / 2 = 500ms) has 
elapsed
+          try {
+            Thread.sleep(700);
+          } catch (InterruptedException e) {
+            Thread.currentThread().interrupt();
+          }
+        }
+      }
+      out.output(row);
+    }
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProviderTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProviderTest.java
index bfb762f8e89..c4ac1e22cbf 100644
--- 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProviderTest.java
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProviderTest.java
@@ -25,6 +25,8 @@ import static 
org.apache.iceberg.util.DateTimeUtil.dateFromDays;
 import static org.apache.iceberg.util.DateTimeUtil.timestampFromMicros;
 import static org.hamcrest.MatcherAssert.assertThat;
 import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertThrows;
 import static org.junit.Assert.assertTrue;
 import static org.junit.Assume.assumeTrue;
 
@@ -117,6 +119,26 @@ public class IcebergWriteSchemaTransformProviderTest {
     new IcebergWriteSchemaTransformProvider().from(transformConfigRow);
   }
 
+  @Test
+  public void testBuildTransformWithRowAndSideInputCache() {
+    Map<String, String> properties = new HashMap<>();
+    properties.put("type", CatalogUtil.ICEBERG_CATALOG_TYPE_HADOOP);
+    properties.put("warehouse", "test_location");
+
+    Row transformConfigRow =
+        Row.withSchema(new 
IcebergWriteSchemaTransformProvider().configurationSchema())
+            .withFieldValue("table", "test_table_identifier")
+            .withFieldValue("catalog_name", "test-name")
+            .withFieldValue("catalog_properties", properties)
+            .withFieldValue("using_side_input_table_cache", true)
+            .withFieldValue("table_refresh_interval_seconds", 60)
+            .withFieldValue("maximum_cache_size", 100)
+            .withFieldValue("polling_buckets", 2)
+            .build();
+
+    new IcebergWriteSchemaTransformProvider().from(transformConfigRow);
+  }
+
   @Test
   public void testSimpleAppend() {
     String identifier = "default.table_" + 
Long.toString(UUID.randomUUID().hashCode(), 16);
@@ -159,6 +181,54 @@ public class IcebergWriteSchemaTransformProviderTest {
     assertThat(writtenRecords, 
Matchers.containsInAnyOrder(TestFixtures.FILE1SNAPSHOT1.toArray()));
   }
 
+  @Test
+  public void testSimpleAppendWithSideInputCache() {
+    String identifier =
+        "default.table_side_input_" + 
Long.toString(UUID.randomUUID().hashCode(), 16);
+
+    Map<String, String> properties = new HashMap<>();
+    properties.put("type", CatalogUtil.ICEBERG_CATALOG_TYPE_HADOOP);
+    properties.put("warehouse", warehouse.location);
+
+    Configuration config =
+        Configuration.builder()
+            .setTable(identifier)
+            .setCatalogName("name")
+            .setCatalogProperties(properties)
+            .setDistributionMode(distributionMode.name())
+            .setUsingSideInputTableCache(true)
+            .setTableRefreshIntervalSeconds(60)
+            .setPollingBuckets(1)
+            .build();
+
+    PCollectionRowTuple input =
+        PCollectionRowTuple.of(
+            INPUT_TAG,
+            testPipeline
+                .apply(
+                    "Records To Add", 
Create.of(TestFixtures.asRows(TestFixtures.FILE1SNAPSHOT1)))
+                
.setRowSchema(IcebergUtils.icebergSchemaToBeamSchema(TestFixtures.SCHEMA)));
+
+    PCollection<Row> result =
+        input
+            .apply(
+                "Append To Table With Cache",
+                new IcebergWriteSchemaTransformProvider().from(config))
+            .get(SNAPSHOTS_TAG);
+
+    PAssert.that(result)
+        .satisfies(new VerifyOutputs(Collections.singletonList(identifier), 
"append"));
+
+    testPipeline.run().waitUntilFinish();
+
+    TableIdentifier tableId = TableIdentifier.parse(identifier);
+    Table table = warehouse.loadTable(tableId);
+
+    List<Record> writtenRecords = 
ImmutableList.copyOf(IcebergGenerics.read(table).build());
+
+    assertThat(writtenRecords, 
Matchers.containsInAnyOrder(TestFixtures.FILE1SNAPSHOT1.toArray()));
+  }
+
   @Test
   public void testWriteUsingManagedTransform() {
     String identifier = "default.table_" + 
Long.toString(UUID.randomUUID().hashCode(), 16);
@@ -194,6 +264,121 @@ public class IcebergWriteSchemaTransformProviderTest {
     assertThat(writtenRecords, 
Matchers.containsInAnyOrder(TestFixtures.FILE1SNAPSHOT1.toArray()));
   }
 
+  @Test
+  public void testWriteUsingManagedTransformWithSideInputCache() {
+    String identifier =
+        "default.table_managed_cache_" + 
Long.toString(UUID.randomUUID().hashCode(), 16);
+
+    String yamlConfig =
+        String.format(
+            "table: %s\n"
+                + "catalog_name: test-name\n"
+                + "distribution_mode: %s\n"
+                + "using_side_input_table_cache: true\n"
+                + "table_refresh_interval_seconds: 60\n"
+                + "maximum_cache_size: 50\n"
+                + "polling_buckets: 1\n"
+                + "catalog_properties: \n"
+                + "  type: %s\n"
+                + "  warehouse: %s",
+            identifier,
+            distributionMode.name(),
+            CatalogUtil.ICEBERG_CATALOG_TYPE_HADOOP,
+            warehouse.location);
+    Map<String, Object> configMap = new Yaml().load(yamlConfig);
+
+    PCollection<Row> inputRows =
+        testPipeline
+            .apply("Records To Add", 
Create.of(TestFixtures.asRows(TestFixtures.FILE1SNAPSHOT1)))
+            
.setRowSchema(IcebergUtils.icebergSchemaToBeamSchema(TestFixtures.SCHEMA));
+
+    Managed.ManagedTransform writeTransform = 
Managed.write(Managed.ICEBERG).withConfig(configMap);
+    PCollectionRowTuple output = PCollectionRowTuple.of(INPUT_TAG, 
inputRows).apply(writeTransform);
+
+    PAssert.that(output.get(SNAPSHOTS_TAG))
+        .satisfies(new VerifyOutputs(Collections.singletonList(identifier), 
"append"));
+
+    testPipeline.run().waitUntilFinish();
+
+    Table table = warehouse.loadTable(TableIdentifier.parse(identifier));
+    List<Record> writtenRecords = 
ImmutableList.copyOf(IcebergGenerics.read(table).build());
+
+    assertThat(writtenRecords, 
Matchers.containsInAnyOrder(TestFixtures.FILE1SNAPSHOT1.toArray()));
+  }
+
+  @Test
+  public void testSideInputCacheSubOptionsRequireExplicitEnablement() {
+    IcebergWriteSchemaTransformProvider provider = new 
IcebergWriteSchemaTransformProvider();
+    Pipeline p = Pipeline.create();
+    PCollectionRowTuple dummyInput =
+        PCollectionRowTuple.of(
+            INPUT_TAG,
+            p.apply("DummyInput", 
Create.of(TestFixtures.asRows(TestFixtures.FILE1SNAPSHOT1)))
+                
.setRowSchema(IcebergUtils.icebergSchemaToBeamSchema(TestFixtures.SCHEMA)));
+
+    // Setting sub-options when using_side_input_table_cache is not set (null) 
must throw
+    // IllegalArgumentException
+    Configuration configWithMaxCacheSizeOnly =
+        Configuration.builder()
+            .setTable("default.table_max_cache")
+            .setCatalogName("name")
+            .setCatalogProperties(Collections.singletonMap("type", "hadoop"))
+            .setMaximumCacheSize(50)
+            .build();
+    assertThrows(
+        IllegalArgumentException.class,
+        () -> dummyInput.apply(provider.from(configWithMaxCacheSizeOnly)));
+
+    Configuration configWithRefreshIntervalOnly =
+        Configuration.builder()
+            .setTable("default.table_refresh")
+            .setCatalogName("name")
+            .setCatalogProperties(Collections.singletonMap("type", "hadoop"))
+            .setTableRefreshIntervalSeconds(60)
+            .build();
+    assertThrows(
+        IllegalArgumentException.class,
+        () -> dummyInput.apply(provider.from(configWithRefreshIntervalOnly)));
+
+    Configuration configWithPollingBucketsOnly =
+        Configuration.builder()
+            .setTable("default.table_buckets")
+            .setCatalogName("name")
+            .setCatalogProperties(Collections.singletonMap("type", "hadoop"))
+            .setPollingBuckets(2)
+            .build();
+    assertThrows(
+        IllegalArgumentException.class,
+        () -> dummyInput.apply(provider.from(configWithPollingBucketsOnly)));
+
+    // Setting using_side_input_table_cache to false while setting sub-options 
must throw
+    // IllegalArgumentException
+    Configuration invalidConfigWithFalse =
+        Configuration.builder()
+            .setTable("default.table_invalid")
+            .setCatalogName("name")
+            .setCatalogProperties(Collections.singletonMap("type", "hadoop"))
+            .setUsingSideInputTableCache(false)
+            .setMaximumCacheSize(50)
+            .build();
+    assertThrows(
+        IllegalArgumentException.class,
+        () -> dummyInput.apply(provider.from(invalidConfigWithFalse)));
+
+    // Explicitly setting using_side_input_table_cache to true with 
sub-options succeeds
+    Configuration validConfig =
+        Configuration.builder()
+            .setTable("default.table_valid")
+            .setCatalogName("name")
+            .setCatalogProperties(Collections.singletonMap("type", "hadoop"))
+            .setUsingSideInputTableCache(true)
+            .setMaximumCacheSize(50)
+            .setTableRefreshIntervalSeconds(60)
+            .setPollingBuckets(2)
+            .build();
+    assertNotNull(dummyInput.apply(provider.from(validConfig)));
+  }
+
   /**
    * @param operation if null, just perform a normal dynamic destination write 
test; otherwise,
    *     performs a simple filter on the record before writing. Valid options 
are "keep", "drop",
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SerializableTableSpecTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SerializableTableSpecTest.java
index a2f2abadfe4..9045b910c65 100644
--- 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SerializableTableSpecTest.java
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SerializableTableSpecTest.java
@@ -21,6 +21,7 @@ import static 
org.apache.iceberg.types.Types.NestedField.optional;
 import static org.apache.iceberg.types.Types.NestedField.required;
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertNull;
 import static org.junit.Assert.assertThrows;
 import static org.junit.Assert.assertTrue;
 import static org.mockito.Mockito.mock;
@@ -42,6 +43,7 @@ import java.util.concurrent.TimeUnit;
 import org.apache.beam.sdk.schemas.SchemaCoder;
 import org.apache.beam.sdk.util.CoderUtils;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
+import org.apache.hadoop.conf.Configurable;
 import org.apache.hadoop.conf.Configuration;
 import org.apache.iceberg.CatalogProperties;
 import org.apache.iceberg.CatalogUtil;
@@ -350,4 +352,79 @@ public class SerializableTableSpecTest {
       executor.shutdown();
     }
   }
+
+  @Test
+  public void testHistoricalSortOrderWithDroppedColumn() {
+    TableIdentifier tableId = TableIdentifier.of("default", 
"historical_sort_table");
+    Schema v1Schema =
+        new Schema(
+            required(1, "id", Types.LongType.get()),
+            optional(2, "name", Types.StringType.get()),
+            optional(3, "dropped_col", Types.StringType.get()));
+
+    SortOrder v1SortOrder =
+        SortOrder.builderFor(v1Schema)
+            .sortBy("dropped_col", SortDirection.ASC, NullOrder.NULLS_FIRST)
+            .build();
+
+    Table table = catalog.buildTable(tableId, 
v1Schema).withSortOrder(v1SortOrder).create();
+
+    int v1OrderId = table.sortOrder().orderId();
+
+    // First replace sort order with one referencing the remaining fields,
+    // making v1SortOrder a historical sort order
+    table.replaceSortOrder().asc("id").commit();
+    int v2OrderId = table.sortOrder().orderId();
+
+    // Then evolve schema by deleting the column that was part of the original 
sort order
+    table.updateSchema().deleteColumn("dropped_col").commit();
+
+    SerializableTableSpec spec = SerializableTableSpec.fromTable(table);
+
+    assertEquals(v2OrderId, spec.getOrderId());
+    assertEquals(table.sortOrder(), spec.getSortOrder());
+    assertEquals(1, spec.getSortOrder().fields().get(0).sourceId());
+
+    // Calling getSortOrders() should successfully bind historical sort orders
+    // without failing with ValidationException: Cannot find source column
+    Map<Integer, SortOrder> sortOrders = spec.getSortOrders();
+    assertNotNull(sortOrders);
+    assertTrue(sortOrders.containsKey(v1OrderId));
+    assertTrue(sortOrders.containsKey(v2OrderId));
+
+    SortOrder historicalOrder = spec.getSortOrder(v1OrderId);
+    assertNotNull(historicalOrder);
+    assertEquals(v1OrderId, historicalOrder.orderId());
+  }
+
+  @Test
+  public void testFileIOWithHadoopConfiguration() {
+    TableIdentifier tableId = TableIdentifier.of("default", 
"hadoop_conf_table");
+    Table table = catalog.createTable(tableId, TestFixtures.SCHEMA);
+
+    SerializableTableSpec spec = SerializableTableSpec.fromTable(table);
+
+    Configuration conf = new Configuration();
+    conf.set("custom.test.prop", "test-value-123");
+
+    FileIO fileIO = spec.getFileIO(conf);
+    assertNotNull(fileIO);
+    assertTrue("FileIO must be Configurable", fileIO instanceof Configurable);
+    assertEquals("test-value-123", ((Configurable) 
fileIO).getConf().get("custom.test.prop"));
+
+    // Also verify that zero-arg getFileIO() returns cached instance
+    assertEquals(fileIO, spec.getFileIO());
+
+    // Verify that calling zero-arg getFileIO() first does not pollute cache:
+    // subsequent getFileIO(conf) must update configuration on Configurable 
FileIO
+    SerializableTableSpec spec2 = SerializableTableSpec.fromTable(table);
+    FileIO fileIO2 = spec2.getFileIO();
+    assertNotNull(fileIO2);
+    assertNull(((Configurable) fileIO2).getConf().get("custom.test.prop"));
+
+    FileIO configuredFileIO2 = spec2.getFileIO(conf);
+    assertEquals(fileIO2, configuredFileIO2);
+    assertEquals(
+        "test-value-123", ((Configurable) 
configuredFileIO2).getConf().get("custom.test.prop"));
+  }
 }
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SideInputTableTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SideInputTableTest.java
index 663c818b587..6b1aa1843f0 100644
--- 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SideInputTableTest.java
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SideInputTableTest.java
@@ -25,6 +25,7 @@ import static org.junit.Assert.assertTrue;
 
 import java.util.Map;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
+import org.apache.hadoop.conf.Configurable;
 import org.apache.hadoop.conf.Configuration;
 import org.apache.iceberg.CatalogProperties;
 import org.apache.iceberg.CatalogUtil;
@@ -40,6 +41,7 @@ import org.apache.iceberg.catalog.TableIdentifier;
 import org.apache.iceberg.data.GenericRecord;
 import org.apache.iceberg.data.Record;
 import org.apache.iceberg.encryption.PlaintextEncryptionManager;
+import org.apache.iceberg.io.FileIO;
 import org.junit.Before;
 import org.junit.Rule;
 import org.junit.Test;
@@ -236,4 +238,50 @@ public class SideInputTableTest {
     assertNotNull(table1a.toString());
     assertTrue(table1a.toString().contains("SideInputTable"));
   }
+
+  @Test
+  public void testSideInputTableWithHadoopConfiguration() {
+    TableIdentifier tableId = TableIdentifier.of("default", 
"side_input_hadoop_conf_table");
+    Table realTable = catalog.createTable(tableId, TestFixtures.SCHEMA);
+    SerializableTableSpec spec = SerializableTableSpec.fromTable(tableId, 
realTable);
+
+    Configuration conf = new Configuration();
+    conf.set("custom.sideinput.prop", "custom-value-456");
+
+    // Test constructor accepting hadoopConf directly
+    SideInputTable tableWithConf = new SideInputTable(spec, ImmutableMap.of(), 
conf);
+    FileIO io = tableWithConf.io();
+    assertNotNull(io);
+    assertTrue("FileIO must implement Configurable", io instanceof 
Configurable);
+    assertEquals("custom-value-456", ((Configurable) 
io).getConf().get("custom.sideinput.prop"));
+
+    // Test constructor accepting IcebergCatalogConfig
+    IcebergCatalogConfig catalogConfig =
+        IcebergCatalogConfig.builder()
+            .setCatalogName("test_catalog")
+            .setConfigProperties(ImmutableMap.of("catalog.conf.prop", 
"val-789"))
+            .build();
+
+    SerializableTableSpec specFromConfig = 
SerializableTableSpec.fromTable(tableId, realTable);
+    SideInputTable tableFromCatalogConfig = new SideInputTable(specFromConfig, 
catalogConfig);
+    FileIO ioFromConfig = tableFromCatalogConfig.io();
+    assertNotNull(ioFromConfig);
+    assertTrue("FileIO must implement Configurable", ioFromConfig instanceof 
Configurable);
+    assertEquals("val-789", ((Configurable) 
ioFromConfig).getConf().get("catalog.conf.prop"));
+
+    // Test constructor accepting EncryptionManager and hadoopConf
+    SerializableTableSpec specWithEncryption = 
SerializableTableSpec.fromTable(tableId, realTable);
+    SideInputTable tableWithEncryptionAndConf =
+        new SideInputTable(specWithEncryption, 
PlaintextEncryptionManager.instance(), conf);
+    FileIO ioWithEncryption = tableWithEncryptionAndConf.io();
+    assertNotNull(ioWithEncryption);
+    assertTrue("FileIO must implement Configurable", ioWithEncryption 
instanceof Configurable);
+    assertEquals(
+        "custom-value-456",
+        ((Configurable) 
ioWithEncryption).getConf().get("custom.sideinput.prop"));
+
+    // Test constructor rejecting null catalogConfig
+    assertThrows(
+        NullPointerException.class, () -> new SideInputTable(spec, 
(IcebergCatalogConfig) null));
+  }
 }
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/TableMetadataDriverTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/TableMetadataDriverTest.java
index c46b155921b..48f86ff3239 100644
--- 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/TableMetadataDriverTest.java
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/TableMetadataDriverTest.java
@@ -18,6 +18,7 @@
 package org.apache.beam.sdk.io.iceberg;
 
 import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull;
+import static org.hamcrest.MatcherAssert.assertThat;
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertNotNull;
 import static org.junit.Assert.assertNull;
@@ -36,10 +37,13 @@ import org.apache.beam.sdk.testing.PAssert;
 import org.apache.beam.sdk.testing.TestPipeline;
 import org.apache.beam.sdk.testing.TestStream;
 import org.apache.beam.sdk.transforms.Create;
+import org.apache.beam.sdk.transforms.Deduplicate;
 import org.apache.beam.sdk.transforms.DoFn;
 import org.apache.beam.sdk.transforms.ParDo;
 import org.apache.beam.sdk.transforms.display.DisplayData;
 import org.apache.beam.sdk.transforms.windowing.FixedWindows;
+import org.apache.beam.sdk.transforms.windowing.GlobalWindow;
+import org.apache.beam.sdk.transforms.windowing.PaneInfo;
 import org.apache.beam.sdk.transforms.windowing.Window;
 import org.apache.beam.sdk.values.KV;
 import org.apache.beam.sdk.values.PCollection;
@@ -62,8 +66,10 @@ import org.apache.iceberg.data.GenericRecord;
 import org.apache.iceberg.data.Record;
 import org.apache.iceberg.types.Types;
 import org.checkerframework.checker.nullness.qual.Nullable;
+import org.hamcrest.Matchers;
 import org.joda.time.Duration;
 import org.joda.time.Instant;
+import org.junit.After;
 import org.junit.Before;
 import org.junit.Rule;
 import org.junit.Test;
@@ -118,6 +124,19 @@ public class TableMetadataDriverTest implements 
Serializable {
         }
       };
 
+  static class ControllableTestClock implements TableMetadataDriver.Clock {
+    private static final AtomicLong CURRENT_TIME = new AtomicLong(0L);
+
+    public static void setTime(long millis) {
+      CURRENT_TIME.set(millis);
+    }
+
+    @Override
+    public long currentTimeMillis() {
+      return CURRENT_TIME.get();
+    }
+  }
+
   @Before
   public void setUp() throws Exception {
     warehouseLocation = "file:" + tempFolder.newFolder().getAbsolutePath();
@@ -126,6 +145,13 @@ public class TableMetadataDriverTest implements 
Serializable {
             .setCatalogName("hadoop")
             .setCatalogProperties(ImmutableMap.of("type", "hadoop", 
"warehouse", warehouseLocation))
             .build();
+    ControllableTestClock.setTime(1000L);
+    TableMetadataDriver.ExtractTableIdsDoFn.setGlobalTestClock(new 
ControllableTestClock());
+  }
+
+  @After
+  public void tearDown() {
+    TableMetadataDriver.ExtractTableIdsDoFn.setGlobalTestClock(null);
   }
 
   private Catalog getCatalog() {
@@ -321,6 +347,9 @@ public class TableMetadataDriverTest implements 
Serializable {
     TableIdentifier tableId = TableIdentifier.of("default", "evolving_table");
     catalog.createTable(tableId, ICEBERG_SCHEMA);
 
+    Duration refreshInterval = Duration.standardSeconds(2);
+    ControllableTestClock.setTime(1000L);
+
     Row row1 =
         Row.withSchema(BEAM_SCHEMA).addValues(1L, "initial_data", 
"default.evolving_table").build();
     Row row2 =
@@ -347,6 +376,7 @@ public class TableMetadataDriverTest implements 
Serializable {
                       @ProcessElement
                       public void processElement(@Element Row row, 
OutputReceiver<Row> out) {
                         if ("trigger_update".equals(row.getString("data"))) {
+                          ControllableTestClock.setTime(5000L);
                           Table table =
                               catalogConfig
                                   .catalog()
@@ -367,7 +397,7 @@ public class TableMetadataDriverTest implements 
Serializable {
             TableMetadataDriver.builder()
                 .setCatalogConfig(catalogConfig)
                 .setDynamicDestinations(DYNAMIC_DESTINATIONS)
-                .setRefreshInterval(Duration.standardSeconds(2))
+                .setRefreshInterval(refreshInterval)
                 .build());
 
     // Downstream consumer transform verifying that updated metadata is 
received
@@ -398,6 +428,10 @@ public class TableMetadataDriverTest implements 
Serializable {
     catalog.createTable(tableId, ICEBERG_SCHEMA);
 
     String tableIdStr = "default.evolving_side_input_table";
+    Duration refreshInterval = Duration.standardSeconds(2);
+    ControllableTestClock.setTime(1000L);
+    ControllableTestClock testClock = new ControllableTestClock();
+
     Row row1 = Row.withSchema(BEAM_SCHEMA).addValues(1L, "initial_data", 
tableIdStr).build();
     Row row2 = Row.withSchema(BEAM_SCHEMA).addValues(2L, "trigger_update", 
tableIdStr).build();
     Row row3 = Row.withSchema(BEAM_SCHEMA).addValues(3L, "post_update_data", 
tableIdStr).build();
@@ -422,7 +456,9 @@ public class TableMetadataDriverTest implements 
Serializable {
                     new DoFn<Row, Row>() {
                       @ProcessElement
                       public void processElement(@Element Row row, 
OutputReceiver<Row> out) {
-                        if ("trigger_update".equals(row.getString("data"))) {
+                        String data = row.getString("data");
+                        if ("trigger_update".equals(data)) {
+                          ControllableTestClock.setTime(5000L);
                           Table table =
                               catalogConfig
                                   .catalog()
@@ -433,6 +469,8 @@ public class TableMetadataDriverTest implements 
Serializable {
                               .updateSchema()
                               .addColumn("new_col", Types.StringType.get())
                               .commit();
+                        } else if ("post_update_data".equals(data)) {
+                          ControllableTestClock.setTime(10000L);
                         }
                         out.output(row);
                       }
@@ -442,12 +480,13 @@ public class TableMetadataDriverTest implements 
Serializable {
     PCollectionView<Map<String, SerializableTableSpec>> metadataView =
         input.apply(
             "CreateMetadataView",
-            TableMetadataDriver.builder()
-                .setCatalogConfig(catalogConfig)
-                .setDynamicDestinations(DYNAMIC_DESTINATIONS)
-                .setRefreshInterval(Duration.standardSeconds(2))
-                .build()
-                .asView());
+            TableMetadataDriver.asView(
+                TableMetadataDriver.builder()
+                    .setCatalogConfig(catalogConfig)
+                    .setDynamicDestinations(DYNAMIC_DESTINATIONS)
+                    .setRefreshInterval(refreshInterval)
+                    .build(),
+                testClock));
 
     PCollection<String> consumerObserved =
         input.apply(
@@ -490,6 +529,9 @@ public class TableMetadataDriverTest implements 
Serializable {
 
     String tableAStr = "default.multi_table_a";
     String tableBStr = "default.multi_table_b";
+    Duration refreshInterval = Duration.standardSeconds(2);
+    ControllableTestClock.setTime(1000L);
+    ControllableTestClock testClock = new ControllableTestClock();
 
     Row rowSeedA = Row.withSchema(BEAM_SCHEMA).addValues(0L, "seed_a", 
tableAStr).build();
     Row rowSeedB = Row.withSchema(BEAM_SCHEMA).addValues(0L, "seed_b", 
tableBStr).build();
@@ -523,6 +565,7 @@ public class TableMetadataDriverTest implements 
Serializable {
                       @ProcessElement
                       public void processElement(@Element Row row, 
OutputReceiver<Row> out) {
                         if ("trigger_update_a".equals(row.getString("data"))) {
+                          ControllableTestClock.setTime(5000L);
                           Table table =
                               catalogConfig
                                   .catalog()
@@ -541,12 +584,13 @@ public class TableMetadataDriverTest implements 
Serializable {
     PCollectionView<Map<String, SerializableTableSpec>> metadataView =
         input.apply(
             "CreateMetadataView",
-            TableMetadataDriver.builder()
-                .setCatalogConfig(catalogConfig)
-                .setDynamicDestinations(DYNAMIC_DESTINATIONS)
-                .setRefreshInterval(Duration.standardSeconds(2))
-                .build()
-                .asView());
+            TableMetadataDriver.asView(
+                TableMetadataDriver.builder()
+                    .setCatalogConfig(catalogConfig)
+                    .setDynamicDestinations(DYNAMIC_DESTINATIONS)
+                    .setRefreshInterval(refreshInterval)
+                    .build(),
+                testClock));
 
     PCollection<String> consumerObserved =
         input.apply(
@@ -706,7 +750,7 @@ public class TableMetadataDriverTest implements 
Serializable {
   }
 
   @Test
-  public void 
testMaximumCacheSizeInStreamingThrowsUnsupportedOperationException() {
+  public void testMaximumCacheSizeInStreamingThrowsIllegalArgumentException() {
     pipeline.enableAbandonedNodeEnforcement(false);
     Row row = Row.withSchema(BEAM_SCHEMA).addValues(1L, "v1", 
"default.test_table").build();
     TestStream<Row> stream =
@@ -718,7 +762,7 @@ public class TableMetadataDriverTest implements 
Serializable {
     PCollection<Row> input = pipeline.apply("StreamInput", stream);
 
     assertThrows(
-        UnsupportedOperationException.class,
+        IllegalArgumentException.class,
         () ->
             input.apply(
                 TableMetadataDriver.builder()
@@ -884,10 +928,7 @@ public class TableMetadataDriverTest implements 
Serializable {
             Row.withSchema(BEAM_SCHEMA).addValues(1L, "v1", null).build(),
             Row.withSchema(BEAM_SCHEMA).addValues(2L, "v2", "").build(),
             Row.withSchema(BEAM_SCHEMA).addValues(3L, "v3", "   ").build(),
-            Row.withSchema(BEAM_SCHEMA).addValues(4L, "v4", 
"default.valid_dest_table").build(),
-            Row.withSchema(BEAM_SCHEMA)
-                .addValues(5L, "v5", "  default.valid_dest_table  ")
-                .build());
+            Row.withSchema(BEAM_SCHEMA).addValues(4L, "v4", 
"default.valid_dest_table").build());
 
     PCollection<Row> input = 
pipeline.apply(Create.of(rows)).setCoder(RowCoder.of(BEAM_SCHEMA));
 
@@ -1234,19 +1275,6 @@ public class TableMetadataDriverTest implements 
Serializable {
     assertEquals(mergedAB, mergedBA);
   }
 
-  static class ControllableTestClock implements TableMetadataDriver.Clock {
-    private static final AtomicLong CURRENT_TIME = new AtomicLong(0L);
-
-    public static void setTime(long millis) {
-      CURRENT_TIME.set(millis);
-    }
-
-    @Override
-    public long currentTimeMillis() {
-      return CURRENT_TIME.get();
-    }
-  }
-
   @Test
   public void testUnusedTablesEvictedFromStreamingCache() {
     TableIdentifier tableIdA = TableIdentifier.of("default", "evict_table_a");
@@ -1498,4 +1526,137 @@ public class TableMetadataDriverTest implements 
Serializable {
 
     pipeline.run();
   }
+
+  @Test
+  public void testExtractTableIdsWorkerLocalPreFiltering() {
+    TestStream.Builder<Row> streamBuilder = TestStream.create(BEAM_SCHEMA);
+    for (int i = 0; i < 1000; i++) {
+      streamBuilder =
+          streamBuilder.addElements(
+              Row.withSchema(BEAM_SCHEMA).addValues((long) i, "data", 
"default.table").build());
+    }
+    TestStream<Row> testStream = streamBuilder.advanceWatermarkToInfinity();
+
+    PCollection<String> tableIds =
+        pipeline
+            .apply(testStream)
+            .apply(
+                ParDo.of(
+                    new TableMetadataDriver.ExtractTableIdsDoFn(
+                        SINGLE_TABLE_DYNAMIC_DESTINATIONS, 
Duration.standardMinutes(5))));
+
+    // With worker-local pre-filtering, 1,000 rows emit at most 1 string per 
worker thread
+    // rather than 1,000 strings
+    PAssert.that(tableIds)
+        .satisfies(
+            actual -> {
+              List<String> list = ImmutableList.copyOf(actual);
+              for (String id : list) {
+                assertEquals("default.table", id);
+              }
+              assertThat(list.size(), Matchers.lessThanOrEqualTo(100));
+              return null;
+            });
+
+    PCollection<String> distinctIds =
+        
tableIds.apply(Deduplicate.<String>values().withDuration(Duration.standardMinutes(5)));
+    PAssert.that(distinctIds).containsInAnyOrder("default.table");
+    pipeline.run();
+  }
+
+  @Test
+  public void testExtractTableIdsMultipleTables() {
+    TestStream.Builder<Row> streamBuilder = TestStream.create(BEAM_SCHEMA);
+    for (int i = 0; i < 100; i++) {
+      streamBuilder =
+          streamBuilder.addElements(
+              Row.withSchema(BEAM_SCHEMA).addValues((long) i, "data", 
"default.table_a").build(),
+              Row.withSchema(BEAM_SCHEMA)
+                  .addValues((long) (i + 100), "data", "default.table_b")
+                  .build());
+    }
+    TestStream<Row> testStream = streamBuilder.advanceWatermarkToInfinity();
+
+    PCollection<String> tableIds =
+        pipeline
+            .apply(testStream)
+            .apply(
+                ParDo.of(
+                    new TableMetadataDriver.ExtractTableIdsDoFn(
+                        DYNAMIC_DESTINATIONS, Duration.standardMinutes(5))));
+
+    PCollection<String> distinctIds =
+        
tableIds.apply(Deduplicate.<String>values().withDuration(Duration.standardMinutes(5)));
+    PAssert.that(distinctIds).containsInAnyOrder("default.table_a", 
"default.table_b");
+    pipeline.run();
+  }
+
+  @Test
+  public void testExtractTableIdsCacheExpiration() {
+    TableMetadataDriver.ExtractTableIdsDoFn doFn =
+        new TableMetadataDriver.ExtractTableIdsDoFn(
+            SINGLE_TABLE_DYNAMIC_DESTINATIONS, Duration.standardMinutes(10));
+    ControllableTestClock testClock = new ControllableTestClock();
+    ControllableTestClock.setTime(1000L);
+    doFn.setClock(testClock);
+
+    List<String> outputs = new ArrayList<>();
+    DoFn.OutputReceiver<String> receiver =
+        new DoFn.OutputReceiver<String>() {
+          @Override
+          public void output(String output) {
+            outputs.add(output);
+          }
+
+          @Override
+          public void outputWithTimestamp(String output, Instant timestamp) {
+            outputs.add(output);
+          }
+
+          @Override
+          public org.apache.beam.sdk.values.OutputBuilder<String> 
builder(String output) {
+            throw new UnsupportedOperationException();
+          }
+        };
+
+    Row row1 = Row.withSchema(BEAM_SCHEMA).addValues(1L, "data", 
"default.table").build();
+    Row row2 = Row.withSchema(BEAM_SCHEMA).addValues(2L, "data", 
"default.table").build();
+    Row row3 = Row.withSchema(BEAM_SCHEMA).addValues(3L, "data", 
"default.table").build();
+
+    doFn.setup();
+    doFn.processElement(row1, GlobalWindow.INSTANCE, PaneInfo.NO_FIRING, 
Instant.now(), receiver);
+    assertEquals(1, outputs.size());
+    assertEquals("default.table", outputs.get(0));
+
+    // Second element should be suppressed by worker-local cache
+    doFn.processElement(row2, GlobalWindow.INSTANCE, PaneInfo.NO_FIRING, 
Instant.now(), receiver);
+    assertEquals(1, outputs.size());
+
+    // Advance clock beyond interval / 2 (5 minutes)
+    ControllableTestClock.setTime(1000L + 
Duration.standardMinutes(6).getMillis());
+
+    // Third element arrives after expiration -> should be emitted
+    doFn.processElement(row3, GlobalWindow.INSTANCE, PaneInfo.NO_FIRING, 
Instant.now(), receiver);
+    assertEquals(2, outputs.size());
+    assertEquals("default.table", outputs.get(1));
+  }
+
+  @Test
+  public void testExtractTableIdsIgnoresNullAndWhitespace() {
+    List<Row> rows = new ArrayList<>();
+    rows.add(Row.withSchema(BEAM_SCHEMA).addValues(1L, "data", (String) 
null).build());
+    rows.add(Row.withSchema(BEAM_SCHEMA).addValues(2L, "data", "").build());
+    rows.add(Row.withSchema(BEAM_SCHEMA).addValues(3L, "data", "   ").build());
+
+    PCollection<String> tableIds =
+        pipeline
+            .apply(Create.of(rows))
+            .apply(
+                ParDo.of(
+                    new TableMetadataDriver.ExtractTableIdsDoFn(
+                        DYNAMIC_DESTINATIONS, Duration.standardMinutes(5))));
+
+    PAssert.that(tableIds).empty();
+    pipeline.run();
+  }
 }

Reply via email to