SGH-1616: Add MessageSourceCustomizer
Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/1616 Resolves #197
This commit is contained in:
committed by
Oleg Zhurakousky
parent
18a6663bb0
commit
fa573a752c
@@ -71,6 +71,7 @@ import org.springframework.cloud.stream.binder.rabbit.properties.RabbitExtendedB
|
||||
import org.springframework.cloud.stream.binder.rabbit.properties.RabbitProducerProperties;
|
||||
import org.springframework.cloud.stream.binder.rabbit.provisioning.RabbitExchangeQueueProvisioner;
|
||||
import org.springframework.cloud.stream.config.ListenerContainerCustomizer;
|
||||
import org.springframework.cloud.stream.config.MessageSourceCustomizer;
|
||||
import org.springframework.cloud.stream.provisioning.ConsumerDestination;
|
||||
import org.springframework.cloud.stream.provisioning.ProducerDestination;
|
||||
import org.springframework.context.support.GenericApplicationContext;
|
||||
@@ -163,14 +164,25 @@ public class RabbitMessageChannelBinder extends
|
||||
public RabbitMessageChannelBinder(ConnectionFactory connectionFactory,
|
||||
RabbitProperties rabbitProperties,
|
||||
RabbitExchangeQueueProvisioner provisioningProvider) {
|
||||
this(connectionFactory, rabbitProperties, provisioningProvider, null);
|
||||
|
||||
this(connectionFactory, rabbitProperties, provisioningProvider, null, null);
|
||||
}
|
||||
|
||||
public RabbitMessageChannelBinder(ConnectionFactory connectionFactory,
|
||||
RabbitProperties rabbitProperties,
|
||||
RabbitExchangeQueueProvisioner provisioningProvider,
|
||||
ListenerContainerCustomizer<AbstractMessageListenerContainer> containerCustomizer) {
|
||||
super(new String[0], provisioningProvider, containerCustomizer);
|
||||
|
||||
this(connectionFactory, rabbitProperties, provisioningProvider, containerCustomizer, null);
|
||||
}
|
||||
|
||||
public RabbitMessageChannelBinder(ConnectionFactory connectionFactory,
|
||||
RabbitProperties rabbitProperties,
|
||||
RabbitExchangeQueueProvisioner provisioningProvider,
|
||||
ListenerContainerCustomizer<AbstractMessageListenerContainer> containerCustomizer,
|
||||
MessageSourceCustomizer<AmqpMessageSource> sourceCustomizer) {
|
||||
|
||||
super(new String[0], provisioningProvider, containerCustomizer, sourceCustomizer);
|
||||
Assert.notNull(connectionFactory, "connectionFactory must not be null");
|
||||
Assert.notNull(rabbitProperties, "rabbitProperties must not be null");
|
||||
this.connectionFactory = connectionFactory;
|
||||
@@ -539,11 +551,13 @@ public class RabbitMessageChannelBinder extends
|
||||
protected PolledConsumerResources createPolledConsumerResources(String name,
|
||||
String group, ConsumerDestination destination,
|
||||
ExtendedConsumerProperties<RabbitConsumerProperties> consumerProperties) {
|
||||
|
||||
Assert.isTrue(!consumerProperties.isMultiplex(),
|
||||
"The Spring Integration polled MessageSource does not currently support muiltiple queues");
|
||||
AmqpMessageSource source = new AmqpMessageSource(this.connectionFactory,
|
||||
destination.getName());
|
||||
source.setRawMessageHeader(true);
|
||||
getMessageSourceCustomizer().configure(source, destination.getName(), group);
|
||||
return new PolledConsumerResources(source, registerErrorInfrastructure(
|
||||
destination, group, consumerProperties, true));
|
||||
}
|
||||
|
||||
@@ -36,9 +36,11 @@ import org.springframework.cloud.stream.binder.rabbit.properties.RabbitBinderCon
|
||||
import org.springframework.cloud.stream.binder.rabbit.properties.RabbitExtendedBindingProperties;
|
||||
import org.springframework.cloud.stream.binder.rabbit.provisioning.RabbitExchangeQueueProvisioner;
|
||||
import org.springframework.cloud.stream.config.ListenerContainerCustomizer;
|
||||
import org.springframework.cloud.stream.config.MessageSourceCustomizer;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.context.annotation.Import;
|
||||
import org.springframework.integration.amqp.inbound.AmqpMessageSource;
|
||||
import org.springframework.lang.Nullable;
|
||||
|
||||
/**
|
||||
@@ -70,11 +72,12 @@ public class RabbitMessageChannelBinderConfiguration {
|
||||
|
||||
@Bean
|
||||
RabbitMessageChannelBinder rabbitMessageChannelBinder(
|
||||
@Nullable ListenerContainerCustomizer<AbstractMessageListenerContainer> listenerContainerCustomizer)
|
||||
throws Exception {
|
||||
@Nullable ListenerContainerCustomizer<AbstractMessageListenerContainer> listenerContainerCustomizer,
|
||||
@Nullable MessageSourceCustomizer<AmqpMessageSource> sourceCustomizer) {
|
||||
|
||||
RabbitMessageChannelBinder binder = new RabbitMessageChannelBinder(
|
||||
this.rabbitConnectionFactory, this.rabbitProperties,
|
||||
provisioningProvider(), listenerContainerCustomizer);
|
||||
provisioningProvider(), listenerContainerCustomizer, sourceCustomizer);
|
||||
binder.setAdminAddresses(
|
||||
this.rabbitBinderConfigurationProperties.getAdminAddresses());
|
||||
binder.setCompressingPostProcessor(gZipPostProcessor());
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2015-2018 the original author or authors.
|
||||
* Copyright 2015-2019 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -43,21 +43,25 @@ import org.springframework.boot.autoconfigure.SpringBootApplication;
|
||||
import org.springframework.boot.builder.SpringApplicationBuilder;
|
||||
import org.springframework.cloud.Cloud;
|
||||
import org.springframework.cloud.stream.annotation.EnableBinding;
|
||||
import org.springframework.cloud.stream.annotation.Input;
|
||||
import org.springframework.cloud.stream.binder.Binder;
|
||||
import org.springframework.cloud.stream.binder.BinderFactory;
|
||||
import org.springframework.cloud.stream.binder.Binding;
|
||||
import org.springframework.cloud.stream.binder.ExtendedConsumerProperties;
|
||||
import org.springframework.cloud.stream.binder.ExtendedProducerProperties;
|
||||
import org.springframework.cloud.stream.binder.ExtendedPropertiesBinder;
|
||||
import org.springframework.cloud.stream.binder.PollableMessageSource;
|
||||
import org.springframework.cloud.stream.binder.rabbit.RabbitMessageChannelBinder;
|
||||
import org.springframework.cloud.stream.binder.rabbit.properties.RabbitConsumerProperties;
|
||||
import org.springframework.cloud.stream.binder.rabbit.properties.RabbitProducerProperties;
|
||||
import org.springframework.cloud.stream.binder.test.junit.rabbit.RabbitTestSupport;
|
||||
import org.springframework.cloud.stream.binding.BindingService;
|
||||
import org.springframework.cloud.stream.config.ListenerContainerCustomizer;
|
||||
import org.springframework.cloud.stream.config.MessageSourceCustomizer;
|
||||
import org.springframework.cloud.stream.messaging.Processor;
|
||||
import org.springframework.context.ConfigurableApplicationContext;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.integration.amqp.inbound.AmqpMessageSource;
|
||||
import org.springframework.integration.channel.DirectChannel;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
@@ -150,6 +154,7 @@ public class RabbitBinderModuleTests {
|
||||
public void testParentConnectionFactoryInheritedByDefaultAndRabbitSettingsPropagated() {
|
||||
context = new SpringApplicationBuilder(SimpleProcessor.class)
|
||||
.web(WebApplicationType.NONE).run("--server.port=0",
|
||||
"--spring.cloud.stream.bindings.source.group=someGroup",
|
||||
"--spring.cloud.stream.bindings.input.group=someGroup",
|
||||
"--spring.cloud.stream.rabbit.bindings.input.consumer.transacted=true",
|
||||
"--spring.cloud.stream.rabbit.bindings.output.producer.transacted=true");
|
||||
@@ -198,6 +203,9 @@ public class RabbitBinderModuleTests {
|
||||
ConnectionNameStrategy cns = TestUtils.getPropertyValue(cf,
|
||||
"connectionNameStrategy", ConnectionNameStrategy.class);
|
||||
assertThat(cns.obtainNewConnectionName(cf)).startsWith("rabbitConnectionFactory");
|
||||
assertThat(TestUtils.getPropertyValue(consumerBindings.get("source").get(0),
|
||||
"target.source.h.advised.targetSource.target.beanName"))
|
||||
.isEqualTo("setByCustomizer:someGroup");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -340,7 +348,7 @@ public class RabbitBinderModuleTests {
|
||||
assertThat(rabbitConsumerProperties.getMaxConcurrency()).isEqualTo(4);
|
||||
}
|
||||
|
||||
@EnableBinding(Processor.class)
|
||||
@EnableBinding({ Processor.class, PMS.class })
|
||||
@SpringBootApplication
|
||||
public static class SimpleProcessor {
|
||||
|
||||
@@ -350,6 +358,11 @@ public class RabbitBinderModuleTests {
|
||||
"setByCustomizerForQueue:" + q + (g == null ? "" : ",andGroup:" + g));
|
||||
}
|
||||
|
||||
@Bean
|
||||
public MessageSourceCustomizer<AmqpMessageSource> sourceCustomizer() {
|
||||
return (s, q, g) -> s.setBeanName("setByCustomizer:" + g);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
public static class ConnectionFactoryConfiguration {
|
||||
@@ -375,4 +388,11 @@ public class RabbitBinderModuleTests {
|
||||
|
||||
}
|
||||
|
||||
public interface PMS {
|
||||
|
||||
@Input
|
||||
PollableMessageSource source();
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user