Performance optimizations for basic use cases

Fixes #1013

- Do not recreate MessageValues if content is byte[]

- Optimize StreamListener dispatching

   - avoid expression evaluation if not necessary
   - use `StreamListenerMessageHandler` directly if there is only one method
This commit is contained in:
Marius Bogoevici
2017-05-31 18:54:48 -04:00
committed by Ilayaperumal Gopinathan
parent 9935822eaf
commit 58172e0b5d
4 changed files with 63 additions and 37 deletions

View File

@@ -68,7 +68,7 @@ public abstract class AbstractBinder<T, C extends ConsumerProperties, P extends
protected final Log logger = LogFactory.getLog(getClass());
private final StringConvertingContentTypeResolver contentTypeResolver = new StringConvertingContentTypeResolver();
protected final StringConvertingContentTypeResolver contentTypeResolver = new StringConvertingContentTypeResolver();
private volatile AbstractApplicationContext applicationContext;

View File

@@ -44,6 +44,7 @@ import org.springframework.messaging.MessageHeaders;
import org.springframework.messaging.SubscribableChannel;
import org.springframework.util.Assert;
import org.springframework.util.MimeType;
import org.springframework.util.MimeTypeUtils;
/**
* {@link AbstractBinder} that serves as base class for {@link MessageChannel} binders.
@@ -480,7 +481,13 @@ public abstract class AbstractMessageChannelBinder<C extends ConsumerProperties,
messageValues = deserializePayloadIfNecessary(messageValues);
}
else {
messageValues = deserializePayloadIfNecessary(requestMessage);
MimeType contentType = AbstractMessageChannelBinder.this.contentTypeResolver.resolve(requestMessage.getHeaders());
if (contentType != null && !MimeTypeUtils.APPLICATION_OCTET_STREAM.equals(contentType)) {
messageValues = deserializePayloadIfNecessary(requestMessage);
}
else {
return requestMessage;
}
}
return messageValues.toMessage();
}

View File

@@ -18,18 +18,18 @@ package org.springframework.cloud.stream.binding;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.List;
import org.springframework.expression.EvaluationContext;
import org.springframework.expression.Expression;
import org.springframework.integration.handler.AbstractReplyProducingMessageHandler;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.MessagingException;
import org.springframework.util.Assert;
/**
* An {@link AbstractReplyProducingMessageHandler} that delegates to a collection of
* internal {@link ConditionalStreamListenerHandler} instances, executing the ones that
* internal {@link ConditionalStreamListenerMessageHandlerWrapper} instances, executing the ones that
* match the given expression.
*
* @author Marius Bogoevici
@@ -37,15 +37,27 @@ import org.springframework.util.Assert;
*/
final class DispatchingStreamListenerMessageHandler extends AbstractReplyProducingMessageHandler {
private final Collection<ConditionalStreamListenerHandler> handlerMethods;
private final List<ConditionalStreamListenerMessageHandlerWrapper> handlerMethods;
private final boolean evaluateExpressions;
private final EvaluationContext evaluationContext;
DispatchingStreamListenerMessageHandler(Collection<ConditionalStreamListenerHandler> handlerMethods,
DispatchingStreamListenerMessageHandler(Collection<ConditionalStreamListenerMessageHandlerWrapper> 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<ConditionalStreamListenerHandler> matchingHandlers = findMatchingHandlers(requestMessage);
List<ConditionalStreamListenerMessageHandlerWrapper> 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<ConditionalStreamListenerHandler> findMatchingHandlers(Message<?> message) {
ArrayList<ConditionalStreamListenerHandler> matchingMethods = new ArrayList<>();
for (ConditionalStreamListenerHandler conditionalStreamListenerHandlerMethod : this.handlerMethods) {
if (conditionalStreamListenerHandlerMethod.getCondition() == null) {
matchingMethods.add(conditionalStreamListenerHandlerMethod);
private List<ConditionalStreamListenerMessageHandlerWrapper> findMatchingHandlers(Message<?> message) {
ArrayList<ConditionalStreamListenerMessageHandlerWrapper> 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;
}
}
}

View File

@@ -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<String, List<StreamListenerHandlerMethodMapping>> mappedBindingEntry : mappedListenerMethods
.entrySet()) {
Collection<DispatchingStreamListenerMessageHandler.ConditionalStreamListenerHandler> handlers = new ArrayList<>();
ArrayList<DispatchingStreamListenerMessageHandler.ConditionalStreamListenerMessageHandlerWrapper> 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();