diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/InboundGatewayTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/InboundGatewayTests.java index aabf054309..29f8155914 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/InboundGatewayTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/InboundGatewayTests.java @@ -189,6 +189,8 @@ class InboundGatewayTests { assertThat(gateway.isPaused()).isFalse(); gateway.stop(); + consumer.close(); + pf.reset(); } @Test @@ -268,6 +270,8 @@ class InboundGatewayTests { assertThat(record).has(value("ERROR")); gateway.stop(); + consumer.close(); + pf.reset(); } @Test @@ -352,6 +356,8 @@ class InboundGatewayTests { assertThat(record).has(value("ERROR")); gateway.stop(); + consumer.close(); + pf.reset(); } @Test @@ -409,6 +415,7 @@ class InboundGatewayTests { gateway.stop(); consumer.close(); + pf.reset(); } } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java index 604a3369f2..44e4200a5a 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java @@ -215,6 +215,7 @@ class MessageDrivenAdapterTests { assertThat(((ConversionException) error.getPayload()).getRecord()).isNotNull(); adapter.stop(); + pf.reset(); } @Test @@ -277,6 +278,7 @@ class MessageDrivenAdapterTests { assertThat(receivedMessageHistory.get().toString()).isEqualTo("myNullChannel"); adapter.stop(); + pf.reset(); } @@ -387,6 +389,7 @@ class MessageDrivenAdapterTests { assertThat(StaticMessageHeaderAccessor.getDeliveryAttempt(originalMessage).get()).isEqualTo(1); adapter.stop(); + pf.reset(); } @Test @@ -474,7 +477,9 @@ class MessageDrivenAdapterTests { assertThat(((ConversionException) error.getPayload()).getMessage()) .contains("Failed to convert to message"); assertThat(((ConversionException) error.getPayload()).getRecords()).hasSize(2); + adapter.stop(); + pf.reset(); } @Test @@ -517,6 +522,7 @@ class MessageDrivenAdapterTests { assertThat(received.getPayload()).isInstanceOf(Map.class); adapter.stop(); + pf.reset(); } @Test @@ -564,6 +570,7 @@ class MessageDrivenAdapterTests { assertThat(received.getPayload()).isEqualTo(new Foo("baz")); adapter.stop(); + pf.reset(); } @SuppressWarnings({ "unchecked", "rawtypes" }) diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageSourceIntegrationTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageSourceIntegrationTests.java index 0360c1136b..2db404a6f2 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageSourceIntegrationTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageSourceIntegrationTests.java @@ -119,7 +119,7 @@ class MessageSourceIntegrationTests { assertThat(received).isNull(); assertThat(KafkaTestUtils.getPropertyValue(source, "consumer.fetcher.minBytes")).isEqualTo(2); source.destroy(); - producerFactory.destroy(); + template.destroy(); } }