From 2e7f9cc729d37651985eff5485bb684f34c1e2a0 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 28 Oct 2015 12:42:12 -0400 Subject: [PATCH] GH-166: Fix `ChannelBinder` config inconsistency Fixes GH-166 (https://github.com/spring-cloud/spring-cloud-stream/issues/166) Since `onInit()` in the `MessageChannelBinderSupport` was called before the `evaluationContext` population, the `RedisMessageChannelBinder#errorAdapter` caused `NPE` for its `evaluationContext`. * Change the initialization order * Use `ExpressionUtils` to populate `evaluationContext` - creates `StandardEvaluationContext` if there is no one in the `applicationContext` * Add more initializations for the `RedisMessageChannelBinder#errorAdapter` * Remove explicit `evaluationContext` population from the `AbstractTestBinder` to rely on the fix for this ticket * Some other simple polishing Address PR comments --- .../local/LocalMessageChannelBinder.java | 3 ++- .../rabbit/RabbitMessageChannelBinder.java | 4 +-- .../redis/RedisMessageChannelBinder.java | 14 ++++++++--- .../binder/MessageChannelBinderSupport.java | 25 +++++++++---------- .../stream/binder/AbstractTestBinder.java | 2 -- 5 files changed, 26 insertions(+), 22 deletions(-) diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-local/src/main/java/org/springframework/cloud/stream/binder/local/LocalMessageChannelBinder.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-local/src/main/java/org/springframework/cloud/stream/binder/local/LocalMessageChannelBinder.java index bbf69c428..39fefd19a 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-local/src/main/java/org/springframework/cloud/stream/binder/local/LocalMessageChannelBinder.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-local/src/main/java/org/springframework/cloud/stream/binder/local/LocalMessageChannelBinder.java @@ -173,7 +173,8 @@ public class LocalMessageChannelBinder extends MessageChannelBinderSupport { } @Override - protected void onInit() { + public void afterPropertiesSet() throws Exception { + super.afterPropertiesSet(); this.executor.setCorePoolSize(this.executorCorePoolSize); this.executor.setMaxPoolSize(this.executorMaxPoolSize); this.executor.setQueueCapacity(this.executorQueueSize); diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java index a3eac3eba..2070a7a6b 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java @@ -420,8 +420,8 @@ public class RabbitMessageChannelBinder extends MessageChannelBinderSupport impl } @Override - protected void onInit() { - super.onInit(); + public void afterPropertiesSet() throws Exception { + super.afterPropertiesSet(); if (this.clustered) { Assert.state(this.addresses.length == this.adminAddresses.length && this.addresses.length == this.nodes.length, diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-redis/src/main/java/org/springframework/cloud/stream/binder/redis/RedisMessageChannelBinder.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-redis/src/main/java/org/springframework/cloud/stream/binder/redis/RedisMessageChannelBinder.java index ac80dc493..48ea26376 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-redis/src/main/java/org/springframework/cloud/stream/binder/redis/RedisMessageChannelBinder.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-redis/src/main/java/org/springframework/cloud/stream/binder/redis/RedisMessageChannelBinder.java @@ -147,8 +147,7 @@ public class RedisMessageChannelBinder extends MessageChannelBinderSupport imple String... headersToMap) { Assert.notNull(connectionFactory, "connectionFactory must not be null"); this.connectionFactory = connectionFactory; - this.errorAdapter = new RedisQueueOutboundChannelAdapter( - parser.parseExpression("headers['" + ERROR_HEADER + "']"), connectionFactory); + if (headersToMap != null && headersToMap.length > 0) { String[] combinedHeadersToMap = Arrays.copyOfRange(BinderHeaders.STANDARD_HEADERS, 0, BinderHeaders.STANDARD_HEADERS.length @@ -160,11 +159,17 @@ public class RedisMessageChannelBinder extends MessageChannelBinderSupport imple else { this.headersToMap = BinderHeaders.STANDARD_HEADERS; } + + this.errorAdapter = new RedisQueueOutboundChannelAdapter( + parser.parseExpression("headers['" + ERROR_HEADER + "']"), connectionFactory); } @Override - protected void onInit() { + public void afterPropertiesSet() throws Exception { + super.afterPropertiesSet(); this.errorAdapter.setIntegrationEvaluationContext(this.evaluationContext); + this.errorAdapter.setBeanFactory(getBeanFactory()); + this.errorAdapter.afterPropertiesSet(); } @Override @@ -424,7 +429,7 @@ public class RedisMessageChannelBinder extends MessageChannelBinderSupport imple transformed.put(PARTITION_HEADER, determinePartition(message, this.partitioningMetadata)); } - + byte[] messageToSend = embeddedHeadersMessageConverter.embedHeaders(transformed, RedisMessageChannelBinder.this.headersToMap); delegate.handleMessage(MessageBuilder.withPayload(messageToSend).copyHeaders(transformed).build()); @@ -458,6 +463,7 @@ public class RedisMessageChannelBinder extends MessageChannelBinderSupport imple // prevent returned message from being copied in superclass return false; } + } private static class RedisPropertiesAccessor extends AbstractBinderPropertiesAccessor { 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 21520fd3b..64485fc74 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 @@ -41,6 +41,7 @@ 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; @@ -55,8 +56,8 @@ import org.springframework.expression.EvaluationContext; import org.springframework.expression.Expression; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.codec.Codec; -import org.springframework.integration.context.IntegrationContextUtils; import org.springframework.integration.endpoint.EventDrivenConsumer; +import org.springframework.integration.expression.ExpressionUtils; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; @@ -377,15 +378,11 @@ public abstract class MessageChannelBinderSupport @Override public void afterPropertiesSet() throws Exception { Assert.notNull(this.applicationContext, "The 'applicationContext' property cannot be null"); - onInit(); if (this.evaluationContext == null) { - this.evaluationContext = IntegrationContextUtils.getEvaluationContext(getBeanFactory()); + this.evaluationContext = ExpressionUtils.createStandardEvaluationContext(getBeanFactory()); } } - protected void onInit() { - } - /** * Dynamically create a producer for the named channel. * @param name The name. @@ -602,18 +599,17 @@ public abstract class MessageChannelBinderSupport } protected final MessageValues deserializePayloadIfNecessary(MessageValues message) { - MessageValues messageToSend = message; Object originalPayload = message.getPayload(); - MimeType contentType = this.contentTypeResolver.resolve(messageToSend); + MimeType contentType = this.contentTypeResolver.resolve(message); Object payload = deserializePayload(originalPayload, contentType); if (payload != null) { - messageToSend.setPayload(payload); + message.setPayload(payload); - Object originalContentType = messageToSend.get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE); - messageToSend.put(MessageHeaders.CONTENT_TYPE, originalContentType); - messageToSend.put(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, null); + Object originalContentType = message.get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE); + message.put(MessageHeaders.CONTENT_TYPE, originalContentType); + message.put(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, null); } - return messageToSend; + return message; } private Object deserializePayload(Object payload, MimeType contentType) { @@ -988,6 +984,7 @@ public abstract class MessageChannelBinderSupport public int getPartitionCount() { return this.partitionCount; } + } /** @@ -1037,6 +1034,7 @@ public abstract class MessageChannelBinderSupport } return channel; } + } /** @@ -1095,6 +1093,7 @@ public abstract class MessageChannelBinderSupport } return className; } + } public static class SetBuilder { diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractTestBinder.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractTestBinder.java index 4147537ae..216486c75 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractTestBinder.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractTestBinder.java @@ -20,7 +20,6 @@ import java.util.HashSet; import java.util.Properties; import java.util.Set; -import org.springframework.expression.spel.support.StandardEvaluationContext; import org.springframework.messaging.MessageChannel; @@ -39,7 +38,6 @@ public abstract class AbstractTestBinder private C binder; public void setBinder(C binder) { - binder.setIntegrationEvaluationContext(new StandardEvaluationContext()); try { binder.afterPropertiesSet(); }