diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java index 0a7f9d26..4318cd52 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/KafkaTemplate.java @@ -31,6 +31,7 @@ import org.springframework.kafka.support.SendResult; import org.springframework.kafka.support.converter.MessageConverter; import org.springframework.kafka.support.converter.MessagingMessageConverter; import org.springframework.messaging.Message; +import org.springframework.util.Assert; import org.springframework.util.concurrent.ListenableFuture; import org.springframework.util.concurrent.SettableListenableFuture; @@ -43,6 +44,7 @@ import org.springframework.util.concurrent.SettableListenableFuture; * * @author Marius Bogoevici * @author Gary Russell + * @author Igor Stepanov */ public class KafkaTemplate implements KafkaOperations { @@ -50,9 +52,10 @@ public class KafkaTemplate implements KafkaOperations { private final ProducerFactory producerFactory; + private final boolean autoFlush; + private MessageConverter messageConverter = new MessagingMessageConverter(); -private final boolean autoFlush; private volatile Producer producer; private volatile String defaultTopic; @@ -172,6 +175,7 @@ private final boolean autoFlush; @Override public void flush() { + Assert.state(this.producer != null, "'producer' must not be null for flushing."); this.producer.flush(); } diff --git a/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTests.java b/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTests.java index 074966fe..760c13cf 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTests.java @@ -17,6 +17,8 @@ package org.springframework.kafka.core; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.fail; +import static org.mockito.Mockito.mock; import static org.springframework.kafka.test.assertj.KafkaConditions.key; import static org.springframework.kafka.test.assertj.KafkaConditions.partition; import static org.springframework.kafka.test.assertj.KafkaConditions.value; @@ -50,7 +52,7 @@ import org.springframework.util.concurrent.ListenableFutureCallback; /** * @author Gary Russell * @author Artem Bilan - * + * @author Igor Stepanov */ public class KafkaTemplateTests { @@ -195,4 +197,18 @@ public class KafkaTemplateTests { pf.createProducer().close(); } + @Test + @SuppressWarnings({"rawtypes", "unchecked"}) + public void flushWithoutSend() throws Exception { + KafkaTemplate template = new KafkaTemplate(mock(ProducerFactory.class)); + try { + template.flush(); + fail("IllegalStateException expected"); + } + catch (Exception e) { + assertThat(e).isInstanceOf(IllegalStateException.class); + assertThat(e.getMessage()).isEqualTo("'producer' must not be null for flushing."); + } + } + }