Skip to content

Commit 55e1ecb

Browse files
Kriti-dev07jrmccluskeygemini-code-assist[bot]
authored
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 <34928439+jrmccluskey@users.noreply.github.com> Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com>
1 parent 83fcb2f commit 55e1ecb

1 file changed

Lines changed: 6 additions & 0 deletions

File tree

sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaCommitOffset.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,8 @@
2525
import org.apache.beam.sdk.coders.KvCoder;
2626
import org.apache.beam.sdk.coders.VarLongCoder;
2727
import org.apache.beam.sdk.coders.VoidCoder;
28+
import org.apache.beam.sdk.metrics.Counter;
29+
import org.apache.beam.sdk.metrics.Metrics;
2830
import org.apache.beam.sdk.schemas.NoSuchSchemaException;
2931
import org.apache.beam.sdk.transforms.DoFn;
3032
import org.apache.beam.sdk.transforms.MapElements;
@@ -62,6 +64,9 @@ public class KafkaCommitOffset<K, V>
6264

6365
static class CommitOffsetDoFn extends DoFn<KV<KafkaSourceDescriptor, Long>, Void> {
6466
private static final Logger LOG = LoggerFactory.getLogger(CommitOffsetDoFn.class);
67+
private static final Counter COMMIT_FAILURES =
68+
Metrics.counter(CommitOffsetDoFn.class, "commit-failures");
69+
6570
private final Map<String, Object> consumerConfig;
6671
private final SerializableFunction<Map<String, Object>, Consumer<byte[], byte[]>>
6772
consumerFactoryFn;
@@ -85,6 +90,7 @@ public void processElement(@Element KV<KafkaSourceDescriptor, Long> element) {
8590
element.getKey().getTopicPartition(),
8691
new OffsetAndMetadata(element.getValue() + 1)));
8792
} catch (Exception e) {
93+
COMMIT_FAILURES.inc();
8894
// TODO: consider retrying.
8995
LOG.warn("Getting exception when committing offset: {}", e.getMessage());
9096
}

0 commit comments

Comments
 (0)