GH-3592: Scatter-Gather: applySeq=true by default
Fixes https://github.com/spring-projects/spring-integration/issues/3592 * Configure XML parser & Java DSL for Scatter-Gather, based on the `RecipientListRouter` to set an `applySequence` to `true` by default. This will make a `gatherer` part to fully rely on the default correlation strategies
This commit is contained in:
committed by
Gary Russell
parent
db287cf98f
commit
6d7aebc65a
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2014-2020 the original author or authors.
|
||||
* Copyright 2014-2022 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.
|
||||
@@ -89,7 +89,7 @@ public class ScatterGatherParser extends AbstractConsumerEndpointParser {
|
||||
builder.addConstructorArgReference(scatterChannel);
|
||||
}
|
||||
else {
|
||||
BeanDefinition scattererDefinition = null;
|
||||
BeanDefinition scattererDefinition;
|
||||
if (!hasScatterer) {
|
||||
scattererDefinition = new RootBeanDefinition(RecipientListRouter.class);
|
||||
}
|
||||
@@ -103,6 +103,11 @@ public class ScatterGatherParser extends AbstractConsumerEndpointParser {
|
||||
if (hasScatterer && scatterer.hasAttribute(ID_ATTRIBUTE)) {
|
||||
scattererId = scatterer.getAttribute(ID_ATTRIBUTE);
|
||||
}
|
||||
|
||||
if (!scatterer.hasAttribute("apply-sequence")) {
|
||||
scattererDefinition.getPropertyValues().addPropertyValue("applySequence", true);
|
||||
}
|
||||
|
||||
parserContext.getRegistry().registerBeanDefinition(scattererId, scattererDefinition); // NOSONAR not null
|
||||
builder.addConstructorArgValue(new RuntimeBeanReference(scattererId));
|
||||
}
|
||||
@@ -113,7 +118,7 @@ public class ScatterGatherParser extends AbstractConsumerEndpointParser {
|
||||
|
||||
Element gatherer = DomUtils.getChildElementByTagName(element, "gatherer");
|
||||
|
||||
BeanDefinition gathererDefinition = null;
|
||||
BeanDefinition gathererDefinition;
|
||||
if (gatherer == null) {
|
||||
try {
|
||||
gatherer = DOCUMENT_BUILDER_FACTORY.newDocumentBuilder().newDocument().createElement("aggregator");
|
||||
@@ -125,7 +130,7 @@ public class ScatterGatherParser extends AbstractConsumerEndpointParser {
|
||||
}
|
||||
gathererDefinition = GATHERER_PARSER.parse(gatherer, // NOSONAR
|
||||
new ParserContext(parserContext.getReaderContext(),
|
||||
parserContext.getDelegate(), scatterGatherDefinition));
|
||||
parserContext.getDelegate(), scatterGatherDefinition));
|
||||
String gathererId = id + ".gatherer";
|
||||
if (gatherer != null && gatherer.hasAttribute(ID_ATTRIBUTE)) {
|
||||
gathererId = gatherer.getAttribute(ID_ATTRIBUTE);
|
||||
|
||||
@@ -2749,6 +2749,7 @@ public abstract class BaseIntegrationFlowDefinition<B extends BaseIntegrationFlo
|
||||
* Populate a {@link ScatterGatherHandler} to the current integration flow position
|
||||
* based on the provided {@link RecipientListRouterSpec} for scattering function
|
||||
* and {@link AggregatorSpec} for gathering function.
|
||||
* For convenience, the {@link RecipientListRouterSpec#applySequence(boolean)} is set to true by default.
|
||||
* @param scatterer the {@link Consumer} for {@link RecipientListRouterSpec} to configure scatterer.
|
||||
* @param gatherer the {@link Consumer} for {@link AggregatorSpec} to configure gatherer.
|
||||
* @param scatterGather the {@link Consumer} for {@link ScatterGatherSpec} to configure
|
||||
@@ -2760,6 +2761,7 @@ public abstract class BaseIntegrationFlowDefinition<B extends BaseIntegrationFlo
|
||||
|
||||
Assert.notNull(scatterer, "'scatterer' must not be null");
|
||||
RecipientListRouterSpec recipientListRouterSpec = new RecipientListRouterSpec();
|
||||
recipientListRouterSpec.applySequence(true);
|
||||
scatterer.accept(recipientListRouterSpec);
|
||||
AggregatorSpec aggregatorSpec = new AggregatorSpec();
|
||||
if (gatherer != null) {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2016-2021 the original author or authors.
|
||||
* Copyright 2016-2022 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.
|
||||
@@ -857,7 +857,6 @@ public class RouterTests {
|
||||
public IntegrationFlow scatterGatherFlow() {
|
||||
return f -> f
|
||||
.scatterGather(scatterer -> scatterer
|
||||
.applySequence(true)
|
||||
.recipientFlow(m -> true, sf -> sf.handle((p, h) -> Math.random() * 10))
|
||||
.recipientFlow(m -> true, sf -> sf.handle((p, h) -> Math.random() * 10))
|
||||
.recipientFlow(m -> true, sf -> sf.handle((p, h) -> Math.random() * 10)),
|
||||
@@ -878,8 +877,7 @@ public class RouterTests {
|
||||
.scatterGather(
|
||||
scatterer -> scatterer
|
||||
.recipientFlow(f1 -> f1.handle((p, h) -> p + " - flow 1"))
|
||||
.recipientFlow(f2 -> f2.handle((p, h) -> p + " - flow 2"))
|
||||
.applySequence(true),
|
||||
.recipientFlow(f2 -> f2.handle((p, h) -> p + " - flow 2")),
|
||||
gatherer -> gatherer
|
||||
.outputProcessor(mg -> mg
|
||||
.getMessages()
|
||||
@@ -900,7 +898,6 @@ public class RouterTests {
|
||||
return f -> f
|
||||
.scatterGather(
|
||||
scatterer -> scatterer
|
||||
.applySequence(true)
|
||||
.recipientFlow(f1 -> f1.transform(p -> "Sub-flow#1"))
|
||||
.recipientFlow(f2 -> f2
|
||||
.channel(c -> c.executor(taskExecutor))
|
||||
@@ -923,7 +920,6 @@ public class RouterTests {
|
||||
public IntegrationFlow propagateErrorFromGatherer(TaskExecutor taskExecutor) {
|
||||
return IntegrationFlows.from(Function.class)
|
||||
.scatterGather(s -> s
|
||||
.applySequence(true)
|
||||
.recipientFlow(subFlow -> subFlow
|
||||
.channel(c -> c.executor(taskExecutor))
|
||||
.transform(p -> "foo")),
|
||||
@@ -943,9 +939,9 @@ public class RouterTests {
|
||||
|
||||
@Bean
|
||||
public IntegrationFlow scatterGatherInSubFlow() {
|
||||
return flow -> flow.scatterGather(s -> s.applySequence(true)
|
||||
return flow -> flow.scatterGather(s -> s
|
||||
.recipientFlow(inflow -> inflow.wireTap(scatterGatherWireTapChannel())
|
||||
.scatterGather(s1 -> s1.applySequence(true)
|
||||
.scatterGather(s1 -> s1
|
||||
.recipientFlow(IntegrationFlowDefinition::bridge)
|
||||
.recipientFlow("sequencetest"::equals,
|
||||
IntegrationFlowDefinition::bridge),
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
<bean id="messageStore" class="org.springframework.integration.store.SimpleMessageStore"/>
|
||||
|
||||
<int:scatter-gather id="scatterGather2" input-channel="input2" gather-channel="gatherChannel" gather-timeout="100">
|
||||
<int:scatterer id="myScatterer" apply-sequence="true">
|
||||
<int:scatterer id="myScatterer">
|
||||
<int:recipient channel="distributionChannel"/>
|
||||
</int:scatterer>
|
||||
<int:gatherer id="myGatherer" message-store="messageStore"/>
|
||||
|
||||
@@ -30,7 +30,7 @@
|
||||
|
||||
<!--Distribution scenario-->
|
||||
<scatter-gather input-channel="inputDistribution" output-channel="output" gather-channel="gatherChannel">
|
||||
<scatterer apply-sequence="true">
|
||||
<scatterer>
|
||||
<recipient channel="distribution1Channel"/>
|
||||
<recipient channel="distribution2Channel"/>
|
||||
<recipient channel="distribution3Channel"/>
|
||||
|
||||
Reference in New Issue
Block a user