From c0c974e891c10448b1f8a4e1c042fa4898c48144 Mon Sep 17 00:00:00 2001 From: Harsh Wadhawe Date: Wed, 12 Aug 2026 12:16:43 -0500 Subject: [PATCH 1/2] Add metrics for Kafka source failures Three failures in the Kafka source were only written to logs: failed offset commits, failed offset resets after a negative acknowledgement, and failed buffer writes. Operators had no way to see them in dashboards or alerts. This adds a counter for each one. Partially addresses #7074. Signed-off-by: Harsh Wadhawe --- .../kafka/consumer/KafkaCustomConsumer.java | 6 +- .../kafka/util/KafkaTopicConsumerMetrics.java | 21 ++++ .../consumer/KafkaCustomConsumerTest.java | 98 +++++++++++++++++++ 3 files changed, 124 insertions(+), 1 deletion(-) diff --git a/data-prepper-plugins/kafka-plugins/src/main/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumer.java b/data-prepper-plugins/kafka-plugins/src/main/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumer.java index 74aa88520e..cc81160f06 100644 --- a/data-prepper-plugins/kafka-plugins/src/main/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumer.java +++ b/data-prepper-plugins/kafka-plugins/src/main/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumer.java @@ -304,7 +304,7 @@ private void addAcknowledgedOffsets(final TopicPartition topicPartition, final R updateOffsetsToCommit(topicPartition, offsetAndMetadata); } - private void resetOffsets() { + void resetOffsets() { // resetting offsets is similar to committing acknowledged offsets. Throttle the frequency of resets by // checking current time with last commit time. Same "lastCommitTime" and commit interval are used in both cases long currentTimeMillis = System.currentTimeMillis(); @@ -326,6 +326,7 @@ private void resetOffsets() { final long epoch = getCurrentTimeNanos(); ownedPartitionsEpoch.put(partition, epoch); } catch (Exception e) { + topicMetrics.getNumberOfOffsetResetFailures().increment(); LOG.error("Failed to seek to last committed offset upon negative acknowledgement {}", partition, e); } }); @@ -381,9 +382,11 @@ private void commitOffsets(boolean forceCommit) { consumer.commitSync(offsetsToCommit); lastCommitTime = currentTimeMillis; } catch (final RebalanceInProgressException ex) { + topicMetrics.getNumberOfCommitFailures().increment(); LOG.error("Failed to commit offsets in topic {} due to rebalance in progress", topicName, ex); return; } catch (Exception e) { + topicMetrics.getNumberOfCommitFailures().increment(); LOG.error("Failed to commit offsets in topic {}", topicName, e); } @@ -541,6 +544,7 @@ private void processRecords(final AcknowledgementSet acknowledgementSet, final L if (e instanceof SizeOverflowException) { topicMetrics.getNumberOfBufferSizeOverflows().increment(); } else { + topicMetrics.getNumberOfBufferWriteFailures().increment(); LOG.debug("Error while adding record to buffer, retrying ", e); } try { diff --git a/data-prepper-plugins/kafka-plugins/src/main/java/org/opensearch/dataprepper/plugins/kafka/util/KafkaTopicConsumerMetrics.java b/data-prepper-plugins/kafka-plugins/src/main/java/org/opensearch/dataprepper/plugins/kafka/util/KafkaTopicConsumerMetrics.java index aba517177e..c23f1c6bad 100644 --- a/data-prepper-plugins/kafka-plugins/src/main/java/org/opensearch/dataprepper/plugins/kafka/util/KafkaTopicConsumerMetrics.java +++ b/data-prepper-plugins/kafka-plugins/src/main/java/org/opensearch/dataprepper/plugins/kafka/util/KafkaTopicConsumerMetrics.java @@ -32,6 +32,9 @@ public class KafkaTopicConsumerMetrics { static final String NUMBER_OF_RECORDS_COMMITTED = "numberOfRecordsCommitted"; static final String NUMBER_OF_RECORDS_CONSUMED = "numberOfRecordsConsumed"; static final String NUMBER_OF_BYTES_CONSUMED = "numberOfBytesConsumed"; + static final String NUMBER_OF_COMMIT_FAILURES = "numberOfCommitFailures"; + static final String NUMBER_OF_OFFSET_RESET_FAILURES = "numberOfOffsetResetFailures"; + static final String NUMBER_OF_BUFFER_WRITE_FAILURES = "numberOfBufferWriteFailures"; static final String ACTUAL_POLL_INTERVAL = "actualPollInterval"; private final String topicName; @@ -49,6 +52,9 @@ public class KafkaTopicConsumerMetrics { private final Counter numberOfRecordsCommitted; private final Counter numberOfRecordsConsumed; private final Counter numberOfBytesConsumed; + private final Counter numberOfCommitFailures; + private final Counter numberOfOffsetResetFailures; + private final Counter numberOfBufferWriteFailures; private final Timer timeBetweenPollCalls; private Instant lastPollTime; @@ -69,6 +75,9 @@ public KafkaTopicConsumerMetrics(final String topicName, final PluginMetrics plu this.numberOfPollAuthErrors = pluginMetrics.counter(getTopicMetricName(NUMBER_OF_POLL_AUTH_ERRORS, topicNameInMetrics)); this.numberOfPositiveAcknowledgements = pluginMetrics.counter(getTopicMetricName(NUMBER_OF_POSITIVE_ACKNOWLEDGEMENTS, topicNameInMetrics)); this.numberOfNegativeAcknowledgements = pluginMetrics.counter(getTopicMetricName(NUMBER_OF_NEGATIVE_ACKNOWLEDGEMENTS, topicNameInMetrics)); + this.numberOfCommitFailures = pluginMetrics.counter(getTopicMetricName(NUMBER_OF_COMMIT_FAILURES, topicNameInMetrics)); + this.numberOfOffsetResetFailures = pluginMetrics.counter(getTopicMetricName(NUMBER_OF_OFFSET_RESET_FAILURES, topicNameInMetrics)); + this.numberOfBufferWriteFailures = pluginMetrics.counter(getTopicMetricName(NUMBER_OF_BUFFER_WRITE_FAILURES, topicNameInMetrics)); this.timeBetweenPollCalls = pluginMetrics.timer(getTopicMetricName(ACTUAL_POLL_INTERVAL, topicNameInMetrics)); lastPollTime = Instant.now(); } @@ -175,6 +184,18 @@ public Counter getNumberOfPositiveAcknowledgements() { return numberOfPositiveAcknowledgements; } + public Counter getNumberOfCommitFailures() { + return numberOfCommitFailures; + } + + public Counter getNumberOfOffsetResetFailures() { + return numberOfOffsetResetFailures; + } + + public Counter getNumberOfBufferWriteFailures() { + return numberOfBufferWriteFailures; + } + public void recordTimeBetweenPolls() { final long timeBetweenPolls = Instant.now().toEpochMilli() - lastPollTime.toEpochMilli(); timeBetweenPollCalls.record(timeBetweenPolls, TimeUnit.MILLISECONDS); diff --git a/data-prepper-plugins/kafka-plugins/src/test/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumerTest.java b/data-prepper-plugins/kafka-plugins/src/test/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumerTest.java index e1f5b020de..1346dfbf14 100644 --- a/data-prepper-plugins/kafka-plugins/src/test/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumerTest.java +++ b/data-prepper-plugins/kafka-plugins/src/test/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumerTest.java @@ -83,6 +83,7 @@ import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.anyMap; +import static org.mockito.Mockito.atLeastOnce; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.times; @@ -176,6 +177,9 @@ public void setUp() throws JsonProcessingException { when(topicMetrics.getNumberOfDeserializationErrors()).thenReturn(counter); when(topicMetrics.getNumberOfInvalidTimeStamps()).thenReturn(counter); when(topicMetrics.getNumberOfPollAuthErrors()).thenReturn(counter); + when(topicMetrics.getNumberOfCommitFailures()).thenReturn(counter); + when(topicMetrics.getNumberOfOffsetResetFailures()).thenReturn(counter); + when(topicMetrics.getNumberOfBufferWriteFailures()).thenReturn(counter); when(topicConfig.getThreadWaitingTime()).thenReturn(Duration.ofSeconds(1)); when(topicConfig.getSerdeFormat()).thenReturn(MessageFormat.PLAINTEXT); when(topicConfig.getAutoCommit()).thenReturn(false); @@ -875,11 +879,105 @@ private ConsumerRecords createJsonRecords(String topic) throws Exception { return new ConsumerRecords(records); } + @ParameterizedTest + @MethodSource("provideExceptionsFromCommit") + public void testCommitOffsets_whenCommitFails_thenIncrementsCommitFailureCounter(final Exception commitException) throws Exception { + final Counter commitFailureCounter = mock(Counter.class); + when(topicMetrics.getNumberOfCommitFailures()).thenReturn(commitFailureCounter); + + final String topic = topicConfig.getName(); + final TopicPartition topicPartition = new TopicPartition(topic, testPartition); + when(topicConfig.getCommitInterval()).thenReturn(Duration.ofMillis(0)); + + consumer = createObjectUnderTest("plaintext", false); + consumer.onPartitionsAssigned(List.of(topicPartition)); + + consumerRecords = createPlainTextRecords(topic, 100L); + when(kafkaConsumer.poll(any(Duration.class))).thenReturn(consumerRecords); + consumer.consumeRecords(); + + doThrow(commitException).when(kafkaConsumer).commitSync(anyMap()); + + // onPartitionsRevoked forces a commit of the offsets gathered above + consumer.onPartitionsRevoked(List.of(topicPartition)); + + verify(commitFailureCounter).increment(); + } + + @Test + public void testResetOffsets_whenSeekFails_thenIncrementsOffsetResetFailureCounter() throws Exception { + final Counter offsetResetFailureCounter = mock(Counter.class); + when(topicMetrics.getNumberOfOffsetResetFailures()).thenReturn(offsetResetFailureCounter); + + final String topic = topicConfig.getName(); + when(topicConfig.getCommitInterval()).thenReturn(Duration.ofMillis(0)); + consumerRecords = createPlainTextRecords(topic, 0L); + when(kafkaConsumer.poll(any(Duration.class))).thenReturn(consumerRecords); + + consumer = createObjectUnderTest("plaintext", true); + consumer.onPartitionsAssigned(List.of(new TopicPartition(topic, testPartition))); + consumer.consumeRecords(); + + final Map.Entry>, CheckpointState> bufferRecords = buffer.read(1000); + for (final Record record : new ArrayList<>(bufferRecords.getKey())) { + record.getData().getEventHandle().release(false); + } + // Negative acknowledgement adds the partition to the set that resetOffsets() seeks + await().atMost(delayTime.plusMillis(5000)) + .until(() -> consumer.getTopicMetrics().getNumberOfNegativeAcknowledgements().count() == 1.0); + + doThrow(new RuntimeException("Failed to look up committed offset")) + .when(kafkaConsumer).committed(any(TopicPartition.class)); + + consumer.resetOffsets(); + + verify(offsetResetFailureCounter).increment(); + } + + @Test + public void testConsumeRecords_whenBufferWriteFails_thenIncrementsBufferWriteFailureCounter() throws Exception { + final Counter bufferWriteFailureCounter = mock(Counter.class); + when(topicMetrics.getNumberOfBufferWriteFailures()).thenReturn(bufferWriteFailureCounter); + + when(topicConfig.getMaxPollInterval()).thenReturn(Duration.ofMillis(4000)); + final String topic = topicConfig.getName(); + consumerRecords = createPlainTextRecords(topic, 0L); + doAnswer((i) -> { + if (!paused && !resumed) { + throw new TimeoutException(); + } + buffer.writeAll(i.getArgument(0), i.getArgument(1)); + return null; + }).when(mockBuffer).writeAll(any(), anyInt()); + doAnswer((i) -> { + if (paused && !resumed) { + return List.of(); + } + return consumerRecords; + }).when(kafkaConsumer).poll(any(Duration.class)); + + consumer = createObjectUnderTestWithMockBuffer("plaintext"); + try { + consumer.onPartitionsAssigned(List.of(new TopicPartition(topic, testPartition))); + consumer.consumeRecords(); + } catch (Exception e) {} + + verify(bufferWriteFailureCounter, atLeastOnce()).increment(); + // A write timeout is not a size overflow, so the overflow counter must stay untouched + assertEquals(0.0, overflowCount); + } + private static Stream provideExceptionsFromBufferWrite() { return Stream.of( Arguments.of(new SizeOverflowException("size overflow")), Arguments.of(new TimeoutException())); } + + private static Stream provideExceptionsFromCommit() { + return Stream.of( + Arguments.of(new RebalanceInProgressException("Rebalance in progress")), + Arguments.of(new RuntimeException("Generic commit failure"))); + } } From 3d8e3ff7116d817b69950de83f78c63f6611ca4b Mon Sep 17 00:00:00 2001 From: Harsh Wadhawe Date: Wed, 12 Aug 2026 12:21:15 -0500 Subject: [PATCH 2/2] Use reflection in test instead of widening method visibility Signed-off-by: Harsh Wadhawe --- .../plugins/kafka/consumer/KafkaCustomConsumer.java | 2 +- .../plugins/kafka/consumer/KafkaCustomConsumerTest.java | 4 +++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/data-prepper-plugins/kafka-plugins/src/main/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumer.java b/data-prepper-plugins/kafka-plugins/src/main/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumer.java index cc81160f06..700b1b4ba8 100644 --- a/data-prepper-plugins/kafka-plugins/src/main/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumer.java +++ b/data-prepper-plugins/kafka-plugins/src/main/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumer.java @@ -304,7 +304,7 @@ private void addAcknowledgedOffsets(final TopicPartition topicPartition, final R updateOffsetsToCommit(topicPartition, offsetAndMetadata); } - void resetOffsets() { + private void resetOffsets() { // resetting offsets is similar to committing acknowledged offsets. Throttle the frequency of resets by // checking current time with last commit time. Same "lastCommitTime" and commit interval are used in both cases long currentTimeMillis = System.currentTimeMillis(); diff --git a/data-prepper-plugins/kafka-plugins/src/test/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumerTest.java b/data-prepper-plugins/kafka-plugins/src/test/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumerTest.java index 1346dfbf14..eefbbfa740 100644 --- a/data-prepper-plugins/kafka-plugins/src/test/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumerTest.java +++ b/data-prepper-plugins/kafka-plugins/src/test/java/org/opensearch/dataprepper/plugins/kafka/consumer/KafkaCustomConsumerTest.java @@ -929,7 +929,9 @@ public void testResetOffsets_whenSeekFails_thenIncrementsOffsetResetFailureCount doThrow(new RuntimeException("Failed to look up committed offset")) .when(kafkaConsumer).committed(any(TopicPartition.class)); - consumer.resetOffsets(); + final java.lang.reflect.Method method = consumer.getClass().getDeclaredMethod("resetOffsets"); + method.setAccessible(true); + method.invoke(consumer); verify(offsetResetFailureCounter).increment(); }