From 2cae7420a3fb5a64c19b36079b74e77b6c2b1846 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Fri, 16 Jan 2015 13:59:39 -0500 Subject: [PATCH] Fixes a race condition in tests https://build.spring.io/browse/INTEXT-KAFKA-JOB1-39 --- .../KafkaMessageDrivenChannelAdapterTests.java | 15 ++++++++------- ...rivenChannelAdapterWithSpecialOffsetTests.java | 5 +++-- ...eDrivenChannelAdapterWithWrongOffsetTests.java | 10 ++++++---- 3 files changed, 17 insertions(+), 13 deletions(-) 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