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();
}