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());
}