From 58ff34d83b14ea5fceb5cce442c8a5ef7d8a78db Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 1 Jul 2014 15:01:22 +0300 Subject: [PATCH] Kafka: KCCtx: return null if consumedData.isEmpty Fixes: https://github.com/spring-projects/spring-integration-extensions/issues/68 --- .../kafka/support/KafkaConsumerContext.java | 15 ++++++++------- 1 file changed, 8 insertions(+), 7 deletions(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java index 5a0f4f2818..e019da4659 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/support/KafkaConsumerContext.java @@ -15,6 +15,11 @@ */ package org.springframework.integration.kafka.support; +import java.util.Collection; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.BeanFactoryAware; @@ -22,11 +27,7 @@ import org.springframework.beans.factory.ListableBeanFactory; import org.springframework.integration.kafka.core.KafkaConsumerDefaults; import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; - -import java.util.Collection; -import java.util.HashMap; -import java.util.List; -import java.util.Map; +import org.springframework.util.CollectionUtils; /** * @author Soby Chacko @@ -54,11 +55,11 @@ public class KafkaConsumerContext implements BeanFactoryAware { for (final ConsumerConfiguration consumerConfiguration : getConsumerConfigurations()) { final Map>> messages = consumerConfiguration.receive(); - if (messages != null){ + if (!CollectionUtils.isEmpty(messages)){ consumedData.putAll(messages); } } - return MessageBuilder.withPayload(consumedData).build(); + return consumedData.isEmpty() ? null : MessageBuilder.withPayload(consumedData).build(); } public String getConsumerTimeout() {