This is an automated email from the ASF dual-hosted git repository.

jrmccluskey 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 55e1ecbb138 Add metric for Kafka offset commit failures (#38889)
55e1ecbb138 is described below

commit 55e1ecbb138ccf8a8ce423f68b618db22bc0b23d
Author: KRITI MITTAL <[email protected]>
AuthorDate: Wed Jul 22 01:36:50 2026 +0530

    Add metric for Kafka offset commit failures (#38889)
    
    * Add metric for Kafka offset commit failures
    
    * Apply suggestions from code review
    
    Co-authored-by: gemini-code-assist[bot] 
<176961590+gemini-code-assist[bot]@users.noreply.github.com>
    
    ---------
    
    Co-authored-by: Jack McCluskey 
<[email protected]>
    Co-authored-by: gemini-code-assist[bot] 
<176961590+gemini-code-assist[bot]@users.noreply.github.com>
---
 .../main/java/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 ac6650c354d..281a5b9a47c 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 java.util.Map;
 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<K, V>
 
   static class CommitOffsetDoFn extends DoFn<KV<KafkaSourceDescriptor, Long>, 
Void> {
     private static final Logger LOG = 
LoggerFactory.getLogger(CommitOffsetDoFn.class);
+    private static final Counter COMMIT_FAILURES =
+        Metrics.counter(CommitOffsetDoFn.class, "commit-failures");
+
     private final Map<String, Object> consumerConfig;
     private final SerializableFunction<Map<String, Object>, Consumer<byte[], 
byte[]>>
         consumerFactoryFn;
@@ -85,6 +90,7 @@ public class KafkaCommitOffset<K, V>
                   element.getKey().getTopicPartition(),
                   new OffsetAndMetadata(element.getValue() + 1)));
         } catch (Exception e) {
+          COMMIT_FAILURES.inc();
           // TODO: consider retrying.
           LOG.warn("Getting exception when committing offset: {}", 
e.getMessage());
         }

Reply via email to