From b4bbb9298e6caf9a0241902ccee6a89d9f27c56c Mon Sep 17 00:00:00 2001 From: Kriti-dev07 Date: Wed, 10 Jun 2026 15:21:45 +0530 Subject: [PATCH 1/2] Add metric for Kafka offset commit failures --- .../org/apache/beam/sdk/io/kafka/KafkaCommitOffset.java | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaCommitOffset.java b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaCommitOffset.java index ac6650c354d4..f84b2df306ca 100644 --- a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaCommitOffset.java +++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaCommitOffset.java @@ -25,6 +25,8 @@ import org.apache.beam.sdk.coders.KvCoder; import org.apache.beam.sdk.coders.VarLongCoder; import org.apache.beam.sdk.coders.VoidCoder; +import org.apache.beam.sdk.metrics.Counter; +import org.apache.beam.sdk.metrics.Metrics; import org.apache.beam.sdk.schemas.NoSuchSchemaException; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.transforms.MapElements; @@ -62,6 +64,9 @@ public class KafkaCommitOffset static class CommitOffsetDoFn extends DoFn, Void> { private static final Logger LOG = LoggerFactory.getLogger(CommitOffsetDoFn.class); + private final Counter commitFailures = + Metrics.counter(CommitOffsetDoFn.class, "commit-failures"); + private final Map consumerConfig; private final SerializableFunction, Consumer> consumerFactoryFn; @@ -85,6 +90,7 @@ public void processElement(@Element KV element) { element.getKey().getTopicPartition(), new OffsetAndMetadata(element.getValue() + 1))); } catch (Exception e) { + commitFailures.inc(); // TODO: consider retrying. LOG.warn("Getting exception when committing offset: {}", e.getMessage()); } From d4f2280d26534380a2f6aed7e6293f61136ac38a Mon Sep 17 00:00:00 2001 From: Jack McCluskey <34928439+jrmccluskey@users.noreply.github.com> Date: Tue, 14 Jul 2026 13:27:11 -0400 Subject: [PATCH 2/2] Apply suggestions from code review Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com> --- .../java/org/apache/beam/sdk/io/kafka/KafkaCommitOffset.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaCommitOffset.java b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaCommitOffset.java index f84b2df306ca..281a5b9a47c6 100644 --- a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaCommitOffset.java +++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaCommitOffset.java @@ -64,7 +64,7 @@ public class KafkaCommitOffset static class CommitOffsetDoFn extends DoFn, Void> { private static final Logger LOG = LoggerFactory.getLogger(CommitOffsetDoFn.class); - private final Counter commitFailures = + private static final Counter COMMIT_FAILURES = Metrics.counter(CommitOffsetDoFn.class, "commit-failures"); private final Map consumerConfig; @@ -90,7 +90,7 @@ public void processElement(@Element KV element) { element.getKey().getTopicPartition(), new OffsetAndMetadata(element.getValue() + 1))); } catch (Exception e) { - commitFailures.inc(); + COMMIT_FAILURES.inc(); // TODO: consider retrying. LOG.warn("Getting exception when committing offset: {}", e.getMessage()); }