From a657c305c656d0ced95cde3a1e007358c5546c0e Mon Sep 17 00:00:00 2001 From: jrmccluskey Date: Thu, 11 Sep 2025 10:33:18 -0400 Subject: [PATCH 1/3] Add ThrottlingSignaler class to the Java SDK --- .../throttling/ThrottlingSignaler.java | 46 +++++++++++++++++++ 1 file changed, 46 insertions(+) create mode 100644 sdks/java/io/components/src/main/java/org/apache/beam/sdk/io/components/throttling/ThrottlingSignaler.java diff --git a/sdks/java/io/components/src/main/java/org/apache/beam/sdk/io/components/throttling/ThrottlingSignaler.java b/sdks/java/io/components/src/main/java/org/apache/beam/sdk/io/components/throttling/ThrottlingSignaler.java new file mode 100644 index 000000000000..4c6f68ed621b --- /dev/null +++ b/sdks/java/io/components/src/main/java/org/apache/beam/sdk/io/components/throttling/ThrottlingSignaler.java @@ -0,0 +1,46 @@ +/* + * 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.components.throttling; + +import org.apache.beam.sdk.metrics.Counter; +import org.apache.beam.sdk.metrics.Metrics; +/** + * The ThrottlingSignaler is a utility class for IOs to signal to the runner + * that a process is being throttled, preventing autoscaling. This is primarily + * used when making calls to a remote service where quotas and rate limiting + * are reasonable considerations. + */ +public class ThrottlingSignaler { + private Counter throttle_counter; + + public ThrottlingSignaler(String namespace) { + this.throttle_counter = Metrics.counter(namespace, Metrics.THROTTLE_TIME_COUNTER_NAME); + } + + public ThrottlingSignaler() { + ThrottlingSignaler(""); + } + + /** + * Signal that a transform has been throttled for an amount of time + * represented in milliseconds. + */ + public void signalThrottling(long milliseconds) { + throttle_counter.inc(milliseconds); + } +} From f435b7942341bab3bdd9870b1d92405d8c0c293c Mon Sep 17 00:00:00 2001 From: Jack McCluskey <34928439+jrmccluskey@users.noreply.github.com> Date: Thu, 11 Sep 2025 10:40:20 -0400 Subject: [PATCH 2/3] Update sdks/java/io/components/src/main/java/org/apache/beam/sdk/io/components/throttling/ThrottlingSignaler.java Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com> --- .../io/components/throttling/ThrottlingSignaler.java | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/sdks/java/io/components/src/main/java/org/apache/beam/sdk/io/components/throttling/ThrottlingSignaler.java b/sdks/java/io/components/src/main/java/org/apache/beam/sdk/io/components/throttling/ThrottlingSignaler.java index 4c6f68ed621b..7805c3befcbf 100644 --- a/sdks/java/io/components/src/main/java/org/apache/beam/sdk/io/components/throttling/ThrottlingSignaler.java +++ b/sdks/java/io/components/src/main/java/org/apache/beam/sdk/io/components/throttling/ThrottlingSignaler.java @@ -26,21 +26,21 @@ * are reasonable considerations. */ public class ThrottlingSignaler { - private Counter throttle_counter; + private final Counter throttleCounter; public ThrottlingSignaler(String namespace) { - this.throttle_counter = Metrics.counter(namespace, Metrics.THROTTLE_TIME_COUNTER_NAME); + this.throttleCounter = Metrics.counter(namespace, Metrics.THROTTLE_TIME_COUNTER_NAME); } public ThrottlingSignaler() { - ThrottlingSignaler(""); + this(""); } - /** + /** * Signal that a transform has been throttled for an amount of time * represented in milliseconds. */ public void signalThrottling(long milliseconds) { - throttle_counter.inc(milliseconds); + throttleCounter.inc(milliseconds); } } From 6d58f8902329a90af82f04e42bca12db193411cd Mon Sep 17 00:00:00 2001 From: jrmccluskey Date: Tue, 7 Oct 2025 09:59:16 -0400 Subject: [PATCH 3/3] set default namespace --- .../beam/sdk/io/components/throttling/ThrottlingSignaler.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdks/java/io/components/src/main/java/org/apache/beam/sdk/io/components/throttling/ThrottlingSignaler.java b/sdks/java/io/components/src/main/java/org/apache/beam/sdk/io/components/throttling/ThrottlingSignaler.java index 7805c3befcbf..894c9294bed4 100644 --- a/sdks/java/io/components/src/main/java/org/apache/beam/sdk/io/components/throttling/ThrottlingSignaler.java +++ b/sdks/java/io/components/src/main/java/org/apache/beam/sdk/io/components/throttling/ThrottlingSignaler.java @@ -33,7 +33,7 @@ public ThrottlingSignaler(String namespace) { } public ThrottlingSignaler() { - this(""); + this(Metrics.THROTTLE_TIME_NAMESPACE); } /**