diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java index ece017ee5a..549ab3fecc 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java @@ -99,6 +99,8 @@ public class ReactiveStreamsConsumerTests { assertThat(stopLatch.await(10, TimeUnit.SECONDS)).isTrue(); assertThat(result).containsExactly(testMessage, testMessage2); + + reactiveConsumer.stop(); } @@ -222,6 +224,8 @@ public class ReactiveStreamsConsumerTests { verify(testSubscriber, never()).onComplete(); assertThat(messages.isEmpty()).isTrue(); + + reactiveConsumer.stop(); } @Test @@ -264,6 +268,8 @@ public class ReactiveStreamsConsumerTests { assertThat(stopLatch.await(10, TimeUnit.SECONDS)).isTrue(); assertThat(result.size()).isEqualTo(3); assertThat(result).containsExactly(testMessage, testMessage2, testMessage2); + + endpointFactoryBean.stop(); } @Test @@ -302,7 +308,8 @@ public class ReactiveStreamsConsumerTests { .expectNext(testMessage, testMessage2) .thenCancel() .verify(); + + reactiveConsumer.stop(); } - }