From b3cf864675fedb20534f378d84cf7a3e7eaf0e17 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Mon, 19 Nov 2018 13:11:20 -0500 Subject: [PATCH] Fix LambdaMessageProcessor for conversion SO: https://stackoverflow.com/questions/53378821 There is no need to jump into the `MessageConverter` if payload type is assignable to the target type --- .../integration/handler/LambdaMessageProcessor.java | 4 +++- .../integration/dsl/flows/IntegrationFlowTests.java | 8 ++++---- 2 files changed, 7 insertions(+), 5 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/LambdaMessageProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/LambdaMessageProcessor.java index 067f1c175d..12d1dbb1c9 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/LambdaMessageProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/LambdaMessageProcessor.java @@ -30,6 +30,7 @@ import org.springframework.messaging.Message; import org.springframework.messaging.MessageHandlingException; import org.springframework.messaging.converter.MessageConverter; import org.springframework.util.Assert; +import org.springframework.util.ClassUtils; import org.springframework.util.ReflectionUtils; /** @@ -104,7 +105,8 @@ public class LambdaMessageProcessor implements MessageProcessor, BeanFac } } else { - if (this.payloadType != null) { + if (this.payloadType != null && + !ClassUtils.isAssignable(this.payloadType, message.getPayload().getClass())) { if (Message.class.isAssignableFrom(this.payloadType)) { args[i] = message; } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dsl/flows/IntegrationFlowTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dsl/flows/IntegrationFlowTests.java index 5d3b9b1838..0b5869573e 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dsl/flows/IntegrationFlowTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dsl/flows/IntegrationFlowTests.java @@ -302,7 +302,7 @@ public class IntegrationFlowTests { assertTrue(this.beanFactory.containsBean("lambdasFlow.transformer#0")); QueueChannel replyChannel = new QueueChannel(); - Message message = MessageBuilder.withPayload("World") + Message message = MessageBuilder.withPayload("World".getBytes()) .setHeader(MessageHeaders.REPLY_CHANNEL, replyChannel) .build(); this.lambdasInput.send(message); @@ -336,7 +336,7 @@ public class IntegrationFlowTests { } @Test - public void testGatewayFlow() throws Exception { + public void testGatewayFlow() { PollableChannel replyChannel = new QueueChannel(); Message message = MessageBuilder.withPayload("foo").setReplyChannel(replyChannel).build(); @@ -815,8 +815,8 @@ public class IntegrationFlowTests { @Bean public IntegrationFlow lambdasFlow() { return IntegrationFlows.from("lambdasInput") - .filter("World"::equals) - .transform("Hello "::concat) + .filter(String.class, "World"::equals) + .transform(String.class, "Hello "::concat) .get(); }