Skip to content

Commit c66a0b5

Browse files
Fix NoSuchElementException from unchecked Optional.get() in Kafka consumer instrumentation (#12007)
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 <noreply@anthropic.com> 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 <noreply@anthropic.com> Co-authored-by: bruce.bujon <bruce.bujon@datadoghq.com>
1 parent ef8a504 commit c66a0b5

2 files changed

Lines changed: 100 additions & 5 deletions

File tree

dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInstrumentationHelper.java

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -2,12 +2,13 @@
22

33
import datadog.trace.bootstrap.ContextStore;
44
import datadog.trace.instrumentation.kafka_common.MetadataState;
5+
import java.util.Optional;
56
import org.apache.kafka.clients.Metadata;
67

78
public class KafkaConsumerInstrumentationHelper {
89
public static String extractGroup(KafkaConsumerInfo kafkaConsumerInfo) {
910
if (kafkaConsumerInfo != null) {
10-
return kafkaConsumerInfo.getConsumerGroup().get();
11+
return kafkaConsumerInfo.getConsumerGroup().orElse(null);
1112
}
1213
return null;
1314
}
@@ -16,16 +17,16 @@ public static String extractClusterId(
1617
KafkaConsumerInfo kafkaConsumerInfo,
1718
ContextStore<Metadata, MetadataState> metadataContextStore) {
1819
if (kafkaConsumerInfo != null) {
19-
Metadata metadata = kafkaConsumerInfo.getmetadata().get();
20-
if (metadata != null) {
21-
MetadataState state = metadataContextStore.get(metadata);
20+
Optional<Metadata> metadata = kafkaConsumerInfo.getmetadata();
21+
if (metadata.isPresent()) {
22+
MetadataState state = metadataContextStore.get(metadata.get());
2223
return state != null ? state.clusterId : null;
2324
}
2425
}
2526
return null;
2627
}
2728

2829
public static String extractBootstrapServers(KafkaConsumerInfo kafkaConsumerInfo) {
29-
return kafkaConsumerInfo == null ? null : kafkaConsumerInfo.getBootstrapServers().get();
30+
return kafkaConsumerInfo == null ? null : kafkaConsumerInfo.getBootstrapServers().orElse(null);
3031
}
3132
}
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,94 @@
1+
package datadog.trace.instrumentation.kafka_clients38;
2+
3+
import static org.junit.jupiter.api.Assertions.assertEquals;
4+
import static org.junit.jupiter.api.Assertions.assertNull;
5+
import static org.mockito.Mockito.mock;
6+
import static org.mockito.Mockito.when;
7+
8+
import datadog.trace.bootstrap.ContextStore;
9+
import datadog.trace.instrumentation.kafka_common.MetadataState;
10+
import org.apache.kafka.clients.Metadata;
11+
import org.junit.jupiter.api.Test;
12+
13+
class KafkaConsumerInstrumentationHelperTest {
14+
15+
@SuppressWarnings("unchecked")
16+
private final ContextStore<Metadata, MetadataState> metadataContextStore =
17+
mock(ContextStore.class);
18+
19+
@Test
20+
void extractGroupReturnsNullForNullKafkaConsumerInfo() {
21+
assertNull(KafkaConsumerInstrumentationHelper.extractGroup(null));
22+
}
23+
24+
@Test
25+
void extractGroupReturnsNullWhenConsumerGroupIsNull() {
26+
KafkaConsumerInfo kafkaConsumerInfo = new KafkaConsumerInfo(null, null, "localhost:9092");
27+
assertNull(KafkaConsumerInstrumentationHelper.extractGroup(kafkaConsumerInfo));
28+
}
29+
30+
@Test
31+
void extractGroupReturnsConsumerGroupWhenPresent() {
32+
KafkaConsumerInfo kafkaConsumerInfo =
33+
new KafkaConsumerInfo("test-group", null, "localhost:9092");
34+
assertEquals("test-group", KafkaConsumerInstrumentationHelper.extractGroup(kafkaConsumerInfo));
35+
}
36+
37+
@Test
38+
void extractBootstrapServersReturnsNullForNullKafkaConsumerInfo() {
39+
assertNull(KafkaConsumerInstrumentationHelper.extractBootstrapServers(null));
40+
}
41+
42+
@Test
43+
void extractBootstrapServersReturnsNullWhenBootstrapServersIsNull() {
44+
KafkaConsumerInfo kafkaConsumerInfo = new KafkaConsumerInfo("test-group", null, null);
45+
assertNull(KafkaConsumerInstrumentationHelper.extractBootstrapServers(kafkaConsumerInfo));
46+
}
47+
48+
@Test
49+
void extractBootstrapServersReturnsValueWhenPresent() {
50+
KafkaConsumerInfo kafkaConsumerInfo =
51+
new KafkaConsumerInfo("test-group", null, "localhost:9092");
52+
assertEquals(
53+
"localhost:9092",
54+
KafkaConsumerInstrumentationHelper.extractBootstrapServers(kafkaConsumerInfo));
55+
}
56+
57+
@Test
58+
void extractClusterIdReturnsNullForNullKafkaConsumerInfo() {
59+
assertNull(KafkaConsumerInstrumentationHelper.extractClusterId(null, metadataContextStore));
60+
}
61+
62+
@Test
63+
void extractClusterIdReturnsNullWhenMetadataIsNull() {
64+
KafkaConsumerInfo kafkaConsumerInfo = new KafkaConsumerInfo("test-group", "localhost:9092");
65+
assertNull(
66+
KafkaConsumerInstrumentationHelper.extractClusterId(
67+
kafkaConsumerInfo, metadataContextStore));
68+
}
69+
70+
@Test
71+
void extractClusterIdReturnsNullWhenNoStateForMetadata() {
72+
Metadata metadata = mock(Metadata.class);
73+
KafkaConsumerInfo kafkaConsumerInfo =
74+
new KafkaConsumerInfo("test-group", metadata, "localhost:9092");
75+
when(metadataContextStore.get(metadata)).thenReturn(null);
76+
assertNull(
77+
KafkaConsumerInstrumentationHelper.extractClusterId(
78+
kafkaConsumerInfo, metadataContextStore));
79+
}
80+
81+
@Test
82+
void extractClusterIdReturnsClusterIdWhenStatePresent() {
83+
Metadata metadata = mock(Metadata.class);
84+
KafkaConsumerInfo kafkaConsumerInfo =
85+
new KafkaConsumerInfo("test-group", metadata, "localhost:9092");
86+
MetadataState state = new MetadataState();
87+
state.clusterId = "cluster-1";
88+
when(metadataContextStore.get(metadata)).thenReturn(state);
89+
assertEquals(
90+
"cluster-1",
91+
KafkaConsumerInstrumentationHelper.extractClusterId(
92+
kafkaConsumerInfo, metadataContextStore));
93+
}
94+
}

0 commit comments

Comments
 (0)