jrmccluskey commented on code in PR #40253:
URL: https://github.com/apache/beam/pull/40253#discussion_r4094625995
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergIO.java:
##########
@@ -579,47 +579,49 @@ public WriteRows withSideInputTableCache() {
* <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();
}
/**
* 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) {
+ 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();
}
/**
* 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();
+ public WriteRows withTableCachePollingBuckets(int pollingBuckets) {
+ Preconditions.checkArgument(
+ pollingBuckets > 0, "tableCachePollingBuckets must be greater than
0");
+ return toBuilder().setTableCachePollingBuckets(pollingBuckets).build();
}
Review Comment:
Precommit breakage is due to this symbol not being updated in test files
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]