diff --git a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/RouterSpec.java b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/RouterSpec.java index 1a610a1..b642934 100644 --- a/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/RouterSpec.java +++ b/spring-integration-java-dsl/src/main/java/org/springframework/integration/dsl/RouterSpec.java @@ -90,7 +90,6 @@ public final class RouterSpec extends Ab if (this.mappingProvider == null) { this.mappingProvider = new RouterSubFlowMappingProvider(this.target); - this.subFlows.add(this.mappingProvider); } this.mappingProvider.addMapping(key, channel); return _this(); @@ -98,6 +97,9 @@ public final class RouterSpec extends Ab @Override public Collection getComponentsToRegister() { + if (this.mappingProvider != null) { + this.subFlows.add(this.mappingProvider); + } return this.subFlows; } diff --git a/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/flows/IntegrationFlowTests.java b/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/flows/IntegrationFlowTests.java index bf54125..7b89955 100644 --- a/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/flows/IntegrationFlowTests.java +++ b/spring-integration-java-dsl/src/test/java/org/springframework/integration/dsl/test/flows/IntegrationFlowTests.java @@ -18,6 +18,7 @@ package org.springframework.integration.dsl.test.flows; import static org.hamcrest.Matchers.containsString; import static org.hamcrest.Matchers.instanceOf; +import static org.junit.Assert.assertArrayEquals; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; @@ -563,7 +564,27 @@ public class IntegrationFlowTests { assertNotNull(receive); assertEquals(payloads[i * 2 + 1], receive.getPayload()); } + } + @Autowired + @Qualifier("routerTwoSubFlows.input") + private MessageChannel routerTwoSubFlowsInput; + + @Autowired + @Qualifier("routerTwoSubFlowsOutput") + private PollableChannel routerTwoSubFlowsOutput; + + @Test + public void testRouterWithTwoSubflows() { + this.routerTwoSubFlowsInput.send(new GenericMessage(Arrays.asList(1, 2, 3, 4, 5, 6))); + Message receive = this.routerTwoSubFlowsOutput.receive(5000); + assertNotNull(receive); + Object payload = receive.getPayload(); + assertThat(payload, instanceOf(List.class)); + @SuppressWarnings("unchecked") + List results = (List) payload; + + assertArrayEquals(new Integer[] {3, 4, 9, 8, 15, 12}, results.toArray(new Integer[results.size()])); } @Test @@ -1191,6 +1212,18 @@ public class IntegrationFlowTests { .get(); } + + @Bean + public IntegrationFlow routerTwoSubFlows() { + return f -> f + .split() + .route(p -> p % 2 == 0, m -> m + .subFlowMapping("true", sf -> sf.handle((p, h) -> p * 2)) + .subFlowMapping("false", sf -> sf.handle((p, h) -> p * 3))) + .aggregate() + .channel(c -> c.queue("routerTwoSubFlowsOutput")); + } + @Bean public RoutingTestBean routingTestBean() { return new RoutingTestBean();