From 9f1e243267699ac3c4813fc9e75433b080f200db Mon Sep 17 00:00:00 2001 From: Andrea Marziali Date: Tue, 15 Sep 2026 11:11:27 +0200 Subject: [PATCH 1/3] Wait for Reactor Kafka test pipelines to finish --- .../src/test/groovy/KafkaReactorForkedTest.groovy | 14 ++++++++++---- 1 file changed, 10 insertions(+), 4 deletions(-) diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaReactorForkedTest.groovy b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaReactorForkedTest.groovy index 49ec8900712..86bee90f840 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaReactorForkedTest.groovy +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaReactorForkedTest.groovy @@ -129,8 +129,9 @@ class KafkaReactorForkedTest extends InstrumentationSpecification { } }) - // create a thread safe queue to store the received message - kafkaReceiver.receive() + // Complete after the expected records and their asynchronous commits have drained. + def receiverCompletion = kafkaReceiver.receive() + .take(100) // publish on another thread to be sure we're propagating that receive span correctly .publishOn(Schedulers.parallel()) .flatMap { @@ -139,7 +140,8 @@ class KafkaReactorForkedTest extends InstrumentationSpecification { } } .subscribeOn(Schedulers.parallel()) - .subscribe() + .then() + .toFuture() // wait until the container has the required number of assigned partitions @@ -152,7 +154,8 @@ class KafkaReactorForkedTest extends InstrumentationSpecification { kafkaSender.send(Mono.just(SenderRecord.create(new ProducerRecord<>(KafkaClientTestBase.SHARED_TOPIC, greeting), null))) } .publishOn(Schedulers.parallel()) - .subscribe() + .blockLast() + receiverCompletion.get(30, TimeUnit.SECONDS) then: // check that the all the consume (100) and the send (100) are reported TEST_WRITER.waitForTraces(200) @@ -174,6 +177,9 @@ class KafkaReactorForkedTest extends InstrumentationSpecification { assert it.get(consumeIndex).getParentId() == it.get(produceIndex).getSpanId() assert it.get(produceIndex).getParentId() == 0 } + cleanup: + receiverCompletion?.cancel(true) + kafkaSender?.close() } def producerSpan( From 297eee1cdda0600c04ec813a63e549fbb219711c Mon Sep 17 00:00:00 2001 From: Andrea Marziali Date: Tue, 15 Sep 2026 13:10:59 +0200 Subject: [PATCH 2/3] Clarify Reactor Kafka completion ordering --- .../src/test/groovy/KafkaReactorForkedTest.groovy | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaReactorForkedTest.groovy b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaReactorForkedTest.groovy index 86bee90f840..c88416d5a1e 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaReactorForkedTest.groovy +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaReactorForkedTest.groovy @@ -129,7 +129,7 @@ class KafkaReactorForkedTest extends InstrumentationSpecification { } }) - // Complete after the expected records and their asynchronous commits have drained. + // Limit reception to 100 records; then() waits for all their asynchronous commits. def receiverCompletion = kafkaReceiver.receive() .take(100) // publish on another thread to be sure we're propagating that receive span correctly From 8281aadd6791bf52974926e49a3bbc7377fd6eb6 Mon Sep 17 00:00:00 2001 From: Andrea Marziali Date: Tue, 15 Sep 2026 13:36:29 +0200 Subject: [PATCH 3/3] Wait for Reactor Kafka commits before stopping receiver --- .../test/groovy/KafkaReactorForkedTest.groovy | 24 +++++++++---------- 1 file changed, 11 insertions(+), 13 deletions(-) diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaReactorForkedTest.groovy b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaReactorForkedTest.groovy index c88416d5a1e..08ec9fb0513 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaReactorForkedTest.groovy +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaReactorForkedTest.groovy @@ -129,19 +129,17 @@ class KafkaReactorForkedTest extends InstrumentationSpecification { } }) - // Limit reception to 100 records; then() waits for all their asynchronous commits. - def receiverCompletion = kafkaReceiver.receive() - .take(100) + def commitsCompleted = new CountDownLatch(100) + def receiverSubscription = kafkaReceiver.receive() // publish on another thread to be sure we're propagating that receive span correctly .publishOn(Schedulers.parallel()) - .flatMap { - receiverRecord -> { - receiverRecord.receiverOffset().commit() - } + .flatMap { receiverRecord -> + receiverRecord.receiverOffset().commit().then(Mono.just(receiverRecord)) } .subscribeOn(Schedulers.parallel()) - .then() - .toFuture() + .subscribe { + commitsCompleted.countDown() + } // wait until the container has the required number of assigned partitions @@ -155,12 +153,12 @@ class KafkaReactorForkedTest extends InstrumentationSpecification { } .publishOn(Schedulers.parallel()) .blockLast() - receiverCompletion.get(30, TimeUnit.SECONDS) + assert commitsCompleted.await(30, TimeUnit.SECONDS) + receiverSubscription.dispose() then: // check that the all the consume (100) and the send (100) are reported TEST_WRITER.waitForTraces(200) - Map> traces = TEST_WRITER.inject([:]) { - map, entry -> + Map> traces = TEST_WRITER.inject([:]) { map, entry -> def key = entry.get(0).getTraceId().toString() map[key] = (map[key] ?: []) + entry return map @@ -178,7 +176,7 @@ class KafkaReactorForkedTest extends InstrumentationSpecification { assert it.get(produceIndex).getParentId() == 0 } cleanup: - receiverCompletion?.cancel(true) + receiverSubscription?.dispose() kafkaSender?.close() }