diff --git a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/integration/MirrorConnectorsIntegrationBaseTest.java b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/integration/MirrorConnectorsIntegrationBaseTest.java index 56d2bf4974093..3ad8897b12d94 100644 --- a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/integration/MirrorConnectorsIntegrationBaseTest.java +++ b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/integration/MirrorConnectorsIntegrationBaseTest.java @@ -1050,8 +1050,11 @@ private void testConnectorMetrics(String format, Supplier assertions) t String topic = "test-topic-metrics"; primary.kafka().createTopic(topic); try (KafkaProducer producer = primary.kafka().createProducer(Map.of())) { - for (int i = 0; i < NUM_RECORDS_PRODUCED; i++) { - producer.send(new ProducerRecord<>(topic, ("value" + i).getBytes())).get(); + var futures = IntStream.range(0, NUM_RECORDS_PRODUCED) + .mapToObj(i -> producer.send(new ProducerRecord(topic, ("value" + i).getBytes()))) + .toList(); + for (var future : futures) { + future.get(); } }