diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java index 46b009c368..21a32be707 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java @@ -115,6 +115,8 @@ public class ReactiveStreamsTests { .take(6) .subscribe(); + final CountDownLatch asyncSetupLatch = new CountDownLatch(1); + Future> future = Executors.newSingleThreadExecutor().submit(() -> Flux.just("11,12,13") @@ -126,7 +128,9 @@ public class ReactiveStreamsTests { .map(Message::getPayload) .take(7) .collectList() - .block(Duration.ofSeconds(10))); + .block(getDuration(asyncSetupLatch))); + + assertTrue(asyncSetupLatch.await(10, TimeUnit.SECONDS)); this.inputChannel.send(new GenericMessage<>("6,7,8,9,10")); @@ -136,6 +140,11 @@ public class ReactiveStreamsTests { assertEquals(7, integers.size()); } + private Duration getDuration(CountDownLatch asyncSetupLatch) { + asyncSetupLatch.countDown(); + return Duration.ofSeconds(10); + } + @Configuration @EnableIntegration public static class ContextConfiguration {