diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java index 5772c1049..abba4c64b 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java @@ -68,7 +68,7 @@ public abstract class AbstractBinder handlerMethods; + private final List handlerMethods; + + private final boolean evaluateExpressions; private final EvaluationContext evaluationContext; - DispatchingStreamListenerMessageHandler(Collection handlerMethods, + DispatchingStreamListenerMessageHandler(Collection handlerMethods, EvaluationContext evaluationContext) { Assert.notEmpty(handlerMethods, "'handlerMethods' cannot be empty"); - Assert.notNull(evaluationContext, "'evaluationContext' cannot be empty"); - this.handlerMethods = handlerMethods; + this.handlerMethods = Collections.unmodifiableList(new ArrayList<>(handlerMethods)); + boolean evaluateExpressions = false; + for (ConditionalStreamListenerMessageHandlerWrapper handlerMethod : handlerMethods) { + if (handlerMethod.getCondition() != null) { + evaluateExpressions = true; + break; + } + } + this.evaluateExpressions = evaluateExpressions; + if (evaluateExpressions) { + Assert.notNull(evaluationContext, "'evaluationContext' cannot be null if conditions are used"); + } this.evaluationContext = evaluationContext; } @@ -56,7 +68,7 @@ final class DispatchingStreamListenerMessageHandler extends AbstractReplyProduci @Override protected Object handleRequestMessage(Message requestMessage) { - Collection matchingHandlers = findMatchingHandlers(requestMessage); + List matchingHandlers = this.evaluateExpressions ? findMatchingHandlers(requestMessage) : this.handlerMethods; if (matchingHandlers.size() == 0) { if (logger.isWarnEnabled()) { logger.warn("Cannot find a @StreamListener matching for message with id: " @@ -65,42 +77,42 @@ final class DispatchingStreamListenerMessageHandler extends AbstractReplyProduci return null; } else if (matchingHandlers.size() > 1) { - for (ConditionalStreamListenerHandler matchingMethod : matchingHandlers) { - matchingMethod.handleMessage(requestMessage); + for (ConditionalStreamListenerMessageHandlerWrapper matchingMethod : matchingHandlers) { + matchingMethod.getStreamListenerMessageHandler().handleMessage(requestMessage); } return null; } else { - final ConditionalStreamListenerHandler singleMatchingHandler = matchingHandlers.iterator().next(); - singleMatchingHandler.handleMessage(requestMessage); + final ConditionalStreamListenerMessageHandlerWrapper singleMatchingHandler = matchingHandlers.get(0); + singleMatchingHandler.getStreamListenerMessageHandler().handleMessage(requestMessage); return null; } } - private Collection findMatchingHandlers(Message message) { - ArrayList matchingMethods = new ArrayList<>(); - for (ConditionalStreamListenerHandler conditionalStreamListenerHandlerMethod : this.handlerMethods) { - if (conditionalStreamListenerHandlerMethod.getCondition() == null) { - matchingMethods.add(conditionalStreamListenerHandlerMethod); + private List findMatchingHandlers(Message message) { + ArrayList matchingMethods = new ArrayList<>(); + for (ConditionalStreamListenerMessageHandlerWrapper conditionalStreamListenerMessageHandlerWrapperMethod : this.handlerMethods) { + if (conditionalStreamListenerMessageHandlerWrapperMethod.getCondition() == null) { + matchingMethods.add(conditionalStreamListenerMessageHandlerWrapperMethod); } else { - boolean conditionMetOnMessage = conditionalStreamListenerHandlerMethod.getCondition().getValue( + boolean conditionMetOnMessage = conditionalStreamListenerMessageHandlerWrapperMethod.getCondition().getValue( this.evaluationContext, message, Boolean.class); if (conditionMetOnMessage) { - matchingMethods.add(conditionalStreamListenerHandlerMethod); + matchingMethods.add(conditionalStreamListenerMessageHandlerWrapperMethod); } } } return matchingMethods; } - static class ConditionalStreamListenerHandler implements MessageHandler { + static class ConditionalStreamListenerMessageHandlerWrapper { - private Expression condition; + private final Expression condition; - private StreamListenerMessageHandler streamListenerMessageHandler; + private final StreamListenerMessageHandler streamListenerMessageHandler; - ConditionalStreamListenerHandler(Expression condition, + ConditionalStreamListenerMessageHandlerWrapper(Expression condition, StreamListenerMessageHandler streamListenerMessageHandler) { Assert.notNull(streamListenerMessageHandler, "the message handler cannot be null"); Assert.isTrue(condition == null || streamListenerMessageHandler.isVoid(), @@ -117,9 +129,8 @@ final class DispatchingStreamListenerMessageHandler extends AbstractReplyProduci return this.streamListenerMessageHandler.isVoid(); } - @Override - public void handleMessage(Message message) throws MessagingException { - this.streamListenerMessageHandler.handleMessage(message); + public StreamListenerMessageHandler getStreamListenerMessageHandler() { + return streamListenerMessageHandler; } } } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/StreamListenerAnnotationBeanPostProcessor.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/StreamListenerAnnotationBeanPostProcessor.java index 660ca6ca0..f3f62c2a6 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/StreamListenerAnnotationBeanPostProcessor.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/StreamListenerAnnotationBeanPostProcessor.java @@ -18,7 +18,6 @@ package org.springframework.cloud.stream.binding; import java.lang.reflect.Method; import java.util.ArrayList; -import java.util.Collection; import java.util.List; import java.util.Map; @@ -52,6 +51,7 @@ import org.springframework.expression.EvaluationContext; import org.springframework.expression.Expression; import org.springframework.expression.spel.standard.SpelExpressionParser; import org.springframework.integration.context.IntegrationContextUtils; +import org.springframework.integration.handler.AbstractReplyProducingMessageHandler; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.SubscribableChannel; import org.springframework.messaging.core.DestinationResolver; @@ -343,7 +343,7 @@ public class StreamListenerAnnotationBeanPostProcessor this.evaluationContext = IntegrationContextUtils.getEvaluationContext(this.applicationContext.getBeanFactory()); for (Map.Entry> mappedBindingEntry : mappedListenerMethods .entrySet()) { - Collection handlers = new ArrayList<>(); + ArrayList handlers = new ArrayList<>(); for (StreamListenerHandlerMethodMapping mapping : mappedBindingEntry.getValue()) { final InvocableHandlerMethod invocableHandlerMethod = this.messageHandlerMethodFactory .createInvocableHandlerMethod(mapping.getTargetBean(), @@ -359,21 +359,29 @@ public class StreamListenerAnnotationBeanPostProcessor if (StringUtils.hasText(mapping.getCondition())) { String conditionAsString = resolveExpressionAsString(mapping.getCondition()); Expression condition = SPEL_EXPRESSION_PARSER.parseExpression(conditionAsString); - handlers.add(new DispatchingStreamListenerMessageHandler.ConditionalStreamListenerHandler( - condition, streamListenerMessageHandler)); + handlers.add( + new DispatchingStreamListenerMessageHandler.ConditionalStreamListenerMessageHandlerWrapper( + condition, streamListenerMessageHandler)); } else { - handlers.add(new DispatchingStreamListenerMessageHandler.ConditionalStreamListenerHandler( - null, streamListenerMessageHandler)); + handlers.add( + new DispatchingStreamListenerMessageHandler.ConditionalStreamListenerMessageHandlerWrapper( + null, streamListenerMessageHandler)); } } if (handlers.size() > 1) { - for (DispatchingStreamListenerMessageHandler.ConditionalStreamListenerHandler handler : handlers) { + for (DispatchingStreamListenerMessageHandler.ConditionalStreamListenerMessageHandlerWrapper handler : handlers) { Assert.isTrue(handler.isVoid(), StreamListenerErrorMessages.MULTIPLE_VALUE_RETURNING_METHODS); } } - DispatchingStreamListenerMessageHandler handler = new DispatchingStreamListenerMessageHandler( - handlers, this.evaluationContext); + AbstractReplyProducingMessageHandler handler; + + if (handlers.size() > 1 || handlers.get(0).getCondition() != null) { + handler = new DispatchingStreamListenerMessageHandler(handlers, this.evaluationContext); + } + else { + handler = handlers.get(0).getStreamListenerMessageHandler(); + } handler.setApplicationContext(this.applicationContext); handler.setChannelResolver(this.binderAwareChannelResolver); handler.afterPropertiesSet();