This is an automated email from the ASF dual-hosted git repository.
ahmedabu98 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 2c2df5a66b6 Rename Iceberg table cache options for clarity (#40253)
2c2df5a66b6 is described below
commit 2c2df5a66b6b3d5f141b3dbfa553aac756c3ccb6
Author: Ahmed Abualsaud <[email protected]>
AuthorDate: Thu Sep 24 14:24:51 2026 -0700
Rename Iceberg table cache options for clarity (#40253)
* rename table cache options
* spotless
* test fixes
* sync
* add to CHANGES
---
CHANGES.md | 1 +
.../org/apache/beam/sdk/io/iceberg/IcebergIO.java | 72 +++++++++++-----------
.../IcebergWriteSchemaTransformProvider.java | 49 ++++++++-------
.../iceberg/IcebergIOSideInputTableCacheTest.java | 50 +++++++--------
.../IcebergWriteSchemaTransformProviderTest.java | 46 +++++++-------
5 files changed, 110 insertions(+), 108 deletions(-)
diff --git a/CHANGES.md b/CHANGES.md
index 63f84e249ec..cc4456bbaae 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -103,6 +103,7 @@
* SolaceIO now supports reading and writing user properties (message metadata)
(Java) ([#40099](https://github.com/apache/beam/issues/40099)).
* [IcebergIO] AddFiles (`IcebergAddFiles` in YAML) can evolve the table schema
before registering files, with `schema_evolution_options`, `required_columns`,
`incompatible_schema_handling` and `unverifiable_file_handling` (Java/YAML,
batch only) ([#40144](https://github.com/apache/beam/issues/40144)).
* [IcebergIO] Added batch and streaming CDC writes that applies
INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE changes to Iceberg V2+ tables by
primary key. Invoke with `IcebergIO.writeCdcRows` (Java) or by setting `mode:
merge-on-read` on the Managed `ICEBERG` write (Java, Python, YAML)
([#39979](https://github.com/apache/beam/issues/39979)).
+* [IcebergIO] Added an optional side-input table cache for writes to
significantly reduce catalog and table requests for large pipelines. A single
worker polls the table and broadcasts it to other workers in the pipeline.
Enable with `IcebergIO.writeRows(...).withSideInputTableCache()` (Java) or by
setting `use_side_input_table_cache: true` on the Managed `ICEBERG` write
(Java, Python) ([#39723](https://github.com/apache/beam/issues/39723)).
## New Features / Improvements
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 bbcaef7a0a6..d98d9944bf3 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
@@ -393,7 +393,7 @@ public class IcebergIO {
.setCatalogConfig(catalog)
.setDistributionMode(DistributionMode.NONE)
.setAutoSharding(false)
- .setUsingSideInputTableCache(false)
+ .setUseSideInputTableCache(false)
.build();
}
@@ -437,13 +437,13 @@ public class IcebergIO {
abstract @Nullable List<String> getSortFields();
- abstract boolean getUsingSideInputTableCache();
+ abstract boolean getUseSideInputTableCache();
- abstract @Nullable Integer getMaximumCacheSize();
+ abstract @Nullable Integer getMaximumTableCacheSize();
- abstract @Nullable Duration getTableRefreshInterval();
+ abstract @Nullable Duration getTableCacheRefreshInterval();
- abstract @Nullable Integer getPollingBuckets();
+ abstract @Nullable Integer getTableCachePollingBuckets();
abstract Builder toBuilder();
@@ -469,13 +469,13 @@ public class IcebergIO {
abstract Builder setSortFields(List<String> sortFields);
- abstract Builder setUsingSideInputTableCache(boolean
usingSideInputTableCache);
+ abstract Builder setUseSideInputTableCache(boolean
useSideInputTableCache);
- abstract Builder setMaximumCacheSize(@Nullable Integer maximumCacheSize);
+ abstract Builder setMaximumTableCacheSize(@Nullable Integer
maximumTableCacheSize);
- abstract Builder setTableRefreshInterval(@Nullable Duration
refreshInterval);
+ abstract Builder setTableCacheRefreshInterval(@Nullable Duration
refreshInterval);
- abstract Builder setPollingBuckets(@Nullable Integer pollingBuckets);
+ abstract Builder setTableCachePollingBuckets(@Nullable Integer
pollingBuckets);
abstract WriteRows build();
}
@@ -569,7 +569,7 @@ public class IcebergIO {
* representations without issuing remote catalog RPCs, drastically
reducing catalog load.
*/
public WriteRows withSideInputTableCache() {
- return toBuilder().setUsingSideInputTableCache(true).build();
+ return toBuilder().setUseSideInputTableCache(true).build();
}
/**
@@ -579,9 +579,10 @@ public class IcebergIO {
* <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();
+ public WriteRows withMaximumTableCacheSize(int maximumTableCacheSize) {
+ Preconditions.checkArgument(
+ maximumTableCacheSize > 0, "maximumTableCacheSize must be greater
than 0");
+ return
toBuilder().setMaximumTableCacheSize(maximumTableCacheSize).build();
}
/**
@@ -589,11 +590,11 @@ public class IcebergIO {
*
* <p>Applicable for unbounded streaming pipelines. Defaults to 5 minutes.
*/
- public WriteRows withTableRefreshInterval(Duration refreshInterval) {
+ public WriteRows withTableCacheRefreshInterval(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();
+ return toBuilder().setTableCacheRefreshInterval(refreshInterval).build();
}
/**
@@ -601,25 +602,26 @@ public class IcebergIO {
* 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();
+ public WriteRows withTableCachePollingBuckets(int pollingBuckets) {
+ Preconditions.checkArgument(
+ pollingBuckets > 0, "tableCachePollingBuckets must be greater than
0");
+ return toBuilder().setTableCachePollingBuckets(pollingBuckets).build();
}
@Override
public void populateDisplayData(DisplayData.Builder builder) {
super.populateDisplayData(builder);
builder.add(
- DisplayData.item("usingSideInputTableCache",
getUsingSideInputTableCache())
+ DisplayData.item("useSideInputTableCache",
getUseSideInputTableCache())
.withLabel("Using Side-Input Table Cache"));
builder.addIfNotNull(
- DisplayData.item("maximumCacheSize", getMaximumCacheSize())
+ DisplayData.item("maximumTableCacheSize", getMaximumTableCacheSize())
.withLabel("Maximum Cache Size"));
builder.addIfNotNull(
- DisplayData.item("tableRefreshInterval", getTableRefreshInterval())
+ DisplayData.item("tableCacheRefreshInterval",
getTableCacheRefreshInterval())
.withLabel("Table Refresh Interval"));
builder.addIfNotNull(
- DisplayData.item("pollingBuckets", getPollingBuckets())
+ DisplayData.item("tableCachePollingBuckets",
getTableCachePollingBuckets())
.withLabel("Catalog Polling Buckets"));
}
@@ -649,29 +651,29 @@ public class IcebergIO {
}
boolean hasSideInputOptions =
- getMaximumCacheSize() != null
- || getTableRefreshInterval() != null
- || getPollingBuckets() != null;
+ getMaximumTableCacheSize() != null
+ || getTableCacheRefreshInterval() != null
+ || getTableCachePollingBuckets() != null;
Preconditions.checkArgument(
- getUsingSideInputTableCache() || !hasSideInputOptions,
- "Cannot specify side-input cache sub-options (maximumCacheSize, "
- + "tableRefreshInterval, pollingBuckets) without enabling
side-input table cache via withSideInputTableCache().");
+ getUseSideInputTableCache() || !hasSideInputOptions,
+ "Cannot specify side-input cache sub-options (maximumTableCacheSize,
"
+ + "tableCacheRefreshInterval, tableCachePollingBuckets) without
enabling side-input table cache via withSideInputTableCache().");
PCollectionView<Map<String, SerializableTableSpec>> metadataView = null;
- if (getUsingSideInputTableCache()) {
+ if (getUseSideInputTableCache()) {
TableMetadataDriver.Builder driverBuilder =
TableMetadataDriver.builder()
.setCatalogConfig(getCatalogConfig())
.setDynamicDestinations(destinations);
- if (getMaximumCacheSize() != null) {
- driverBuilder.setMaximumCacheSize(getMaximumCacheSize());
+ if (getMaximumTableCacheSize() != null) {
+ driverBuilder.setMaximumCacheSize(getMaximumTableCacheSize());
}
- if (getTableRefreshInterval() != null) {
- driverBuilder.setRefreshInterval(getTableRefreshInterval());
+ if (getTableCacheRefreshInterval() != null) {
+ driverBuilder.setRefreshInterval(getTableCacheRefreshInterval());
}
- if (getPollingBuckets() != null) {
- driverBuilder.setPollingBuckets(getPollingBuckets());
+ if (getTableCachePollingBuckets() != null) {
+ driverBuilder.setPollingBuckets(getTableCachePollingBuckets());
}
metadataView = input.apply("GenerateTableMetadataView",
driverBuilder.build().asView());
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 4f568fb6562..91d0ed73835 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
@@ -212,20 +212,20 @@ public class IcebergWriteSchemaTransformProvider
@SchemaFieldDescription(
"Enables expirable side-input caching of Iceberg table metadata across
workers to reduce catalog load.")
- public abstract @Nullable Boolean getUsingSideInputTableCache();
+ public abstract @Nullable Boolean getUseSideInputTableCache();
@SchemaFieldDescription(
"For a streaming pipeline, sets the interval in seconds at which table
metadata is refreshed from the catalog.")
- public abstract @Nullable Integer getTableRefreshIntervalSeconds();
+ public abstract @Nullable Integer getTableCacheRefreshIntervalSeconds();
@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();
+ public abstract @Nullable Integer getMaximumTableCacheSize();
@SchemaFieldDescription(
"Sets the number of parallel buckets/workers used to query the Iceberg
catalog during refreshes. Defaults to 1.")
- public abstract @Nullable Integer getPollingBuckets();
+ public abstract @Nullable Integer getTableCachePollingBuckets();
@SchemaFieldDescription(
"Columns defining row identity (equality-delete fields). Defaults to
the destination table's "
@@ -316,13 +316,14 @@ public class IcebergWriteSchemaTransformProvider
public abstract Builder setWriteProperties(Map<String, String>
writeProperties);
- public abstract Builder setUsingSideInputTableCache(Boolean
usingSideInputTableCache);
+ public abstract Builder setUseSideInputTableCache(Boolean
useSideInputTableCache);
- public abstract Builder setTableRefreshIntervalSeconds(Integer
tableRefreshIntervalSeconds);
+ public abstract Builder setTableCacheRefreshIntervalSeconds(
+ Integer tableCacheRefreshIntervalSeconds);
- public abstract Builder setMaximumCacheSize(Integer maximumCacheSize);
+ public abstract Builder setMaximumTableCacheSize(Integer
maximumTableCacheSize);
- public abstract Builder setPollingBuckets(Integer pollingBuckets);
+ public abstract Builder setTableCachePollingBuckets(Integer
pollingBuckets);
public abstract Builder setMode(String mode);
@@ -419,14 +420,14 @@ public class IcebergWriteSchemaTransformProvider
requireMode(Mode.APPEND, "autosharding", getAutosharding(), unsupported);
requireMode(Mode.APPEND, "write_properties", getWriteProperties(),
unsupported);
requireMode(
- Mode.APPEND, "using_side_input_table_cache",
getUsingSideInputTableCache(), unsupported);
+ Mode.APPEND, "using_side_input_table_cache",
getUseSideInputTableCache(), unsupported);
requireMode(
Mode.APPEND,
"table_refresh_interval_seconds",
- getTableRefreshIntervalSeconds(),
+ getTableCacheRefreshIntervalSeconds(),
unsupported);
- requireMode(Mode.APPEND, "maximum_cache_size", getMaximumCacheSize(),
unsupported);
- requireMode(Mode.APPEND, "polling_buckets", getPollingBuckets(),
unsupported);
+ requireMode(Mode.APPEND, "maximum_cache_size",
getMaximumTableCacheSize(), unsupported);
+ requireMode(Mode.APPEND, "polling_buckets",
getTableCachePollingBuckets(), unsupported);
if (!unsupported.isEmpty()) {
throw new IllegalArgumentException(
String.format(
@@ -535,34 +536,32 @@ public class IcebergWriteSchemaTransformProvider
}
boolean hasSideInputOptions =
- configuration.getTableRefreshIntervalSeconds() != null
- || configuration.getMaximumCacheSize() != null
- || configuration.getPollingBuckets() != null;
+ configuration.getTableCacheRefreshIntervalSeconds() != null
+ || configuration.getMaximumTableCacheSize() != null
+ || configuration.getTableCachePollingBuckets() != null;
- if (!Boolean.TRUE.equals(configuration.getUsingSideInputTableCache())
- && hasSideInputOptions) {
+ if (!Boolean.TRUE.equals(configuration.getUseSideInputTableCache()) &&
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());
+ boolean enableSideInputCache =
Boolean.TRUE.equals(configuration.getUseSideInputTableCache());
if (enableSideInputCache) {
writeTransform = writeTransform.withSideInputTableCache();
- @Nullable Integer refreshSec =
configuration.getTableRefreshIntervalSeconds();
+ @Nullable Integer refreshSec =
configuration.getTableCacheRefreshIntervalSeconds();
if (refreshSec != null) {
writeTransform =
-
writeTransform.withTableRefreshInterval(Duration.standardSeconds(refreshSec));
+
writeTransform.withTableCacheRefreshInterval(Duration.standardSeconds(refreshSec));
}
- @Nullable Integer maxCacheSize = configuration.getMaximumCacheSize();
+ @Nullable Integer maxCacheSize =
configuration.getMaximumTableCacheSize();
if (maxCacheSize != null) {
- writeTransform = writeTransform.withMaximumCacheSize(maxCacheSize);
+ writeTransform =
writeTransform.withMaximumTableCacheSize(maxCacheSize);
}
- @Nullable Integer pollingBuckets = configuration.getPollingBuckets();
+ @Nullable Integer pollingBuckets =
configuration.getTableCachePollingBuckets();
if (pollingBuckets != null) {
- writeTransform = writeTransform.withPollingBuckets(pollingBuckets);
+ writeTransform =
writeTransform.withTableCachePollingBuckets(pollingBuckets);
}
}
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
index 289fabc6991..957dafa7847 100644
---
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
@@ -159,7 +159,7 @@ public class IcebergIOSideInputTableCacheTest implements
Serializable {
IcebergIO.writeRows(catalogConfig)
.to(tableId)
.withSideInputTableCache()
- .withPollingBuckets(1);
+ .withTableCachePollingBuckets(1);
input.apply("WriteToTable", applyDistribution(write));
PipelineResult result = testPipeline.run();
@@ -251,7 +251,7 @@ public class IcebergIOSideInputTableCacheTest implements
Serializable {
IcebergIO.writeRows(catalogConfig)
.to(dynamicDestinations)
.withSideInputTableCache()
- .withPollingBuckets(1);
+ .withTableCachePollingBuckets(1);
input.apply("WriteDynamic", applyDistribution(write));
PipelineResult result = testPipeline.run();
@@ -352,8 +352,8 @@ public class IcebergIOSideInputTableCacheTest implements
Serializable {
IcebergIO.writeRows(catalogConfig)
.to(dynamicDestinations)
.withSideInputTableCache()
- .withMaximumCacheSize(2)
- .withPollingBuckets(1);
+ .withMaximumTableCacheSize(2)
+ .withTableCachePollingBuckets(1);
input.apply("WriteWithSampleCap", applyDistribution(write));
PipelineResult result = testPipeline.run();
@@ -397,8 +397,8 @@ public class IcebergIOSideInputTableCacheTest implements
Serializable {
.to(tableId)
.withSideInputTableCache()
.withTriggeringFrequency(Duration.standardSeconds(1))
- .withTableRefreshInterval(Duration.standardSeconds(2))
- .withPollingBuckets(1);
+ .withTableCacheRefreshInterval(Duration.standardSeconds(2))
+ .withTableCachePollingBuckets(1);
input.apply("StreamingWrite", applyDistribution(write));
PipelineResult result = testPipeline.run();
@@ -467,8 +467,8 @@ public class IcebergIOSideInputTableCacheTest implements
Serializable {
.to(tableId)
.withSideInputTableCache()
.withTriggeringFrequency(Duration.standardSeconds(1))
- .withTableRefreshInterval(Duration.standardSeconds(1))
- .withPollingBuckets(1);
+ .withTableCacheRefreshInterval(Duration.standardSeconds(1))
+ .withTableCachePollingBuckets(1);
input.apply("StreamingWriteEvolved", write);
PipelineResult result = testPipeline.run();
@@ -522,16 +522,16 @@ public class IcebergIOSideInputTableCacheTest implements
Serializable {
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.withMaximumTableCacheSize(0));
+ assertThrows(IllegalArgumentException.class, () ->
write.withMaximumTableCacheSize(-1));
assertThrows(
- IllegalArgumentException.class, () ->
write.withTableRefreshInterval(Duration.ZERO));
+ IllegalArgumentException.class, () ->
write.withTableCacheRefreshInterval(Duration.ZERO));
- assertThrows(IllegalArgumentException.class, () ->
write.withPollingBuckets(0));
- assertThrows(IllegalArgumentException.class, () ->
write.withPollingBuckets(-1));
+ assertThrows(IllegalArgumentException.class, () ->
write.withTableCachePollingBuckets(0));
+ assertThrows(IllegalArgumentException.class, () ->
write.withTableCachePollingBuckets(-1));
- // Unbounded streaming pipeline with maximumCacheSize must fail at expand
+ // Unbounded streaming pipeline with maximumTableCacheSize must fail at
expand
Pipeline p = Pipeline.create();
Schema schema = Schema.builder().addInt64Field("id").build();
TestStream<Row> testStream =
@@ -543,18 +543,18 @@ public class IcebergIOSideInputTableCacheTest implements
Serializable {
// Sub-options specified without withSideInputTableCache() must fail at
expand
IcebergIO.WriteRows writeWithMaxCacheOnly =
- IcebergIO.writeRows(catalogConfig).to(tableId).withMaximumCacheSize(5);
+
IcebergIO.writeRows(catalogConfig).to(tableId).withMaximumTableCacheSize(5);
assertThrows(IllegalArgumentException.class, () ->
streamInput.apply(writeWithMaxCacheOnly));
IcebergIO.WriteRows writeWithRefreshIntervalOnly =
IcebergIO.writeRows(catalogConfig)
.to(tableId)
- .withTableRefreshInterval(Duration.standardMinutes(1));
+ .withTableCacheRefreshInterval(Duration.standardMinutes(1));
assertThrows(
IllegalArgumentException.class, () ->
streamInput.apply(writeWithRefreshIntervalOnly));
IcebergIO.WriteRows writeWithPollingBucketsOnly =
- IcebergIO.writeRows(catalogConfig).to(tableId).withPollingBuckets(2);
+
IcebergIO.writeRows(catalogConfig).to(tableId).withTableCachePollingBuckets(2);
assertThrows(
IllegalArgumentException.class, () ->
streamInput.apply(writeWithPollingBucketsOnly));
@@ -562,7 +562,7 @@ public class IcebergIOSideInputTableCacheTest implements
Serializable {
IcebergIO.writeRows(catalogConfig)
.to(tableId)
.withSideInputTableCache()
- .withMaximumCacheSize(5);
+ .withMaximumTableCacheSize(5);
assertThrows(IllegalArgumentException.class, () ->
streamInput.apply(streamWrite));
}
@@ -574,9 +574,9 @@ public class IcebergIOSideInputTableCacheTest implements
Serializable {
IcebergIO.writeRows(catalogConfig)
.to(tableId)
.withSideInputTableCache()
- .withMaximumCacheSize(100)
- .withTableRefreshInterval(Duration.standardMinutes(10))
- .withPollingBuckets(3);
+ .withMaximumTableCacheSize(100)
+ .withTableCacheRefreshInterval(Duration.standardMinutes(10))
+ .withTableCachePollingBuckets(3);
DisplayData displayData = DisplayData.from(write);
Map<String, String> items = new HashMap<>();
@@ -584,10 +584,10 @@ public class IcebergIOSideInputTableCacheTest implements
Serializable {
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"));
+ assertEquals("true", items.get("useSideInputTableCache"));
+ assertEquals("100", items.get("maximumTableCacheSize"));
+ assertEquals("600000", items.get("tableCacheRefreshInterval"));
+ assertEquals("3", items.get("tableCachePollingBuckets"));
}
private static class EvolveSpecMidExecutionDoFn extends DoFn<Row, 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 c4ac1e22cbf..b408445053f 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
@@ -130,10 +130,10 @@ public class IcebergWriteSchemaTransformProviderTest {
.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)
+ .withFieldValue("use_side_input_table_cache", true)
+ .withFieldValue("table_cache_refresh_interval_seconds", 60)
+ .withFieldValue("maximum_table_cache_size", 100)
+ .withFieldValue("table_cache_polling_buckets", 2)
.build();
new IcebergWriteSchemaTransformProvider().from(transformConfigRow);
@@ -196,9 +196,9 @@ public class IcebergWriteSchemaTransformProviderTest {
.setCatalogName("name")
.setCatalogProperties(properties)
.setDistributionMode(distributionMode.name())
- .setUsingSideInputTableCache(true)
- .setTableRefreshIntervalSeconds(60)
- .setPollingBuckets(1)
+ .setUseSideInputTableCache(true)
+ .setTableCacheRefreshIntervalSeconds(60)
+ .setTableCachePollingBuckets(1)
.build();
PCollectionRowTuple input =
@@ -274,10 +274,10 @@ public class IcebergWriteSchemaTransformProviderTest {
"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"
+ + "use_side_input_table_cache: true\n"
+ + "table_cache_refresh_interval_seconds: 60\n"
+ + "maximum_table_cache_size: 50\n"
+ + "table_cache_polling_buckets: 1\n"
+ "catalog_properties: \n"
+ " type: %s\n"
+ " warehouse: %s",
@@ -316,14 +316,14 @@ public class IcebergWriteSchemaTransformProviderTest {
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
+ // Setting sub-options when use_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)
+ .setMaximumTableCacheSize(50)
.build();
assertThrows(
IllegalArgumentException.class,
@@ -334,7 +334,7 @@ public class IcebergWriteSchemaTransformProviderTest {
.setTable("default.table_refresh")
.setCatalogName("name")
.setCatalogProperties(Collections.singletonMap("type", "hadoop"))
- .setTableRefreshIntervalSeconds(60)
+ .setTableCacheRefreshIntervalSeconds(60)
.build();
assertThrows(
IllegalArgumentException.class,
@@ -345,36 +345,36 @@ public class IcebergWriteSchemaTransformProviderTest {
.setTable("default.table_buckets")
.setCatalogName("name")
.setCatalogProperties(Collections.singletonMap("type", "hadoop"))
- .setPollingBuckets(2)
+ .setTableCachePollingBuckets(2)
.build();
assertThrows(
IllegalArgumentException.class,
() -> dummyInput.apply(provider.from(configWithPollingBucketsOnly)));
- // Setting using_side_input_table_cache to false while setting sub-options
must throw
+ // Setting use_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)
+ .setUseSideInputTableCache(false)
+ .setMaximumTableCacheSize(50)
.build();
assertThrows(
IllegalArgumentException.class,
() -> dummyInput.apply(provider.from(invalidConfigWithFalse)));
- // Explicitly setting using_side_input_table_cache to true with
sub-options succeeds
+ // Explicitly setting use_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)
+ .setUseSideInputTableCache(true)
+ .setMaximumTableCacheSize(50)
+ .setTableCacheRefreshIntervalSeconds(60)
+ .setTableCachePollingBuckets(2)
.build();
assertNotNull(dummyInput.apply(provider.from(validConfig)));
}