From be016f0f7d1e74f3e0904ee35fd65f01fa877c86 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Fri, 27 May 2016 12:38:46 -0400 Subject: [PATCH] GH-93: Resolve Placeholder in @TopicPartition Fixes #93 Polishing for test config --- .../KafkaListenerAnnotationBeanPostProcessor.java | 2 +- .../annotation/EnableKafkaIntegrationTests.java | 13 ++++++++++--- 2 files changed, 11 insertions(+), 4 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java index 2e95f5b4..5cdfc198 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java @@ -455,7 +455,7 @@ public class KafkaListenerAnnotationBeanPostProcessor List result = new ArrayList<>(); if (partitions.length > 0) { for (int i = 0; i < partitions.length; i++) { - resolvePartitionAsInteger((String) topic, partitions[i], result); + resolvePartitionAsInteger((String) topic, resolveExpression(partitions[i]), result); } } return result; diff --git a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java index cd6faa86..75d4063d 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java @@ -36,6 +36,7 @@ import org.mockito.Mockito; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.context.support.PropertySourcesPlaceholderConfigurer; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.config.KafkaListenerContainerFactory; import org.springframework.kafka.config.KafkaListenerEndpointRegistry; @@ -189,6 +190,11 @@ public class EnableKafkaIntegrationTests { @EnableTransactionManagement(proxyTargetClass = true) public static class Config { + @Bean + public static PropertySourcesPlaceholderConfigurer ppc() { + return new PropertySourcesPlaceholderConfigurer(); + } + @Bean public PlatformTransactionManager transactionManager() { return Mockito.mock(PlatformTransactionManager.class); @@ -358,12 +364,12 @@ public class EnableKafkaIntegrationTests { public void manualStart(String foo) { } - @KafkaListener(id = "foo", topics = "annotated1") + @KafkaListener(id = "foo", topics = "${topicOne:annotated1}") public void listen1(String foo) { this.latch1.countDown(); } - @KafkaListener(id = "bar", topicPattern = "annotated2") + @KafkaListener(id = "bar", topicPattern = "${topicTwo:annotated2}") public void listen2(@Payload String foo, @Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) Integer key, @Header(KafkaHeaders.RECEIVED_PARTITION_ID) int partition, @@ -374,7 +380,8 @@ public class EnableKafkaIntegrationTests { this.latch2.countDown(); } - @KafkaListener(id = "baz", topicPartitions = @TopicPartition(topic = "annotated3", partitions = "0")) + @KafkaListener(id = "baz", topicPartitions = @TopicPartition(topic = "${topicThree:annotated3}", + partitions = "${zero:0}")) public void listen3(ConsumerRecord record) { this.record = record; this.latch3.countDown();