From 3c5db06ed7265ae6095ca248409d4ce0d6f2328d Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Mon, 26 May 2014 18:33:12 +0300 Subject: [PATCH] DSL: Add `LambdaMessageProcessor` --- .../dsl/IntegrationFlowBuilder.java | 140 ++++++++++++------ .../dsl/LambdaMessageProcessor.java | 111 ++++++++++++++ .../dsl/support/GenericHandler.java | 28 ++++ .../dsl/test/IntegrationFlowTests.java | 66 ++++----- 4 files changed, 263 insertions(+), 82 deletions(-) create mode 100644 spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/LambdaMessageProcessor.java create mode 100644 spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/GenericHandler.java diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlowBuilder.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlowBuilder.java index 362a2bc..8404296 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlowBuilder.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/IntegrationFlowBuilder.java @@ -36,6 +36,7 @@ import org.springframework.integration.dsl.support.BeanNameMessageProcessor; import org.springframework.integration.dsl.support.ComponentConfigurer; import org.springframework.integration.dsl.support.EndpointConfigurer; import org.springframework.integration.dsl.support.FixedSubscriberChannelPrototype; +import org.springframework.integration.dsl.support.GenericHandler; import org.springframework.integration.dsl.support.GenericRouter; import org.springframework.integration.dsl.support.GenericSplitter; import org.springframework.integration.dsl.support.MessageChannelReference; @@ -136,14 +137,25 @@ public final class IntegrationFlowBuilder { } public IntegrationFlowBuilder transform(GenericTransformer genericTransformer) { - return this.transform(genericTransformer, null); + return this.transform(null, genericTransformer); + } + + public IntegrationFlowBuilder transform(Class

payloadType, GenericTransformer genericTransformer) { + return this.transform(payloadType, genericTransformer, null); } public IntegrationFlowBuilder transform(GenericTransformer genericTransformer, EndpointConfigurer> endpointConfigurer) { + return this.transform(null, genericTransformer, endpointConfigurer); + } + + public IntegrationFlowBuilder transform(Class

payloadType, GenericTransformer genericTransformer, + EndpointConfigurer> endpointConfigurer) { Assert.notNull(genericTransformer); - Transformer transformer = genericTransformer instanceof Transformer - ? (Transformer) genericTransformer : new MethodInvokingTransformer(genericTransformer); + Transformer transformer = genericTransformer instanceof Transformer ? (Transformer) genericTransformer : + (isLambda(genericTransformer) + ? new MethodInvokingTransformer(new LambdaMessageProcessor(genericTransformer, payloadType)) + : new MethodInvokingTransformer(genericTransformer)); return this.handle(new MessageTransformingHandler(transformer), endpointConfigurer); } @@ -153,14 +165,25 @@ public final class IntegrationFlowBuilder { } public IntegrationFlowBuilder filter(GenericSelector genericSelector) { - return this.filter(genericSelector, null); + return this.filter(null, genericSelector); } - public IntegrationFlowBuilder filter(GenericSelector genericSelector, + public

IntegrationFlowBuilder filter(Class

payloadType, GenericSelector

genericSelector) { + return this.filter(payloadType, genericSelector, null); + } + + public

IntegrationFlowBuilder filter(GenericSelector

genericSelector, + EndpointConfigurer endpointConfigurer) { + return filter(null, genericSelector, endpointConfigurer); + } + + public

IntegrationFlowBuilder filter(Class

payloadType, GenericSelector

genericSelector, EndpointConfigurer endpointConfigurer) { Assert.notNull(genericSelector); - MessageSelector selector = genericSelector instanceof MessageSelector - ? (MessageSelector) genericSelector : new MethodInvokingSelector(genericSelector); + MessageSelector selector = genericSelector instanceof MessageSelector ? (MessageSelector) genericSelector : + (isLambda(genericSelector) + ? new MethodInvokingSelector(new LambdaMessageProcessor(genericSelector, payloadType)) + : new MethodInvokingSelector(genericSelector)); return this.register(new FilterEndpointSpec(new MessageFilter(selector)), endpointConfigurer); } @@ -178,6 +201,32 @@ public final class IntegrationFlowBuilder { endpointConfigurer); } + public

IntegrationFlowBuilder handle(GenericHandler

handler) { + return this.handle(null, handler); + } + + public

IntegrationFlowBuilder handle(GenericHandler

handler, + EndpointConfigurer> endpointConfigurer) { + return this.handle(null, handler, endpointConfigurer); + } + + + public

IntegrationFlowBuilder handle(Class

payloadType, GenericHandler

handler) { + return this.handle(payloadType, handler, null); + } + + public

IntegrationFlowBuilder handle(Class

payloadType, GenericHandler

handler, + EndpointConfigurer> endpointConfigurer) { + ServiceActivatingHandler serviceActivatingHandler = null; + if (isLambda(handler)) { + serviceActivatingHandler = new ServiceActivatingHandler(new LambdaMessageProcessor(handler, payloadType)); + } + else { + serviceActivatingHandler = new ServiceActivatingHandler(handler); + } + return this.handle(serviceActivatingHandler, endpointConfigurer); + } + public IntegrationFlowBuilder handle(H messageHandler, EndpointConfigurer> endpointConfigurer) { Assert.notNull(messageHandler); @@ -243,19 +292,10 @@ public final class IntegrationFlowBuilder { return this.addComponent(headerEnricher).transform(headerEnricher, endpointConfigurer); } - public IntegrationFlowBuilder split() { - return this.split((EndpointConfigurer>) null); - } - - public - IntegrationFlowBuilder split(EndpointConfigurer> endpointConfigurer) { + public IntegrationFlowBuilder split(EndpointConfigurer> endpointConfigurer) { return this.split(new DefaultMessageSplitter(), endpointConfigurer); } - public IntegrationFlowBuilder split(String expression) { - return this.split(expression, (EndpointConfigurer>) null); - } - public IntegrationFlowBuilder split(String expression, EndpointConfigurer> endpointConfigurer) { return this.split(new ExpressionEvaluatingSplitter(PARSER.parseExpression(expression)), endpointConfigurer); @@ -271,17 +311,21 @@ public final class IntegrationFlowBuilder { endpointConfigurer); } - public IntegrationFlowBuilder split(AbstractMessageSplitter splitter) { - return this.split(splitter, null); - } - - public IntegrationFlowBuilder split(GenericSplitter splitter) { - return this.split(splitter, null); + public

IntegrationFlowBuilder split(Class

payloadType, GenericSplitter

splitter) { + return this.split(payloadType, splitter, null); } public IntegrationFlowBuilder split(GenericSplitter splitter, EndpointConfigurer> endpointConfigurer) { - return this.split(new MethodInvokingSplitter(splitter, "split"), endpointConfigurer); + return split(null, splitter, endpointConfigurer); + } + + public

IntegrationFlowBuilder split(Class

payloadType, GenericSplitter

splitter, + EndpointConfigurer> endpointConfigurer) { + MethodInvokingSplitter split = isLambda(splitter) + ? new MethodInvokingSplitter(new LambdaMessageProcessor(splitter, payloadType)) + : new MethodInvokingSplitter(splitter, "split"); + return this.split(split, endpointConfigurer); } public IntegrationFlowBuilder split(S splitter, @@ -348,8 +392,7 @@ public final class IntegrationFlowBuilder { return this.resequence((EndpointConfigurer>) null); } - public - IntegrationFlowBuilder resequence(EndpointConfigurer> endpointConfigurer) { + public IntegrationFlowBuilder resequence(EndpointConfigurer> endpointConfigurer) { return this.resequence(new ResequencingMessageHandler(new ResequencingMessageGroupProcessor()), endpointConfigurer); } @@ -375,20 +418,11 @@ public final class IntegrationFlowBuilder { return this.handle(resequencer, endpointConfigurer); } - public IntegrationFlowBuilder aggregate() { - return this.aggregate((EndpointConfigurer>) null); - } - - public - IntegrationFlowBuilder aggregate(EndpointConfigurer> endpointConfigurer) { + public IntegrationFlowBuilder aggregate(EndpointConfigurer> endpointConfigurer) { return this.aggregate(new AggregatingMessageHandler(new DefaultAggregatingMessageGroupProcessor()), endpointConfigurer); } - public IntegrationFlowBuilder aggregate(ComponentConfigurer aggregatorConfigurer) { - return this.aggregate(aggregatorConfigurer, null); - } - public IntegrationFlowBuilder aggregate(ComponentConfigurer aggregatorConfigurer, EndpointConfigurer> endpointConfigurer) { Assert.notNull(aggregatorConfigurer); @@ -397,10 +431,6 @@ public final class IntegrationFlowBuilder { return this.aggregate(spec.get(), endpointConfigurer); } - public IntegrationFlowBuilder aggregate(AggregatingMessageHandler aggregator) { - return this.aggregate(aggregator, null); - } - public IntegrationFlowBuilder aggregate(AggregatingMessageHandler aggregator, EndpointConfigurer> endpointConfigurer) { return this.handle(aggregator, endpointConfigurer); @@ -441,23 +471,36 @@ public final class IntegrationFlowBuilder { } public IntegrationFlowBuilder route(GenericRouter router) { - return this.route(router, (ComponentConfigurer>) null); + return this.route(null, router); } public IntegrationFlowBuilder route(GenericRouter router, ComponentConfigurer> routerConfigurer) { - return this.route(router, routerConfigurer, null); + return this.route(null, router, routerConfigurer); + } + + public IntegrationFlowBuilder route(Class

payloadType, GenericRouter router) { + return this.route(payloadType, router, null, null); + } + + public IntegrationFlowBuilder route(Class

payloadType, GenericRouter router, + ComponentConfigurer> routerConfigurer) { + return this.route(payloadType, router, routerConfigurer, null); } public IntegrationFlowBuilder route(GenericRouter router, ComponentConfigurer> routerConfigurer, EndpointConfigurer> endpointConfigurer) { - return this.route(new MethodInvokingRouter(router), routerConfigurer, endpointConfigurer); + return route(null, router, routerConfigurer, endpointConfigurer); } - public IntegrationFlowBuilder route(R router, - ComponentConfigurer> routerConfigurer) { - return this.route(router, routerConfigurer, null); + public IntegrationFlowBuilder route(Class

payloadType, GenericRouter router, + ComponentConfigurer> routerConfigurer, + EndpointConfigurer> endpointConfigurer) { + MethodInvokingRouter methodInvokingRouter = isLambda(router) + ? new MethodInvokingRouter(new LambdaMessageProcessor(router, payloadType)) + : new MethodInvokingRouter(router); + return this.route(methodInvokingRouter, routerConfigurer, endpointConfigurer); } public IntegrationFlowBuilder route(R router, @@ -613,4 +656,9 @@ public final class IntegrationFlowBuilder { return this.flow; } + private static boolean isLambda(Object o) { + Class aClass = o.getClass(); + return aClass.isSynthetic() && !aClass.isAnonymousClass() && aClass.getDeclaredMethods().length == 1; + } + } diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/LambdaMessageProcessor.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/LambdaMessageProcessor.java new file mode 100644 index 0000000..42b7c1b --- /dev/null +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/LambdaMessageProcessor.java @@ -0,0 +1,111 @@ +/* + * Copyright 2014 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.dsl; + +import java.lang.reflect.Method; +import java.util.Map; + +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.BeanFactory; +import org.springframework.beans.factory.BeanFactoryAware; +import org.springframework.core.convert.ConversionService; +import org.springframework.core.convert.TypeDescriptor; +import org.springframework.core.convert.support.DefaultConversionService; +import org.springframework.integration.handler.MessageProcessor; +import org.springframework.integration.support.utils.IntegrationUtils; +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageHandlingException; +import org.springframework.util.Assert; + +/** + * @author Artem Bilan + */ +class LambdaMessageProcessor implements MessageProcessor, BeanFactoryAware { + + private final Object target; + + private final Method method; + + private final TypeDescriptor payloadType; + + private final Class[] parameterTypes; + + + private ConversionService conversionService; + + public LambdaMessageProcessor(Object target, Class payloadType) { + Assert.notNull(target); + this.target = target; + Method[] declaredMethods = target.getClass().getDeclaredMethods(); + Assert.isTrue(declaredMethods.length == 1, "LambdaMessageProcessor is applicable for inline or lambda" + + "classes with single method - functional interfaces implementations."); + this.method = declaredMethods[0]; + this.method.setAccessible(true); + this.parameterTypes = this.method.getParameterTypes(); + this.payloadType = TypeDescriptor.valueOf(payloadType); + } + + @Override + public void setBeanFactory(BeanFactory beanFactory) throws BeansException { + ConversionService conversionService = IntegrationUtils.getConversionService(beanFactory); + if (conversionService == null) { + conversionService = new DefaultConversionService(); + } + this.conversionService = conversionService; + } + + @Override + public Object processMessage(Message message) { + Object[] args = new Object[this.parameterTypes.length]; + for (int i = 0; i < this.parameterTypes.length; i++) { + Class parameterType = this.parameterTypes[i]; + if (Message.class.isAssignableFrom(parameterType)) { + args[i] = message; + } + if (Map.class.isAssignableFrom(parameterType)) { + if (message.getPayload() instanceof Map) { + args[i] = message.getPayload(); + } + else { + args[i] = message.getHeaders(); + } + } + else { + if (this.payloadType != null) { + if (Message.class.isAssignableFrom(this.payloadType.getType())) { + args[i] = message; + } + else { + args[i] = this.conversionService.convert(message.getPayload(), + TypeDescriptor.forObject(message.getPayload()), this.payloadType); + } + + } + else { + args[i] = message.getPayload(); + } + } + } + + try { + return this.method.invoke(this.target, args); + } + catch (Exception e) { + throw new MessageHandlingException(message, e); + } + } +} diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/GenericHandler.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/GenericHandler.java new file mode 100644 index 0000000..5bf275e --- /dev/null +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/support/GenericHandler.java @@ -0,0 +1,28 @@ +/* + * Copyright 2014 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.dsl.support; + +import java.util.Map; + +/** + * @author Artem Bilan + */ +public interface GenericHandler

{ + + Object handle(P payload, Map headers); + +} diff --git a/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/IntegrationFlowTests.java b/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/IntegrationFlowTests.java index c82fe01..bcdd86d 100644 --- a/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/IntegrationFlowTests.java +++ b/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/IntegrationFlowTests.java @@ -82,6 +82,7 @@ import org.springframework.integration.channel.FixedSubscriberChannel; import org.springframework.integration.channel.PublishSubscribeChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.config.EnableIntegration; +import org.springframework.integration.config.EnableMessageHistory; import org.springframework.integration.config.GlobalChannelInterceptor; import org.springframework.integration.core.MessageSource; import org.springframework.integration.dsl.IntegrationFlow; @@ -91,6 +92,7 @@ import org.springframework.integration.dsl.SplitterEndpointSpec; import org.springframework.integration.dsl.amqp.Amqp; import org.springframework.integration.dsl.channel.DirectChannelSpec; import org.springframework.integration.dsl.channel.MessageChannels; +import org.springframework.integration.dsl.support.GenericHandler; import org.springframework.integration.dsl.support.GenericSplitter; import org.springframework.integration.dsl.support.Pollers; import org.springframework.integration.endpoint.MethodInvokingMessageSource; @@ -372,9 +374,9 @@ public class IntegrationFlowTests { @Test public void testHandle() { assertNull(this.eventHolder.get()); - this.flow3Input.send(new GenericMessage<>("foo")); + this.flow3Input.send(new GenericMessage<>("2")); assertNotNull(this.eventHolder.get()); - assertEquals("foo", this.eventHolder.get()); + assertEquals(4, this.eventHolder.get()); } @Test @@ -548,7 +550,7 @@ public class IntegrationFlowTests { assertFalse(receive.getHeaders().containsKey("foo")); assertTrue(receive.getHeaders().containsKey("FOO")); assertEquals("BAR", receive.getHeaders().get("FOO")); - assertEquals(new Integer(i + 1), receive.getPayload()); + assertEquals(i + 1, receive.getPayload()); } } @@ -603,29 +605,29 @@ public class IntegrationFlowTests { @Test public void testRouter() { - int[] payloads = new int[]{1, 2, 3, 4, 5, 6}; + int[] payloads = new int[] {1, 2, 3, 4, 5, 6}; for (int payload : payloads) { - this.routerInput.send(new GenericMessage(payload)); + this.routerInput.send(new GenericMessage<>(payload)); } for (int i = 0; i < 3; i++) { Message receive = this.oddChannel.receive(2000); assertNotNull(receive); - assertEquals(new Integer(i * 2 + 1), receive.getPayload()); + assertEquals(i * 2 + 1, receive.getPayload()); receive = this.evenChannel.receive(2000); assertNotNull(receive); - assertEquals(new Integer(i * 2 + 2), receive.getPayload()); + assertEquals(i * 2 + 2, receive.getPayload()); } } @Test public void testMethodInvokingRouter() { - Message fooMessage = new GenericMessage("foo"); - Message barMessage = new GenericMessage("bar"); - Message badMessage = new GenericMessage("bad"); + Message fooMessage = new GenericMessage<>("foo"); + Message barMessage = new GenericMessage<>("bar"); + Message badMessage = new GenericMessage<>("bad"); this.routerMethodInput.send(fooMessage); @@ -684,9 +686,9 @@ public class IntegrationFlowTests { @Test public void testMethodInvokingRouter3() { - Message fooMessage = new GenericMessage("foo"); - Message barMessage = new GenericMessage("bar"); - Message badMessage = new GenericMessage("bad"); + Message fooMessage = new GenericMessage<>("foo"); + Message barMessage = new GenericMessage<>("bar"); + Message badMessage = new GenericMessage<>("bad"); this.routerMethod3Input.send(fooMessage); @@ -715,9 +717,9 @@ public class IntegrationFlowTests { @Test public void testMultiRouter() { - Message fooMessage = new GenericMessage("foo"); - Message barMessage = new GenericMessage("bar"); - Message badMessage = new GenericMessage("bad"); + Message fooMessage = new GenericMessage<>("foo"); + Message barMessage = new GenericMessage<>("bar"); + Message badMessage = new GenericMessage<>("bad"); this.routerMultiInput.send(fooMessage); Message result1a = this.fooChannel.receive(2000); @@ -750,7 +752,7 @@ public class IntegrationFlowTests { Message fooMessage = MessageBuilder.withPayload("fooPayload").setHeader("recipient", true).build(); Message barMessage = MessageBuilder.withPayload("barPayload").setHeader("recipient", true).build(); - Message badMessage = new GenericMessage("badPayload"); + Message badMessage = new GenericMessage<>("badPayload"); this.recipientListInput.send(fooMessage); Message result1a = this.fooChannel.receive(2000); @@ -986,7 +988,7 @@ public class IntegrationFlowTests { @Bean public MongoDbFactory mongoDbFactory() throws Exception { - return new SimpleMongoDbFactory(new MongoClient("localhost",mongoPort), "local"); + return new SimpleMongoDbFactory(new MongoClient("localhost", mongoPort), "local"); } @Bean @@ -999,7 +1001,7 @@ public class IntegrationFlowTests { @Bean public IntegrationFlow priorityFlow(PriorityCapableChannelMessageStore mongoDbChannelMessageStore) { return IntegrationFlows.from(MessageChannels.priority("priorityChannel", - mongoDbChannelMessageStore, "priorityGroup").interceptor()) + mongoDbChannelMessageStore, "priorityGroup").interceptor()) .bridge(s -> s.poller(Pollers.fixedDelay(1000, 2000))) .channel(MessageChannels.queue("priorityReplyChannel")) .get(); @@ -1047,6 +1049,7 @@ public class IntegrationFlowTests { @Bean public IntegrationFlow flow3() { return IntegrationFlows.from("flow3Input") + .handle(Integer.class, (p, h) -> p * 2) .handle(new ApplicationEventPublishingMessageHandler()) .get(); } @@ -1206,16 +1209,9 @@ public class IntegrationFlowTests { .enrichHeaders(s -> s.header("FOO", "BAR")) .split("testSplitterData", "buildList", c -> c.applySequence(false)) .channel(MessageChannels.executor(this.taskExecutor())) - .split(new GenericSplitter>>() { - - @Override - public Collection split(Message> target) { - return target.getPayload(); - } - }, c -> c.applySequence(false)) + .split(Message.class, target -> (List) target.getPayload(), c -> c.applySequence(false)) .channel(MessageChannels.executor(this.taskExecutor())) - .split((SplitterEndpointSpec s) -> - s.applySequence(false).get().getT2().setDelimiters(",")) + .split(s -> s.applySequence(false).get().getT2().setDelimiters(",")) .channel(MessageChannels.executor(this.taskExecutor())) .transform(Integer::parseInt) .enrichHeaders(s -> s.headerExpression(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER, "payload")) @@ -1227,10 +1223,10 @@ public class IntegrationFlowTests { @Bean public IntegrationFlow splitAggregateFlow() { return IntegrationFlows.fromFixedMessageChannel("splitAggregateInput") - .split() + .split(null) .channel(MessageChannels.executor(this.taskExecutor())) .resequence() - .aggregate() + .aggregate(null) .get(); } @@ -1261,8 +1257,7 @@ public class IntegrationFlowTests { .route(p -> p % 2 == 0, m -> m.suffix("Channel") .channelMapping("true", "even") - .channelMapping("false", "odd") - ) + .channelMapping("false", "odd")) .get(); } @@ -1295,9 +1290,8 @@ public class IntegrationFlowTests { @Bean public IntegrationFlow routeMultiMethodInvocationFlow() { return IntegrationFlows.from("routerMultiInput") - .route(p -> p.equals("foo") || p.equals("bar") ? new String[] {"foo", "bar"} : null, - s -> s.suffix("-channel") - ) + .route(String.class, p -> p.equals("foo") || p.equals("bar") ? new String[] {"foo", "bar"} : null, + s -> s.suffix("-channel")) .get(); } @@ -1324,7 +1318,7 @@ public class IntegrationFlowTests { public IntegrationFlow amqpFlow() { return IntegrationFlows.from(Amqp.inboundGateway(this.rabbitConnectionFactory, queue()).get()) .transform("hello "::concat) - .transform((String p) -> p.toUpperCase()) + .transform(String.class, String::toUpperCase) .get(); }