Configure consumer customizer
Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/1841 * Add test
This commit is contained in:
committed by
Artem Bilan
parent
bf30ecdea1
commit
278ba795d0
@@ -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<AbstractMessageListenerContainer<?, ?>> listenerContainerCustomizer,
|
||||
@Nullable MessageSourceCustomizer<KafkaMessageSource<?, ?>> sourceCustomizer,
|
||||
@Nullable ProducerMessageHandlerCustomizer<KafkaProducerMessageHandler<?, ?>> messageHandlerCustomizer,
|
||||
@Nullable ConsumerEndpointCustomizer<KafkaMessageDrivenChannelAdapter<?, ?>> consumerCustomizer,
|
||||
ObjectProvider<KafkaBindingRebalanceListener> rebalanceListener,
|
||||
ObjectProvider<DlqPartitionFunction> dlqPartitionFunction) {
|
||||
|
||||
@@ -118,6 +121,7 @@ public class KafkaBinderConfiguration {
|
||||
kafkaMessageChannelBinder
|
||||
.setExtendedBindingProperties(this.kafkaExtendedBindingProperties);
|
||||
kafkaMessageChannelBinder.setProducerMessageHandlerCustomizer(messageHandlerCustomizer);
|
||||
kafkaMessageChannelBinder.setConsumerEndpointCustomizer(consumerCustomizer);
|
||||
return kafkaMessageChannelBinder;
|
||||
}
|
||||
|
||||
|
||||
@@ -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<String, List<Binding<MessageChannel>>> consumerBindings = (Map<String, List<Binding<MessageChannel>>>) channelBindingServiceAccessor
|
||||
Map<String, List<Binding<MessageChannel>>> consumerBindings =
|
||||
(Map<String, List<Binding<MessageChannel>>>) 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<String, Binding<MessageChannel>> producerBindings = (Map<String, Binding<MessageChannel>>) channelBindingServiceAccessor
|
||||
@SuppressWarnings("unchecked")
|
||||
Map<String, Binding<MessageChannel>> producerBindings =
|
||||
(Map<String, Binding<MessageChannel>>) channelBindingServiceAccessor
|
||||
.getPropertyValue("producerBindings");
|
||||
|
||||
assertThat(new DirectFieldAccessor(
|
||||
@@ -155,6 +162,11 @@ public class KafkaBinderActuatorTests {
|
||||
return (s, q, g) -> s.setBeanName("setByCustomizer:" + q);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ConsumerEndpointCustomizer<KafkaMessageDrivenChannelAdapter<?, ?>> consumerCustomizer() {
|
||||
return (p, q, g) -> p.setBeanName("setByCustomizer:" + q);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ProducerMessageHandlerCustomizer<KafkaProducerMessageHandler<?, ?>> handlerCustomizer() {
|
||||
return (handler, destinationName) -> handler.setBeanName("setByCustomizer:" + destinationName);
|
||||
|
||||
Reference in New Issue
Block a user