From 0998a9502ee1d90a7a0492888114b53b900c7900 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Thu, 3 Oct 2024 12:17:25 -0400 Subject: [PATCH] Fix `KafkaDslTests` for splitter race condition It looks like `Stream.generate()` might producer items not in the expected order. So, use `Stream.toList()` to be sure in the sequence size, and then `resequence()` to be sure that items are emitted downstream in the proper sequence order. --- .../springframework/integration/kafka/dsl/KafkaDslTests.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java index a60fe504aa..b03a703637 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java @@ -349,7 +349,8 @@ public class KafkaDslTests { public IntegrationFlow sendToKafkaFlow( KafkaProducerMessageHandlerSpec kafkaMessageHandlerTopic2) { return f -> f - .splitWith(s -> s.function(p -> Stream.generate(() -> p).limit(101).iterator())) + .splitWith(s -> s.function(p -> Stream.generate(() -> p).limit(101).toList())) + .resequence() .enrichHeaders(h -> h.header(KafkaIntegrationHeaders.FUTURE_TOKEN, "foo")) .publishSubscribeChannel(c -> c .subscribe(sf -> sf.handle(