From f51cc43d52276e6ccf358e5a4b801c7eca718590 Mon Sep 17 00:00:00 2001 From: Dave Syer Date: Tue, 20 Oct 2015 07:58:34 -0400 Subject: [PATCH] Add @Conditional to allow running in Cloud Foundry without connectors --- .../RabbitServiceAutoConfiguration.java | 2 + .../config/RedisServiceAutoConfiguration.java | 2 + .../binder/MessageChannelBinderSupport.java | 59 +++++++++---------- 3 files changed, 33 insertions(+), 30 deletions(-) diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitServiceAutoConfiguration.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitServiceAutoConfiguration.java index 5b171a085..eb6d12941 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitServiceAutoConfiguration.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitServiceAutoConfiguration.java @@ -20,6 +20,7 @@ import org.springframework.amqp.rabbit.connection.ConnectionFactory; import org.springframework.boot.autoconfigure.AutoConfigureBefore; import org.springframework.boot.autoconfigure.amqp.RabbitAutoConfiguration; import org.springframework.boot.autoconfigure.cloud.CloudAutoConfiguration; +import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.cloud.Cloud; import org.springframework.cloud.CloudFactory; @@ -47,6 +48,7 @@ import org.springframework.context.annotation.PropertySource; public class RabbitServiceAutoConfiguration { @Configuration @Profile("cloud") + @ConditionalOnClass(Cloud.class) protected static class CloudConfig { @Bean public Cloud cloud() { diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-redis/src/main/java/org/springframework/cloud/stream/binder/redis/config/RedisServiceAutoConfiguration.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-redis/src/main/java/org/springframework/cloud/stream/binder/redis/config/RedisServiceAutoConfiguration.java index c858d3072..4bb3c3413 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-redis/src/main/java/org/springframework/cloud/stream/binder/redis/config/RedisServiceAutoConfiguration.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-redis/src/main/java/org/springframework/cloud/stream/binder/redis/config/RedisServiceAutoConfiguration.java @@ -18,6 +18,7 @@ package org.springframework.cloud.stream.binder.redis.config; import org.springframework.boot.autoconfigure.AutoConfigureBefore; import org.springframework.boot.autoconfigure.cloud.CloudAutoConfiguration; +import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.redis.RedisAutoConfiguration; import org.springframework.cloud.Cloud; @@ -46,6 +47,7 @@ import org.springframework.data.redis.connection.RedisConnectionFactory; public class RedisServiceAutoConfiguration { @Configuration + @ConditionalOnClass(Cloud.class) @Profile("cloud") protected static class CloudConfig { @Bean diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-spi/src/main/java/org/springframework/cloud/stream/binder/MessageChannelBinderSupport.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-spi/src/main/java/org/springframework/cloud/stream/binder/MessageChannelBinderSupport.java index 9f9a62c75..21520fd3b 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-spi/src/main/java/org/springframework/cloud/stream/binder/MessageChannelBinderSupport.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-spi/src/main/java/org/springframework/cloud/stream/binder/MessageChannelBinderSupport.java @@ -16,6 +16,12 @@ package org.springframework.cloud.stream.binder; +import static org.springframework.util.MimeTypeUtils.ALL; +import static org.springframework.util.MimeTypeUtils.APPLICATION_OCTET_STREAM; +import static org.springframework.util.MimeTypeUtils.APPLICATION_OCTET_STREAM_VALUE; +import static org.springframework.util.MimeTypeUtils.TEXT_PLAIN; +import static org.springframework.util.MimeTypeUtils.TEXT_PLAIN_VALUE; + import java.io.ByteArrayOutputStream; import java.io.IOException; import java.io.UnsupportedEncodingException; @@ -35,7 +41,6 @@ import java.util.concurrent.ConcurrentMap; import org.slf4j.Logger; import org.slf4j.LoggerFactory; - import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.InitializingBean; @@ -68,12 +73,6 @@ import org.springframework.util.IdGenerator; import org.springframework.util.MimeType; import org.springframework.util.StringUtils; -import static org.springframework.util.MimeTypeUtils.ALL; -import static org.springframework.util.MimeTypeUtils.APPLICATION_OCTET_STREAM; -import static org.springframework.util.MimeTypeUtils.APPLICATION_OCTET_STREAM_VALUE; -import static org.springframework.util.MimeTypeUtils.TEXT_PLAIN; -import static org.springframework.util.MimeTypeUtils.TEXT_PLAIN_VALUE; - /** * @author David Turanski * @author Gary Russell @@ -265,7 +264,7 @@ public abstract class MessageChannelBinderSupport } protected IdGenerator getIdGenerator() { - return idGenerator; + return this.idGenerator; } public void setIntegrationEvaluationContext(EvaluationContext evaluationContext) { @@ -377,7 +376,7 @@ public abstract class MessageChannelBinderSupport @Override public void afterPropertiesSet() throws Exception { - Assert.notNull(applicationContext, "The 'applicationContext' property cannot be null"); + Assert.notNull(this.applicationContext, "The 'applicationContext' property cannot be null"); onInit(); if (this.evaluationContext == null) { this.evaluationContext = IntegrationContextUtils.getEvaluationContext(getBeanFactory()); @@ -555,8 +554,8 @@ public abstract class MessageChannelBinderSupport bean.stop(); } catch (Exception e) { - if (logger.isWarnEnabled()) { - logger.warn("failed to stop adapter", e); + if (this.logger.isWarnEnabled()) { + this.logger.warn("failed to stop adapter", e); } } } @@ -605,7 +604,7 @@ public abstract class MessageChannelBinderSupport protected final MessageValues deserializePayloadIfNecessary(MessageValues message) { MessageValues messageToSend = message; Object originalPayload = message.getPayload(); - MimeType contentType = contentTypeResolver.resolve(messageToSend); + MimeType contentType = this.contentTypeResolver.resolve(messageToSend); Object payload = deserializePayload(originalPayload, contentType); if (payload != null) { messageToSend.setPayload(payload); @@ -643,12 +642,12 @@ public abstract class MessageChannelBinderSupport String className = JavaClassMimeTypeConversion.classNameFromMimeType(contentType); try { // Cache types to avoid unnecessary ClassUtils.forName calls. - Class targetType = payloadTypeCache.get(className); + Class targetType = this.payloadTypeCache.get(className); if (targetType == null) { targetType = ClassUtils.forName(className, null); - payloadTypeCache.put(className, targetType); + this.payloadTypeCache.put(className, targetType); } - return codec.decode(bytes, targetType); + return this.codec.decode(bytes, targetType); } catch (ClassNotFoundException e) { throw new SerializationFailedException("unable to deserialize [" + className + "]. Class not found.", @@ -708,7 +707,7 @@ public abstract class MessageChannelBinderSupport clazz = ClassUtils.forName(partitionKeyExtractorClassName, this.applicationContext.getClassLoader()); } catch (Exception e) { - logger.error("Failed to load key extractor", e); + this.logger.error("Failed to load key extractor", e); throw new BinderException("Failed to load key extractor: " + partitionKeyExtractorClassName, e); } try { @@ -719,7 +718,7 @@ public abstract class MessageChannelBinderSupport return ((PartitionKeyExtractorStrategy) extractor).extractKey(message); } catch (Exception e) { - logger.error("Failed to instantiate key extractor", e); + this.logger.error("Failed to instantiate key extractor", e); throw new BinderException("Failed to instantiate key extractor: " + partitionKeyExtractorClassName, e); } } @@ -734,7 +733,7 @@ public abstract class MessageChannelBinderSupport clazz = ClassUtils.forName(partitionSelectorClassName, this.applicationContext.getClassLoader()); } catch (Exception e) { - logger.error("Failed to load partition selector", e); + this.logger.error("Failed to load partition selector", e); throw new BinderException("Failed to load partition selector: " + partitionSelectorClassName, e); } try { @@ -745,7 +744,7 @@ public abstract class MessageChannelBinderSupport return ((PartitionSelectorStrategy) extractor).selectPartition(key, partitionCount); } catch (Exception e) { - logger.error("Failed to instantiate partition selector", e); + this.logger.error("Failed to instantiate partition selector", e); throw new BinderException("Failed to instantiate partition selector: " + partitionSelectorClassName, e); } @@ -882,8 +881,8 @@ public abstract class MessageChannelBinderSupport Binding binding = Binding.forDirectProducer(name, producerChannel, consumer, properties); addBinding(binding); binding.start(); - if (logger.isInfoEnabled()) { - logger.info("Producer bound directly: " + binding); + if (this.logger.isInfoEnabled()) { + this.logger.info("Producer bound directly: " + binding); } } @@ -934,14 +933,14 @@ public abstract class MessageChannelBinderSupport if (directBinding != null) { directBinding.stop(); this.bindings.remove(directBinding); - if (logger.isInfoEnabled()) { - logger.info("direct binding reverted: " + directBinding); + if (this.logger.isInfoEnabled()) { + this.logger.info("direct binding reverted: " + directBinding); } } } } catch (Exception e) { - logger.error("Could not revert direct binding: " + binding, e); + this.logger.error("Could not revert direct binding: " + binding, e); } } @@ -987,7 +986,7 @@ public abstract class MessageChannelBinderSupport } public int getPartitionCount() { - return partitionCount; + return this.partitionCount; } } @@ -1014,11 +1013,11 @@ public abstract class MessageChannelBinderSupport @SuppressWarnings("unchecked") public T createAndRegisterChannel(String name) { T channel = createSharedChannel(name); - ConfigurableListableBeanFactory beanFactory = applicationContext.getBeanFactory(); + ConfigurableListableBeanFactory beanFactory = MessageChannelBinderSupport.this.applicationContext.getBeanFactory(); beanFactory.registerSingleton(name, channel); channel = (T) beanFactory.initializeBean(channel, name); - if (logger.isDebugEnabled()) { - logger.debug("Registered channel:" + name); + if (MessageChannelBinderSupport.this.logger.isDebugEnabled()) { + MessageChannelBinderSupport.this.logger.debug("Registered channel:" + name); } return channel; } @@ -1027,9 +1026,9 @@ public abstract class MessageChannelBinderSupport public T lookupSharedChannel(String name) { T channel = null; - if (applicationContext.containsBean(name)) { + if (MessageChannelBinderSupport.this.applicationContext.containsBean(name)) { try { - channel = applicationContext.getBean(name, requiredType); + channel = MessageChannelBinderSupport.this.applicationContext.getBean(name, this.requiredType); } catch (Exception e) { throw new IllegalArgumentException("bean '" + name