From e20ee61f84a3d91d251bb8d65dc253ec7fd105a2 Mon Sep 17 00:00:00 2001 From: Piotr Wolski Date: Mon, 20 Jul 2026 11:02:23 -0600 Subject: [PATCH 1/2] Fix NoSuchElementException from unchecked Optional.get() in Kafka consumer instrumentation extractGroup, extractClusterId, and extractBootstrapServers called Optional.get() without checking isPresent()/using orElse(), which threw NoSuchElementException when the underlying consumer group, metadata, or bootstrap servers were not captured. Introduced in 1.64 and observed in production error telemetry. Co-Authored-By: Claude Sonnet 5 --- .../KafkaConsumerInstrumentationHelper.java | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInstrumentationHelper.java b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInstrumentationHelper.java index 11c020e78e3..57e0baa7edd 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInstrumentationHelper.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInstrumentationHelper.java @@ -2,12 +2,13 @@ import datadog.trace.bootstrap.ContextStore; import datadog.trace.instrumentation.kafka_common.MetadataState; +import java.util.Optional; import org.apache.kafka.clients.Metadata; public class KafkaConsumerInstrumentationHelper { public static String extractGroup(KafkaConsumerInfo kafkaConsumerInfo) { if (kafkaConsumerInfo != null) { - return kafkaConsumerInfo.getConsumerGroup().get(); + return kafkaConsumerInfo.getConsumerGroup().orElse(null); } return null; } @@ -16,9 +17,9 @@ public static String extractClusterId( KafkaConsumerInfo kafkaConsumerInfo, ContextStore metadataContextStore) { if (kafkaConsumerInfo != null) { - Metadata metadata = kafkaConsumerInfo.getmetadata().get(); - if (metadata != null) { - MetadataState state = metadataContextStore.get(metadata); + Optional metadata = kafkaConsumerInfo.getmetadata(); + if (metadata.isPresent()) { + MetadataState state = metadataContextStore.get(metadata.get()); return state != null ? state.clusterId : null; } } @@ -26,6 +27,6 @@ public static String extractClusterId( } public static String extractBootstrapServers(KafkaConsumerInfo kafkaConsumerInfo) { - return kafkaConsumerInfo == null ? null : kafkaConsumerInfo.getBootstrapServers().get(); + return kafkaConsumerInfo == null ? null : kafkaConsumerInfo.getBootstrapServers().orElse(null); } } From 62918831b8a593c4a3da44642cf620e49f8f892b Mon Sep 17 00:00:00 2001 From: Piotr Wolski Date: Mon, 20 Jul 2026 11:26:21 -0600 Subject: [PATCH 2/2] Add unit tests for KafkaConsumerInstrumentationHelper Covers the null/empty-Optional cases for extractGroup, extractClusterId, and extractBootstrapServers to prevent regressions of the previously fixed NoSuchElementException. Co-Authored-By: Claude Sonnet 5 --- ...afkaConsumerInstrumentationHelperTest.java | 94 +++++++++++++++++++ 1 file changed, 94 insertions(+) create mode 100644 dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/test/java/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInstrumentationHelperTest.java diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/test/java/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInstrumentationHelperTest.java b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/test/java/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInstrumentationHelperTest.java new file mode 100644 index 00000000000..dee19831317 --- /dev/null +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/test/java/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInstrumentationHelperTest.java @@ -0,0 +1,94 @@ +package datadog.trace.instrumentation.kafka_clients38; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import datadog.trace.bootstrap.ContextStore; +import datadog.trace.instrumentation.kafka_common.MetadataState; +import org.apache.kafka.clients.Metadata; +import org.junit.jupiter.api.Test; + +class KafkaConsumerInstrumentationHelperTest { + + @SuppressWarnings("unchecked") + private final ContextStore metadataContextStore = + mock(ContextStore.class); + + @Test + void extractGroupReturnsNullForNullKafkaConsumerInfo() { + assertNull(KafkaConsumerInstrumentationHelper.extractGroup(null)); + } + + @Test + void extractGroupReturnsNullWhenConsumerGroupIsNull() { + KafkaConsumerInfo kafkaConsumerInfo = new KafkaConsumerInfo(null, null, "localhost:9092"); + assertNull(KafkaConsumerInstrumentationHelper.extractGroup(kafkaConsumerInfo)); + } + + @Test + void extractGroupReturnsConsumerGroupWhenPresent() { + KafkaConsumerInfo kafkaConsumerInfo = + new KafkaConsumerInfo("test-group", null, "localhost:9092"); + assertEquals("test-group", KafkaConsumerInstrumentationHelper.extractGroup(kafkaConsumerInfo)); + } + + @Test + void extractBootstrapServersReturnsNullForNullKafkaConsumerInfo() { + assertNull(KafkaConsumerInstrumentationHelper.extractBootstrapServers(null)); + } + + @Test + void extractBootstrapServersReturnsNullWhenBootstrapServersIsNull() { + KafkaConsumerInfo kafkaConsumerInfo = new KafkaConsumerInfo("test-group", null, null); + assertNull(KafkaConsumerInstrumentationHelper.extractBootstrapServers(kafkaConsumerInfo)); + } + + @Test + void extractBootstrapServersReturnsValueWhenPresent() { + KafkaConsumerInfo kafkaConsumerInfo = + new KafkaConsumerInfo("test-group", null, "localhost:9092"); + assertEquals( + "localhost:9092", + KafkaConsumerInstrumentationHelper.extractBootstrapServers(kafkaConsumerInfo)); + } + + @Test + void extractClusterIdReturnsNullForNullKafkaConsumerInfo() { + assertNull(KafkaConsumerInstrumentationHelper.extractClusterId(null, metadataContextStore)); + } + + @Test + void extractClusterIdReturnsNullWhenMetadataIsNull() { + KafkaConsumerInfo kafkaConsumerInfo = new KafkaConsumerInfo("test-group", "localhost:9092"); + assertNull( + KafkaConsumerInstrumentationHelper.extractClusterId( + kafkaConsumerInfo, metadataContextStore)); + } + + @Test + void extractClusterIdReturnsNullWhenNoStateForMetadata() { + Metadata metadata = mock(Metadata.class); + KafkaConsumerInfo kafkaConsumerInfo = + new KafkaConsumerInfo("test-group", metadata, "localhost:9092"); + when(metadataContextStore.get(metadata)).thenReturn(null); + assertNull( + KafkaConsumerInstrumentationHelper.extractClusterId( + kafkaConsumerInfo, metadataContextStore)); + } + + @Test + void extractClusterIdReturnsClusterIdWhenStatePresent() { + Metadata metadata = mock(Metadata.class); + KafkaConsumerInfo kafkaConsumerInfo = + new KafkaConsumerInfo("test-group", metadata, "localhost:9092"); + MetadataState state = new MetadataState(); + state.clusterId = "cluster-1"; + when(metadataContextStore.get(metadata)).thenReturn(state); + assertEquals( + "cluster-1", + KafkaConsumerInstrumentationHelper.extractClusterId( + kafkaConsumerInfo, metadataContextStore)); + } +}