GH-833: Add container configurer callback

Resolves https://github.com/spring-projects/spring-amqp/issues/833

Allow setting any properties not directly exposed by the factory.
This commit is contained in:
Gary Russell
2018-10-23 14:14:52 -04:00
committed by Artem Bilan
parent 969f095960
commit 0dca5bcded
2 changed files with 27 additions and 3 deletions

View File

@@ -19,6 +19,7 @@ package org.springframework.amqp.rabbit.config;
import java.util.concurrent.Executor;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Consumer;
import org.aopalliance.aop.Advice;
import org.apache.commons.logging.Log;
@@ -64,6 +65,8 @@ public abstract class AbstractRabbitListenerContainerFactory<C extends AbstractM
protected final Log logger = LogFactory.getLog(getClass());
protected final AtomicInteger counter = new AtomicInteger();
private ConnectionFactory connectionFactory;
private ErrorHandler errorHandler;
@@ -112,7 +115,7 @@ public abstract class AbstractRabbitListenerContainerFactory<C extends AbstractM
private RecoveryCallback<?> recoveryCallback;
protected final AtomicInteger counter = new AtomicInteger();
private Consumer<C> containerConfigurer;
/**
* @param connectionFactory The connection factory.
@@ -328,6 +331,16 @@ public abstract class AbstractRabbitListenerContainerFactory<C extends AbstractM
this.recoveryCallback = recoveryCallback;
}
/**
* A {@link Consumer} that is invoked to enable setting other container properties not
* exposed by this container factory.
* @param configurer the configurer;
* @since 2.1.1
*/
public void setContainerConfigurer(Consumer<C> configurer) {
this.containerConfigurer = configurer;
}
@SuppressWarnings("deprecation")
@Override
public C createListenerContainer(RabbitListenerEndpoint endpoint) {
@@ -340,8 +353,13 @@ public abstract class AbstractRabbitListenerContainerFactory<C extends AbstractM
instance.setErrorHandler(this.errorHandler);
}
if (this.messageConverter != null) {
endpoint.setMessageConverter(this.messageConverter);
if (endpoint.getMessageConverter() == null) {
if (endpoint != null) {
endpoint.setMessageConverter(this.messageConverter);
if (endpoint.getMessageConverter() == null) {
instance.setMessageConverter(this.messageConverter);
}
}
else {
instance.setMessageConverter(this.messageConverter);
}
}
@@ -422,6 +440,10 @@ public abstract class AbstractRabbitListenerContainerFactory<C extends AbstractM
}
initializeContainer(instance, endpoint);
if (this.containerConfigurer != null) {
this.containerConfigurer.accept(instance);
}
return instance;
}

View File

@@ -109,6 +109,7 @@ public class RabbitListenerContainerFactoryTests {
this.factory.setRecoveryBackOff(recoveryBackOff);
this.factory.setMissingQueuesFatal(true);
this.factory.setAfterReceivePostProcessors(afterReceivePostProcessor);
this.factory.setContainerConfigurer(c -> c.setShutdownTimeout(10_000));
assertArrayEquals(new Advice[] {advice}, this.factory.getAdviceChain());
@@ -131,6 +132,7 @@ public class RabbitListenerContainerFactoryTests {
assertEquals(6, fieldAccessor.getPropertyValue("consecutiveIdleTrigger"));
assertEquals(3, fieldAccessor.getPropertyValue("prefetchCount"));
assertEquals(1500L, fieldAccessor.getPropertyValue("receiveTimeout"));
assertEquals(10_000L, fieldAccessor.getPropertyValue("shutdownTimeout"));
assertEquals(false, fieldAccessor.getPropertyValue("defaultRequeueRejected"));
Advice[] actualAdviceChain = (Advice[]) fieldAccessor.getPropertyValue("adviceChain");
assertEquals("Wrong number of advice", 1, actualAdviceChain.length);