Skip to content

Commit 48ac26b

Browse files
authored
remove flakiness on kafka tests (#11972)
remove flakiness on kafka tests fix codex/autotest issues Co-authored-by: andrea.marziali <andrea.marziali@datadoghq.com>
1 parent 48bc030 commit 48ac26b

5 files changed

Lines changed: 10 additions & 30 deletions

File tree

dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/latestDepTest/groovy/KafkaClientTestBase.groovy

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -165,7 +165,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
165165
container.setupMessageListener(new MessageListener<String, String>() {
166166
@Override
167167
void onMessage(ConsumerRecord<String, String> record) {
168-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
169168
records.add(record)
170169
}
171170
})

dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaClientCustomPropagationConfigTest.groovy

Lines changed: 0 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -97,31 +97,27 @@ class KafkaClientCustomPropagationConfigTest extends InstrumentationSpecificatio
9797
container1.setupMessageListener(new MessageListener<String, String>() {
9898
@Override
9999
void onMessage(ConsumerRecord<String, String> record) {
100-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
101100
records1.add(record)
102101
}
103102
})
104103

105104
container2.setupMessageListener(new MessageListener<String, String>() {
106105
@Override
107106
void onMessage(ConsumerRecord<String, String> record) {
108-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
109107
records2.add(record)
110108
}
111109
})
112110

113111
container3.setupMessageListener(new MessageListener<String, String>() {
114112
@Override
115113
void onMessage(ConsumerRecord<String, String> record) {
116-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
117114
records3.add(record)
118115
}
119116
})
120117

121118
container4.setupMessageListener(new MessageListener<String, String>() {
122119
@Override
123120
void onMessage(ConsumerRecord<String, String> record) {
124-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
125121
records4.add(record)
126122
}
127123
})
@@ -202,31 +198,27 @@ class KafkaClientCustomPropagationConfigTest extends InstrumentationSpecificatio
202198
container1.setupMessageListener(new MessageListener<String, String>() {
203199
@Override
204200
void onMessage(ConsumerRecord<String, String> record) {
205-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
206201
records1.add(activeSpan())
207202
}
208203
})
209204

210205
container2.setupMessageListener(new MessageListener<String, String>() {
211206
@Override
212207
void onMessage(ConsumerRecord<String, String> record) {
213-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
214208
records2.add(activeSpan())
215209
}
216210
})
217211

218212
container3.setupMessageListener(new MessageListener<String, String>() {
219213
@Override
220214
void onMessage(ConsumerRecord<String, String> record) {
221-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
222215
records3.add(activeSpan())
223216
}
224217
})
225218

226219
container4.setupMessageListener(new MessageListener<String, String>() {
227220
@Override
228221
void onMessage(ConsumerRecord<String, String> record) {
229-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
230222
records4.add(activeSpan())
231223
}
232224
})

dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaClientTestBase.groovy

Lines changed: 5 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -252,7 +252,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
252252
container.setupMessageListener(new MessageListener<String, String>() {
253253
@Override
254254
void onMessage(ConsumerRecord<String, String> record) {
255-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
256255
records.add(record)
257256
}
258257
})
@@ -420,7 +419,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
420419
container.setupMessageListener(new MessageListener<String, String>() {
421420
@Override
422421
void onMessage(ConsumerRecord<String, String> record) {
423-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
424422
records.add(record)
425423
}
426424
})
@@ -471,10 +469,13 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
471469
}
472470
}
473471

472+
// sort a snapshot so the producer trace is deterministically first, regardless of write order
473+
def sortedTraces = new ArrayList<>(TEST_WRITER)
474+
sortedTraces.sort(SORT_TRACES_BY_ID)
474475
def headers = received.headers()
475476
headers.iterator().hasNext()
476-
new String(headers.headers("x-datadog-trace-id").iterator().next().value()) == "${TEST_WRITER[0][2].traceId}"
477-
new String(headers.headers("x-datadog-parent-id").iterator().next().value()) == "${TEST_WRITER[0][2].spanId}"
477+
new String(headers.headers("x-datadog-trace-id").iterator().next().value()) == "${sortedTraces[0][2].traceId}"
478+
new String(headers.headers("x-datadog-parent-id").iterator().next().value()) == "${sortedTraces[0][2].spanId}"
478479

479480
if (isDataStreamsEnabled()) {
480481
StatsGroup first = TEST_DATA_STREAMS_WRITER.groups.find { it.parentHash == 0 }
@@ -553,7 +554,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
553554
container.setupMessageListener(new MessageListener<String, String>() {
554555
@Override
555556
void onMessage(ConsumerRecord<String, String> record) {
556-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
557557
records.add(record)
558558
}
559559
})
@@ -889,7 +889,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
889889
container.setupMessageListener(new BatchMessageListener<String, String>() {
890890
@Override
891891
void onMessage(List<ConsumerRecord<String, String>> consumerRecords) {
892-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
893892
consumerRecords.each {
894893
records.add(it)
895894
}
@@ -1026,7 +1025,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
10261025
container.setupMessageListener(new MessageListener<String, String>() {
10271026
@Override
10281027
void onMessage(ConsumerRecord<String, String> record) {
1029-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
10301028
records.add(record)
10311029
if (isDataStreamsEnabled()) {
10321030
// even if header propagation is disabled, we want data streams to work.

dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/test/groovy/KafkaClientCustomPropagationConfigTest.groovy

Lines changed: 0 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -101,31 +101,27 @@ class KafkaClientCustomPropagationConfigTest extends InstrumentationSpecificatio
101101
container1.setupMessageListener(new MessageListener<String, String>() {
102102
@Override
103103
void onMessage(ConsumerRecord<String, String> record) {
104-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
105104
records1.add(record)
106105
}
107106
})
108107

109108
container2.setupMessageListener(new MessageListener<String, String>() {
110109
@Override
111110
void onMessage(ConsumerRecord<String, String> record) {
112-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
113111
records2.add(record)
114112
}
115113
})
116114

117115
container3.setupMessageListener(new MessageListener<String, String>() {
118116
@Override
119117
void onMessage(ConsumerRecord<String, String> record) {
120-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
121118
records3.add(record)
122119
}
123120
})
124121

125122
container4.setupMessageListener(new MessageListener<String, String>() {
126123
@Override
127124
void onMessage(ConsumerRecord<String, String> record) {
128-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
129125
records4.add(record)
130126
}
131127
})
@@ -205,31 +201,27 @@ class KafkaClientCustomPropagationConfigTest extends InstrumentationSpecificatio
205201
container1.setupMessageListener(new MessageListener<String, String>() {
206202
@Override
207203
void onMessage(ConsumerRecord<String, String> record) {
208-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
209204
records1.add(activeSpan())
210205
}
211206
})
212207

213208
container2.setupMessageListener(new MessageListener<String, String>() {
214209
@Override
215210
void onMessage(ConsumerRecord<String, String> record) {
216-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
217211
records2.add(activeSpan())
218212
}
219213
})
220214

221215
container3.setupMessageListener(new MessageListener<String, String>() {
222216
@Override
223217
void onMessage(ConsumerRecord<String, String> record) {
224-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
225218
records3.add(activeSpan())
226219
}
227220
})
228221

229222
container4.setupMessageListener(new MessageListener<String, String>() {
230223
@Override
231224
void onMessage(ConsumerRecord<String, String> record) {
232-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
233225
records4.add(activeSpan())
234226
}
235227
})

dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/test/groovy/KafkaClientTestBase.groovy

Lines changed: 5 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -190,7 +190,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
190190
container.setupMessageListener(new MessageListener<String, String>() {
191191
@Override
192192
void onMessage(ConsumerRecord<String, String> record) {
193-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
194193
records.add(record)
195194
}
196195
})
@@ -349,7 +348,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
349348
container.setupMessageListener(new MessageListener<String, String>() {
350349
@Override
351350
void onMessage(ConsumerRecord<String, String> record) {
352-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
353351
records.add(record)
354352
}
355353
})
@@ -404,10 +402,13 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
404402
}
405403
}
406404

405+
// sort a snapshot so the producer trace is deterministically first, regardless of write order
406+
def sortedTraces = new ArrayList<>(TEST_WRITER)
407+
sortedTraces.sort(SORT_TRACES_BY_ID)
407408
def headers = received.headers()
408409
headers.iterator().hasNext()
409-
new String(headers.headers("x-datadog-trace-id").iterator().next().value()) == "${TEST_WRITER[0][2].traceId}"
410-
new String(headers.headers("x-datadog-parent-id").iterator().next().value()) == "${TEST_WRITER[0][2].spanId}"
410+
new String(headers.headers("x-datadog-trace-id").iterator().next().value()) == "${sortedTraces[0][2].traceId}"
411+
new String(headers.headers("x-datadog-parent-id").iterator().next().value()) == "${sortedTraces[0][2].spanId}"
411412

412413
if (isDataStreamsEnabled()) {
413414
StatsGroup first = TEST_DATA_STREAMS_WRITER.groups.find { it.parentHash == 0 }
@@ -480,7 +481,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
480481
container.setupMessageListener(new MessageListener<String, String>() {
481482
@Override
482483
void onMessage(ConsumerRecord<String, String> record) {
483-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
484484
records.add(record)
485485
}
486486
})
@@ -824,7 +824,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
824824
container.setupMessageListener(new MessageListener<String, String>() {
825825
@Override
826826
void onMessage(ConsumerRecord<String, String> record) {
827-
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
828827
records.add(record)
829828
if (isDataStreamsEnabled()) {
830829
// even if header propagation is disabled, we want data streams to work.

0 commit comments

Comments
 (0)