From 1df6d872fa5a66424ee7c7e0638c4c159c4fc066 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 17 Jun 2014 12:31:10 -0400 Subject: [PATCH] INT-3430 AMQP Outbound: Add Eager Connect JIRA: https://jira.spring.io/browse/INT-3430 Add an option to eagerly connect to rabbit when only using outbound endpoints. Log an ERROR if the connection can not be established during context initialization. Polishing --- .../AmqpOutboundChannelAdapterParser.java | 1 + .../config/AmqpOutboundGatewayParser.java | 2 + .../amqp/outbound/AmqpOutboundEndpoint.java | 53 ++++++++++++++++--- .../config/spring-integration-amqp-4.1.xsd | 12 +++++ ...boundChannelAdapterParserTests-context.xml | 1 + ...AmqpOutboundChannelAdapterParserTests.java | 33 ++++++++++++ ...AmqpOutboundGatewayParserTests-context.xml | 1 + .../AmqpOutboundGatewayParserTests.java | 3 ++ .../context/IntegrationObjectSupport.java | 7 +++ src/reference/docbook/amqp.xml | 22 ++++++-- src/reference/docbook/whats-new.xml | 14 +++++ 11 files changed, 140 insertions(+), 9 deletions(-) diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java index 8c37ca3771..6f65114b0e 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java @@ -50,6 +50,7 @@ public class AmqpOutboundChannelAdapterParser extends AbstractOutboundChannelAda IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "routing-key", true); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "routing-key-expression"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "default-delivery-mode"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "lazy-connect"); IntegrationNamespaceUtils.configureHeaderMapper(element, builder, parserContext, DefaultAmqpHeaderMapper.class, null); diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java index 7e51d7b3bb..318ffbe054 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java @@ -30,6 +30,7 @@ import org.springframework.util.StringUtils; * @author Oleg Zhurakousky * @author Gunnar Hillert * @author Artem Bilan + * @author Gary Russell * * @since 2.1 */ @@ -56,6 +57,7 @@ public class AmqpOutboundGatewayParser extends AbstractConsumerEndpointParser { IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "reply-timeout", "sendTimeout"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "requires-reply"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "default-delivery-mode"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "lazy-connect"); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "reply-channel", "outputChannel"); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "return-channel"); diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpoint.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpoint.java index 5e8710d857..1b02207cdc 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpoint.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpoint.java @@ -20,11 +20,15 @@ import org.springframework.amqp.core.AmqpTemplate; import org.springframework.amqp.core.MessageDeliveryMode; import org.springframework.amqp.core.MessagePostProcessor; import org.springframework.amqp.core.MessageProperties; +import org.springframework.amqp.rabbit.connection.Connection; +import org.springframework.amqp.rabbit.connection.ConnectionFactory; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.amqp.rabbit.core.RabbitTemplate.ReturnCallback; import org.springframework.amqp.rabbit.support.CorrelationData; import org.springframework.amqp.support.converter.MessageConverter; import org.springframework.beans.factory.BeanFactory; +import org.springframework.context.ApplicationListener; +import org.springframework.context.event.ContextRefreshedEvent; import org.springframework.expression.Expression; import org.springframework.expression.ExpressionParser; import org.springframework.expression.spel.SpelParserConfiguration; @@ -50,9 +54,10 @@ import org.springframework.util.Assert; * @since 2.1 */ public class AmqpOutboundEndpoint extends AbstractReplyProducingMessageHandler - implements RabbitTemplate.ConfirmCallback, ReturnCallback { + implements RabbitTemplate.ConfirmCallback, ReturnCallback, ApplicationListener { - private static final ExpressionParser expressionParser = new SpelExpressionParser(new SpelParserConfiguration(true, true)); + private static final ExpressionParser expressionParser = + new SpelExpressionParser(new SpelParserConfiguration(true, true)); private final AmqpTemplate amqpTemplate; @@ -85,6 +90,8 @@ public class AmqpOutboundEndpoint extends AbstractReplyProducingMessageHandler private volatile MessageDeliveryMode defaultDeliveryMode; + private volatile boolean lazyConnect = true; + public AmqpOutboundEndpoint(AmqpTemplate amqpTemplate) { Assert.notNull(amqpTemplate, "amqpTemplate must not be null"); this.amqpTemplate = amqpTemplate; @@ -137,6 +144,17 @@ public class AmqpOutboundEndpoint extends AbstractReplyProducingMessageHandler this.defaultDeliveryMode = defaultDeliveryMode; } + /** + * Set to {@code false} to attempt to connect during endpoint start; + * default {@code true}, meaning the connection will be attempted + * to be established on the arrival of the first message. + * @param lazyConnect the lazyConnect to set + * @since 4.1 + */ + public void setLazyConnect(boolean lazyConnect) { + this.lazyConnect = lazyConnect; + } + @Override public String getComponentType() { return expectReply ? "amqp:outbound-gateway" : "amqp:outbound-channel-adapter"; @@ -188,6 +206,25 @@ public class AmqpOutboundEndpoint extends AbstractReplyProducingMessageHandler } } + @Override + public void onApplicationEvent(ContextRefreshedEvent event) { + if (!this.lazyConnect && event.getApplicationContext().equals(getApplicationContext()) + && this.amqpTemplate instanceof RabbitTemplate) { + ConnectionFactory connectionFactory = ((RabbitTemplate) this.amqpTemplate).getConnectionFactory(); + if (connectionFactory != null) { + try { + Connection connection = connectionFactory.createConnection(); + if (connection != null) { + connection.close(); + } + } + catch (RuntimeException e) { + logger.error("Failed to eagerly establish the connection.", e); + } + } + } + } + @Override protected Object handleRequestMessage(Message requestMessage) { String exchangeName = this.exchangeName; @@ -228,7 +265,8 @@ public class AmqpOutboundEndpoint extends AbstractReplyProducingMessageHandler @Override public org.springframework.amqp.core.Message postProcessMessage( org.springframework.amqp.core.Message message) throws AmqpException { - headerMapper.fromHeadersToRequest(requestMessage.getHeaders(), message.getMessageProperties()); + headerMapper.fromHeadersToRequest(requestMessage.getHeaders(), + message.getMessageProperties()); checkDeliveryMode(requestMessage, message.getMessageProperties()); return message; } @@ -241,7 +279,8 @@ public class AmqpOutboundEndpoint extends AbstractReplyProducingMessageHandler @Override public org.springframework.amqp.core.Message postProcessMessage( org.springframework.amqp.core.Message message) throws AmqpException { - headerMapper.fromHeadersToRequest(requestMessage.getHeaders(), message.getMessageProperties()); + headerMapper.fromHeadersToRequest(requestMessage.getHeaders(), + message.getMessageProperties()); return message; } }); @@ -253,10 +292,12 @@ public class AmqpOutboundEndpoint extends AbstractReplyProducingMessageHandler "RabbitTemplate implementation is required for publisher confirms"); MessageConverter converter = ((RabbitTemplate) this.amqpTemplate).getMessageConverter(); MessageProperties amqpMessageProperties = new MessageProperties(); - org.springframework.amqp.core.Message amqpMessage = converter.toMessage(requestMessage.getPayload(), amqpMessageProperties); + org.springframework.amqp.core.Message amqpMessage = + converter.toMessage(requestMessage.getPayload(), amqpMessageProperties); this.headerMapper.fromHeadersToRequest(requestMessage.getHeaders(), amqpMessageProperties); checkDeliveryMode(requestMessage, amqpMessageProperties); - org.springframework.amqp.core.Message amqpReplyMessage = this.amqpTemplate.sendAndReceive(exchangeName, routingKey, amqpMessage); + org.springframework.amqp.core.Message amqpReplyMessage = + this.amqpTemplate.sendAndReceive(exchangeName, routingKey, amqpMessage); if (amqpReplyMessage == null) { return null; } diff --git a/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-4.1.xsd b/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-4.1.xsd index 95d528ac1b..1396f02df7 100644 --- a/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-4.1.xsd +++ b/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-4.1.xsd @@ -489,6 +489,18 @@ property set to TRUE. + + + + By default, the connection is established lazily, when the first message is sent. If you wish to detect + connection configuration problems during application initialization, set this to 'false'. + If the eager connection fails, an ERROR log will be emitted. + + + + + + diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests-context.xml b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests-context.xml index 96c79e8da1..2abb6c8574 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests-context.xml +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests-context.xml @@ -24,6 +24,7 @@ diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java index 2a417c6817..4dca240ce0 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java @@ -17,12 +17,18 @@ package org.springframework.integration.amqp.config; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNull; import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; +import static org.mockito.Matchers.any; +import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import java.io.IOException; @@ -32,8 +38,10 @@ import java.util.List; import java.util.Map; import java.util.concurrent.atomic.AtomicBoolean; +import org.apache.commons.logging.Log; import org.junit.Test; import org.junit.runner.RunWith; +import org.mockito.Matchers; import org.mockito.Mockito; import org.mockito.invocation.InvocationOnMock; import org.mockito.stubbing.Answer; @@ -47,10 +55,12 @@ import org.springframework.amqp.rabbit.support.CorrelationData; import org.springframework.amqp.rabbit.support.PublisherCallbackChannel; import org.springframework.amqp.rabbit.support.PublisherCallbackChannelImpl; import org.springframework.beans.BeansException; +import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.parsing.BeanDefinitionParsingException; import org.springframework.context.ApplicationContext; +import org.springframework.context.event.ContextRefreshedEvent; import org.springframework.context.support.ClassPathXmlApplicationContext; import org.springframework.integration.amqp.AmqpHeaders; import org.springframework.integration.amqp.outbound.AmqpOutboundEndpoint; @@ -108,6 +118,7 @@ public class AmqpOutboundChannelAdapterParserTests { assertEquals("amqp:outbound-channel-adapter", ((NamedComponent) handler).getComponentType()); handler.handleMessage(new GenericMessage("foo")); assertEquals(1, adviceCalled); + assertTrue(TestUtils.getPropertyValue(handler, "lazyConnect", Boolean.class)); } @Test @@ -116,6 +127,7 @@ public class AmqpOutboundChannelAdapterParserTests { AmqpOutboundEndpoint endpoint = TestUtils.getPropertyValue(eventDrivenConsumer, "handler", AmqpOutboundEndpoint.class); assertNotNull(TestUtils.getPropertyValue(endpoint, "defaultDeliveryMode")); + assertFalse(TestUtils.getPropertyValue(endpoint, "lazyConnect", Boolean.class)); Field amqpTemplateField = ReflectionUtils.findField(AmqpOutboundEndpoint.class, "amqpTemplate"); amqpTemplateField.setAccessible(true); @@ -335,6 +347,27 @@ public class AmqpOutboundChannelAdapterParserTests { assertSame(this.context.getBean("customHeaderMapper"), headerMapper); } + @Test + public void testInt3430FailForNotLazyConnect() { + RabbitTemplate amqpTemplate = mock(RabbitTemplate.class); + ConnectionFactory connectionFactory = mock(ConnectionFactory.class); + RuntimeException toBeThrown = new RuntimeException("Test Connection Exception"); + doThrow(toBeThrown).when(connectionFactory).createConnection(); + when(amqpTemplate.getConnectionFactory()).thenReturn(connectionFactory); + AmqpOutboundEndpoint handler = new AmqpOutboundEndpoint(amqpTemplate); + Log logger = spy(TestUtils.getPropertyValue(handler, "logger", Log.class)); + new DirectFieldAccessor(handler).setPropertyValue("logger", logger); + ApplicationContext context = mock(ApplicationContext.class); + handler.setApplicationContext(context); + handler.afterPropertiesSet(); + ContextRefreshedEvent event = new ContextRefreshedEvent(context); + handler.onApplicationEvent(event); + verify(logger, never()).error(Matchers.anyString(), any(RuntimeException.class)); + handler.setLazyConnect(false); + handler.onApplicationEvent(event); + verify(logger).error("Failed to eagerly establish the connection.", toBeThrown); + } + public static class FooAdvice extends AbstractRequestHandlerAdvice { diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests-context.xml b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests-context.xml index add67812cc..afa225d3a0 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests-context.xml +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests-context.xml @@ -39,6 +39,7 @@ exchange-name="si.test.exchange" routing-key="si.test.binding" amqp-template="amqpTemplate" + lazy-connect="false" order="5" default-delivery-mode="NON_PERSISTENT" requires-reply="false" diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests.java index 2b3706107e..71c75b94c7 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests.java @@ -76,6 +76,8 @@ public class AmqpOutboundGatewayParserTests { Long sendTimeout = TestUtils.getPropertyValue(gateway, "messagingTemplate.sendTimeout", Long.class); assertEquals(Long.valueOf(777), sendTimeout); + assertTrue(TestUtils.getPropertyValue(gateway, "lazyConnect", Boolean.class)); + context.close(); } @@ -88,6 +90,7 @@ public class AmqpOutboundGatewayParserTests { AmqpOutboundEndpoint endpoint = TestUtils.getPropertyValue(eventDrivernConsumer, "handler", AmqpOutboundEndpoint.class); assertNotNull(TestUtils.getPropertyValue(endpoint, "defaultDeliveryMode")); + assertFalse(TestUtils.getPropertyValue(endpoint, "lazyConnect", Boolean.class)); assertFalse(TestUtils.getPropertyValue(endpoint, "requiresReply", Boolean.class)); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationObjectSupport.java b/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationObjectSupport.java index 04bffc9158..9bf5f51ee1 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationObjectSupport.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationObjectSupport.java @@ -192,6 +192,13 @@ public abstract class IntegrationObjectSupport implements BeanNameAware, NamedCo return this.applicationContext == null ? null : this.applicationContext.getId(); } + /** + * @return the applicationContext + */ + protected ApplicationContext getApplicationContext() { + return applicationContext; + } + /** * @see IntegrationContextUtils#getIntegrationProperties(BeanFactory) * @return The global integration properties. diff --git a/src/reference/docbook/amqp.xml b/src/reference/docbook/amqp.xml index 9ae0c62b84..ef49aed04b 100644 --- a/src/reference/docbook/amqp.xml +++ b/src/reference/docbook/amqp.xml @@ -418,7 +418,8 @@ public Object handle(@Payload String payload, @Header(AmqpHeaders.CHANNEL) Chann confirm-correlation-expression=""]]>]]> + return-channel=""]]>]]> @@ -505,6 +506,13 @@ public Object handle(@Payload String payload, @Header(AmqpHeaders.CHANNEL) Chann for each endpoint. + + When set to false, the endpoint will attempt to connect to the + broker during application context initialization. This allows "fail fast" detection of + bad configuration, but will also cause initialization to fail if the broker is down. + When true (default), the connection is established (if it doesn't already exist because + some other component established it) when the first message is sent. + @@ -521,8 +529,9 @@ public Object handle(@Payload String payload, @Header(AmqpHeaders.CHANNEL) Chann reply-channel=""]]>]]> + default-delivery-mode""]]>]]> @@ -593,6 +602,13 @@ public Object handle(@Payload String payload, @Header(AmqpHeaders.CHANNEL) Chann for each endpoint. + + When set to false, the endpoint will attempt to connect to the + broker during application context initialization. This allows "fail fast" detection of + bad configuration, but will also cause initialization to fail if the broker is down. + When true (default), the connection is established (if it doesn't already exist because + some other component established it) when the first message is sent. + diff --git a/src/reference/docbook/whats-new.xml b/src/reference/docbook/whats-new.xml index 4eb2389d1e..27b926935b 100644 --- a/src/reference/docbook/whats-new.xml +++ b/src/reference/docbook/whats-new.xml @@ -9,4 +9,18 @@ in more details, please see the Issue Tracker tickets that were resolved as part of the 4.1 development process. +
+ General Changes +
+ AMQP Outbound Endpoints + + The AMQP outbound endpoints support a new property lazy-connect + (default true). When true, the connection to the broker is not established + until the first message arrives (assuming there are no inbound endpoints, which + always attempt to establish the connection during startup). When set the 'false' an + attempt to establish the connection is made during application startup. + See for more information. + +
+