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..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 @@ -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 static final Counter COMMIT_FAILURES = + 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) { + COMMIT_FAILURES.inc(); // TODO: consider retrying. LOG.warn("Getting exception when committing offset: {}", e.getMessage()); }