From bd5cc2f89d42fd72bd5a4a68ca25a5fce83306ba Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Tue, 26 Jun 2018 12:22:58 -0400 Subject: [PATCH] GH-398 added support for container customization PLease see https://github.com/spring-cloud/spring-cloud-stream-binder-rabbit/issues/139 for more details Resolves #398 Added test and polishing --- .../kafka/KafkaMessageChannelBinder.java | 10 +++++++-- .../config/KafkaBinderConfiguration.java | 7 ++++-- .../integration/KafkaBinderActuatorTests.java | 22 +++++++++++++++++++ 3 files changed, 35 insertions(+), 4 deletions(-) diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index 9ea126b76..cbd2fbcef 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -67,6 +67,7 @@ import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerPro import org.springframework.cloud.stream.binder.kafka.properties.KafkaExtendedBindingProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaProducerProperties; import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; +import org.springframework.cloud.stream.config.ListenerContainerCustomizer; import org.springframework.cloud.stream.provisioning.ConsumerDestination; import org.springframework.cloud.stream.provisioning.ProducerDestination; import org.springframework.context.Lifecycle; @@ -146,9 +147,13 @@ public class KafkaMessageChannelBinder extends private KafkaExtendedBindingProperties extendedBindingProperties = new KafkaExtendedBindingProperties(); + public KafkaMessageChannelBinder(KafkaBinderConfigurationProperties configurationProperties, KafkaTopicProvisioner provisioningProvider) { + this(configurationProperties, provisioningProvider, null); + } + public KafkaMessageChannelBinder(KafkaBinderConfigurationProperties configurationProperties, - KafkaTopicProvisioner provisioningProvider) { - super(headersToMap(configurationProperties), provisioningProvider); + KafkaTopicProvisioner provisioningProvider, ListenerContainerCustomizer> containerCustomizer) { + super(headersToMap(configurationProperties), provisioningProvider, containerCustomizer); this.configurationProperties = configurationProperties; if (StringUtils.hasText(configurationProperties.getTransaction().getTransactionIdPrefix())) { this.transactionManager = new KafkaTransactionManager<>( @@ -407,6 +412,7 @@ public class KafkaMessageChannelBinder extends this.logger.debug( "Listened partitions: " + StringUtils.collectionToCommaDelimitedString(listenedPartitions)); } + this.getContainerCustomizer().configure(messageListenerContainer, destination.getName(), group); final KafkaMessageDrivenChannelAdapter kafkaMessageDrivenChannelAdapter = new KafkaMessageDrivenChannelAdapter<>(messageListenerContainer); kafkaMessageDrivenChannelAdapter.setMessageConverter(getMessageConverter(extendedConsumerProperties)); diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java index 625bfddea..3772829b3 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java @@ -36,12 +36,15 @@ import org.springframework.cloud.stream.binder.kafka.properties.JaasLoginModuleC import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaExtendedBindingProperties; import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; +import org.springframework.cloud.stream.config.ListenerContainerCustomizer; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; +import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.security.jaas.KafkaJaasLoginModuleInitializer; import org.springframework.kafka.support.LoggingProducerListener; import org.springframework.kafka.support.ProducerListener; +import org.springframework.lang.Nullable; /** * @author David Turanski @@ -81,10 +84,10 @@ public class KafkaBinderConfiguration { @Bean KafkaMessageChannelBinder kafkaMessageChannelBinder(KafkaBinderConfigurationProperties configurationProperties, - KafkaTopicProvisioner provisioningProvider) { + KafkaTopicProvisioner provisioningProvider, @Nullable ListenerContainerCustomizer> listenerContainerCustomizer) { KafkaMessageChannelBinder kafkaMessageChannelBinder = new KafkaMessageChannelBinder( - configurationProperties, provisioningProvider); + configurationProperties, provisioningProvider, listenerContainerCustomizer); kafkaMessageChannelBinder.setProducerListener(producerListener); kafkaMessageChannelBinder.setExtendedBindingProperties(this.kafkaExtendedBindingProperties); return kafkaMessageChannelBinder; diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java index c54eadbac..daba1b48e 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java @@ -16,6 +16,9 @@ package org.springframework.cloud.stream.binder.kafka.integration; +import java.util.List; +import java.util.Map; + import io.micrometer.core.instrument.MeterRegistry; import io.micrometer.core.instrument.binder.MeterBinder; @@ -25,6 +28,7 @@ import org.junit.ClassRule; import org.junit.Test; import org.junit.runner.RunWith; +import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.test.context.FilteredClassLoader; @@ -32,9 +36,15 @@ import org.springframework.boot.test.context.SpringBootTest; import org.springframework.boot.test.context.runner.ApplicationContextRunner; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.annotation.StreamListener; +import org.springframework.cloud.stream.binder.Binding; +import org.springframework.cloud.stream.binding.BindingService; +import org.springframework.cloud.stream.config.ListenerContainerCustomizer; import org.springframework.cloud.stream.messaging.Sink; +import org.springframework.context.annotation.Bean; import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.test.rule.KafkaEmbedded; +import org.springframework.messaging.MessageChannel; import org.springframework.test.context.junit4.SpringRunner; import static org.assertj.core.api.Assertions.assertThat; @@ -93,6 +103,13 @@ public class KafkaBinderActuatorTests { .run(context -> { assertThat(context.getBeanNamesForType(MeterRegistry.class)).isEmpty(); assertThat(context.getBeanNamesForType(MeterBinder.class)).isEmpty(); + + DirectFieldAccessor channelBindingServiceAccessor = new DirectFieldAccessor(context.getBean(BindingService.class)); + @SuppressWarnings("unchecked") + Map>> consumerBindings = (Map>>) channelBindingServiceAccessor + .getPropertyValue("consumerBindings"); + assertThat(new DirectFieldAccessor(consumerBindings.get("input").get(0)).getPropertyValue("lifecycle.messageListenerContainer.beanName")) + .isEqualTo("setByCustomizer:input"); }); } @@ -100,6 +117,11 @@ public class KafkaBinderActuatorTests { @EnableAutoConfiguration public static class KafkaMetricsTestConfig { + @Bean + public ListenerContainerCustomizer> containerCustomizer() { + return (c, q, g) -> c.setBeanName("setByCustomizer:" + q); + } + @StreamListener(Sink.INPUT) public void process(String payload) throws InterruptedException { // Artificial slow listener to emulate consumer lag