GH-93: Resolve Placeholder in @TopicPartition

Fixes #93

Polishing for test config
This commit is contained in:
Gary Russell
2016-05-27 12:38:46 -04:00
committed by Artem Bilan
parent 6a0813a534
commit be016f0f7d
2 changed files with 11 additions and 4 deletions

View File

@@ -455,7 +455,7 @@ public class KafkaListenerAnnotationBeanPostProcessor<K, V>
List<org.apache.kafka.common.TopicPartition> 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;

View File

@@ -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();