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 df5a4e340..6cd6360cd 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 @@ -40,6 +40,7 @@ import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfi import org.springframework.cloud.stream.binder.kafka.properties.KafkaExtendedBindingProperties; import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; import org.springframework.cloud.stream.binder.kafka.utils.DlqPartitionFunction; +import org.springframework.cloud.stream.config.ConsumerEndpointCustomizer; import org.springframework.cloud.stream.config.ListenerContainerCustomizer; import org.springframework.cloud.stream.config.MessageSourceCustomizer; import org.springframework.cloud.stream.config.ProducerMessageHandlerCustomizer; @@ -48,6 +49,7 @@ import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; +import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter; import org.springframework.integration.kafka.inbound.KafkaMessageSource; import org.springframework.integration.kafka.outbound.KafkaProducerMessageHandler; import org.springframework.kafka.listener.AbstractMessageListenerContainer; @@ -107,6 +109,7 @@ public class KafkaBinderConfiguration { @Nullable ListenerContainerCustomizer> listenerContainerCustomizer, @Nullable MessageSourceCustomizer> sourceCustomizer, @Nullable ProducerMessageHandlerCustomizer> messageHandlerCustomizer, + @Nullable ConsumerEndpointCustomizer> consumerCustomizer, ObjectProvider rebalanceListener, ObjectProvider dlqPartitionFunction) { @@ -118,6 +121,7 @@ public class KafkaBinderConfiguration { kafkaMessageChannelBinder .setExtendedBindingProperties(this.kafkaExtendedBindingProperties); kafkaMessageChannelBinder.setProducerMessageHandlerCustomizer(messageHandlerCustomizer); + kafkaMessageChannelBinder.setConsumerEndpointCustomizer(consumerCustomizer); 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 0be08139c..f722dc510 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 @@ -39,12 +39,14 @@ import org.springframework.cloud.stream.annotation.StreamListener; import org.springframework.cloud.stream.binder.Binding; import org.springframework.cloud.stream.binder.PollableMessageSource; import org.springframework.cloud.stream.binding.BindingService; +import org.springframework.cloud.stream.config.ConsumerEndpointCustomizer; import org.springframework.cloud.stream.config.ListenerContainerCustomizer; import org.springframework.cloud.stream.config.MessageSourceCustomizer; import org.springframework.cloud.stream.config.ProducerMessageHandlerCustomizer; import org.springframework.cloud.stream.messaging.Processor; import org.springframework.cloud.stream.messaging.Sink; import org.springframework.context.annotation.Bean; +import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter; import org.springframework.integration.kafka.inbound.KafkaMessageSource; import org.springframework.integration.kafka.outbound.KafkaProducerMessageHandler; import org.springframework.kafka.core.KafkaTemplate; @@ -117,21 +119,26 @@ public class KafkaBinderActuatorTests { DirectFieldAccessor channelBindingServiceAccessor = new DirectFieldAccessor( context.getBean(BindingService.class)); - // @checkstyle:off @SuppressWarnings("unchecked") - Map>> consumerBindings = (Map>>) channelBindingServiceAccessor + Map>> consumerBindings = + (Map>>) channelBindingServiceAccessor .getPropertyValue("consumerBindings"); - // @checkstyle:on assertThat(new DirectFieldAccessor( consumerBindings.get("input").get(0)).getPropertyValue( "lifecycle.messageListenerContainer.beanName")) .isEqualTo("setByCustomizer:input"); + assertThat(new DirectFieldAccessor( + consumerBindings.get("input").get(0)).getPropertyValue( + "lifecycle.beanName")) + .isEqualTo("setByCustomizer:input"); assertThat(new DirectFieldAccessor( consumerBindings.get("source").get(0)).getPropertyValue( "lifecycle.beanName")) .isEqualTo("setByCustomizer:source"); - Map> producerBindings = (Map>) channelBindingServiceAccessor + @SuppressWarnings("unchecked") + Map> producerBindings = + (Map>) channelBindingServiceAccessor .getPropertyValue("producerBindings"); assertThat(new DirectFieldAccessor( @@ -155,6 +162,11 @@ public class KafkaBinderActuatorTests { return (s, q, g) -> s.setBeanName("setByCustomizer:" + q); } + @Bean + public ConsumerEndpointCustomizer> consumerCustomizer() { + return (p, q, g) -> p.setBeanName("setByCustomizer:" + q); + } + @Bean public ProducerMessageHandlerCustomizer> handlerCustomizer() { return (handler, destinationName) -> handler.setBeanName("setByCustomizer:" + destinationName);