diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java index e4162f35..3ee8aaca 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java @@ -20,6 +20,7 @@ import java.util.HashMap; import java.util.Map; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; /** @@ -41,7 +42,7 @@ public class DefaultKafkaConsumerFactory implements ConsumerFactory @Override public boolean isAutoCommit() { - Object auto = this.configs.get("enable.auto.commit"); + Object auto = this.configs.get(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG); return auto instanceof Boolean ? (Boolean) auto : auto instanceof String ? Boolean.valueOf((String) auto) : false; } diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java index ba95afbf..419f3f39 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java @@ -37,6 +37,7 @@ import org.apache.commons.logging.LogFactory; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.common.TopicPartition; import org.junit.ClassRule; import org.junit.Test; @@ -57,6 +58,7 @@ import org.springframework.kafka.test.utils.KafkaTestUtils; /** * @author Gary Russell + * @author Artem Bilan * */ public class ConcurrentMessageListenerContainerTests { @@ -150,9 +152,32 @@ public class ConcurrentMessageListenerContainerTests { @Test public void testDefinedPartitions() throws Exception { logger.info("Start auto parts"); - Map props = KafkaTestUtils.consumerProps("test3", "true", embeddedKafka); + final Map props = KafkaTestUtils.consumerProps("test3", "true", embeddedKafka); TopicPartition topic1Partition0 = new TopicPartition(topic3, 0); - DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(props); + + final CountDownLatch initialConsumersLatch = new CountDownLatch(2); + + DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory(props) { + + @Override + public Consumer createConsumer() { + return new KafkaConsumer(props) { + + @Override + public ConsumerRecords poll(long timeout) { + try { + return super.poll(timeout); + } + finally { + initialConsumersLatch.countDown(); + } + } + + }; + } + + }; + ConcurrentMessageListenerContainer container1 = new ConcurrentMessageListenerContainer<>(cf, topic1Partition0); final CountDownLatch latch1 = new CountDownLatch(2); @@ -163,10 +188,11 @@ public class ConcurrentMessageListenerContainerTests { logger.info("auto part: " + message); latch1.countDown(); } + }); container1.setBeanName("b1"); container1.start(); - Thread.sleep(1000); + TopicPartition topic1Partition1 = new TopicPartition(topic3, 1); ConcurrentMessageListenerContainer container2 = new ConcurrentMessageListenerContainer<>(cf, topic1Partition1); @@ -178,10 +204,13 @@ public class ConcurrentMessageListenerContainerTests { logger.info("auto part: " + message); latch2.countDown(); } + }); container2.setBeanName("b2"); container2.start(); - Thread.sleep(1000); + + assertTrue(initialConsumersLatch.await(20, TimeUnit.SECONDS)); + Map senderProps = KafkaTestUtils.senderProps(embeddedKafka); ProducerFactory pf = new DefaultKafkaProducerFactory(senderProps); KafkaTemplate template = new KafkaTemplate<>(pf); @@ -191,7 +220,9 @@ public class ConcurrentMessageListenerContainerTests { template.convertAndSend(0, "baz"); template.convertAndSend(2, "qux"); template.flush(); + assertTrue(latch1.await(60, TimeUnit.SECONDS)); + assertTrue(latch2.await(60, TimeUnit.SECONDS)); container1.stop(); container2.stop(); diff --git a/spring-kafka/src/test/resources/log4j.properties b/spring-kafka/src/test/resources/log4j.properties new file mode 100644 index 00000000..b4e3588b --- /dev/null +++ b/spring-kafka/src/test/resources/log4j.properties @@ -0,0 +1,8 @@ +log4j.rootCategory=WARN, stdout + +log4j.appender.stdout=org.apache.log4j.ConsoleAppender +log4j.appender.stdout.layout=org.apache.log4j.PatternLayout +log4j.appender.stdout.layout.ConversionPattern=%d{HH:mm:ss.SSS} %-5p [%t][%c] %m%n + +log4j.category.org.springframework.kafka=TRACE +log4j.category.org.apache.kafka.clients=WARN