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();
+ }
}