diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterTests.java index 3760444eb8..1723ab2067 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterTests.java @@ -89,15 +89,16 @@ public class KafkaMessageDrivenChannelAdapterTests extends AbstractMessageListen kafkaMessageDrivenChannelAdapter.setOutputChannel(new MessageChannel() { @Override public boolean send(Message message) { - latch.countDown(); - return receivedData.put( - (Integer)message.getHeaders().get(KafkaHeaders.PARTITION_ID), + boolean addedSuccessfully = receivedData.put( + (Integer) message.getHeaders().get(KafkaHeaders.PARTITION_ID), new KeyedMessageWithOffset( - (String)message.getHeaders().get(KafkaHeaders.MESSAGE_KEY), - (String)message.getPayload(), - (Long)message.getHeaders().get(KafkaHeaders.OFFSET), + (String) message.getHeaders().get(KafkaHeaders.MESSAGE_KEY), + (String) message.getPayload(), + (Long) message.getHeaders().get(KafkaHeaders.OFFSET), Thread.currentThread().getName(), - (Integer)message.getHeaders().get(KafkaHeaders.PARTITION_ID))); + (Integer) message.getHeaders().get(KafkaHeaders.PARTITION_ID))); + latch.countDown(); + return addedSuccessfully; } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterWithSpecialOffsetTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterWithSpecialOffsetTests.java index e8665f9073..c18ca4b05f 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterWithSpecialOffsetTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterWithSpecialOffsetTests.java @@ -110,8 +110,7 @@ public class KafkaMessageDrivenChannelAdapterWithSpecialOffsetTests extends Abst kafkaMessageDrivenChannelAdapter.setOutputChannel(new MessageChannel() { @Override public boolean send(Message message) { - latch.countDown(); - return receivedData.put( + boolean addedSuccessfully = receivedData.put( (Integer) message.getHeaders().get(KafkaHeaders.PARTITION_ID), new KeyedMessageWithOffset( (String) message.getHeaders().get(KafkaHeaders.MESSAGE_KEY), @@ -119,6 +118,8 @@ public class KafkaMessageDrivenChannelAdapterWithSpecialOffsetTests extends Abst (Long) message.getHeaders().get(KafkaHeaders.OFFSET), Thread.currentThread().getName(), (Integer) message.getHeaders().get(KafkaHeaders.PARTITION_ID))); + latch.countDown(); + return addedSuccessfully; } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterWithWrongOffsetTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterWithWrongOffsetTests.java index ad85926afd..caeca2a2fa 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterWithWrongOffsetTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterWithWrongOffsetTests.java @@ -113,8 +113,7 @@ public class KafkaMessageDrivenChannelAdapterWithWrongOffsetTests extends Abstra kafkaMessageDrivenChannelAdapter.setOutputChannel(new MessageChannel() { @Override public boolean send(Message message) { - latch.countDown(); - return receivedData.put( + boolean addedSuccessfully = receivedData.put( (Integer) message.getHeaders().get(KafkaHeaders.PARTITION_ID), new KeyedMessageWithOffset( (String) message.getHeaders().get(KafkaHeaders.MESSAGE_KEY), @@ -122,6 +121,8 @@ public class KafkaMessageDrivenChannelAdapterWithWrongOffsetTests extends Abstra (Long) message.getHeaders().get(KafkaHeaders.OFFSET), Thread.currentThread().getName(), (Integer) message.getHeaders().get(KafkaHeaders.PARTITION_ID))); + latch.countDown(); + return addedSuccessfully; } @@ -254,8 +255,7 @@ public class KafkaMessageDrivenChannelAdapterWithWrongOffsetTests extends Abstra kafkaMessageDrivenChannelAdapter.setOutputChannel(new MessageChannel() { @Override public boolean send(Message message) { - latch.countDown(); - return receivedData.put( + boolean addedSuccessfully = receivedData.put( (Integer) message.getHeaders().get(KafkaHeaders.PARTITION_ID), new KeyedMessageWithOffset( (String) message.getHeaders().get(KafkaHeaders.MESSAGE_KEY), @@ -263,6 +263,8 @@ public class KafkaMessageDrivenChannelAdapterWithWrongOffsetTests extends Abstra (Long) message.getHeaders().get(KafkaHeaders.OFFSET), Thread.currentThread().getName(), (Integer) message.getHeaders().get(KafkaHeaders.PARTITION_ID))); + latch.countDown(); + return addedSuccessfully; } @Override