diff --git a/src/docbkx/aggregator.xml b/src/docbkx/aggregator.xml
index 209406979c..5efa549a06 100644
--- a/src/docbkx/aggregator.xml
+++ b/src/docbkx/aggregator.xml
@@ -14,9 +14,8 @@
Technically, the Aggregator is more complex than a Splitter, because
it is required to maintain state (the Messages to be aggregated), to
- decide when the complete group of Messages is available. In order to do this
- it requires a MessageStore
-
+ decide when the complete group of Messages is available. In order to do
+ this it requires a MessageStore
@@ -27,20 +26,20 @@
Aggregator will create a single message by processing the whole group, and
will send that aggregated message as output.
- An main aspect of
- implementing an Aggregator is providing the logic that has to be executed
- when the aggregation (creation of a single message out of many) takes
- place. The other two aspects are correlation and release
+ An main aspect of implementing an Aggregator is providing the logic
+ that has to be executed when the aggregation (creation of a single message
+ out of many) takes place. The other two aspects are correlation and
+ release
- In Spring Integration, the grouping of the messages for aggregation (correlation)
- is done by default based on their CORRELATION_ID message header (i.e. the
- messages with the same CORRELATION_ID will be grouped together). However,
- this can be customized, and the users can opt for other ways of
- specifying how the messages should be grouped together, by using a
- CorrelationStrategy (see below).
+ In Spring Integration, the grouping of the messages for aggregation
+ (correlation) is done by default based on their CORRELATION_ID message
+ header (i.e. the messages with the same CORRELATION_ID will be grouped
+ together). However, this can be customized, and the users can opt for
+ other ways of specifying how the messages should be grouped together, by
+ using a CorrelationStrategy (see below).
- To determine whether or not a group of messages may be processed,
- a ReleaseStrategy is consulted. The default release strategy for aggregator
+ To determine whether or not a group of messages may be processed, a
+ ReleaseStrategy is consulted. The default release strategy for aggregator
will release groups that have all messages from the sequence, but this can
be entirely customized
@@ -52,9 +51,10 @@
- The interface MessageGroupProcessor and related
- base class AbstractAggregatingMessageGroupProcessor and its
- subclass MethodInvokingAggregatingMessageGroupProcessor
+ The interface MessageGroupProcessor and related
+ base class AbstractAggregatingMessageGroupProcessor and
+ its subclass
+ MethodInvokingAggregatingMessageGroupProcessor
@@ -69,76 +69,106 @@
+
+
CorrelatingMessageHandler
+
+
The CorrelatingMessageHandler is a
MessageHandler implementation, encapsulating the common
- functionalities of an Aggregator (and other correlating use cases), which are:
-
-
- correlating messages into a group to be aggregated
-
-
- maintaining those messages in a MessageStore until the group may be released
-
-
- deciding when the group is in fact may be released
-
-
- processing the released group into a single aggregated message
-
-
- recognizing and responding to an expired group
-
-
- The responsibility of deciding how the messages should be grouped together
- is delegated to a CorrelationStrategy instance. The responsibility
- of deciding whether the message group can be released is delegated to a
- ReleaseStrategy instance.
+ functionalities of an Aggregator (and other correlating use cases),
+ which are:
+
+ correlating messages into a group to be aggregated
+
+
+
+ maintaining those messages in a MessageStore until the group
+ may be released
+
+
+
+ deciding when the group is in fact may be released
+
+
+
+ processing the released group into a single aggregated
+ message
+
+
+
+ recognizing and responding to an expired group
+
+ The responsibility of deciding how the messages should
+ be grouped together is delegated to a CorrelationStrategy
+ instance. The responsibility of deciding whether the message group can
+ be released is delegated to a ReleaseStrategy
+ instance.
+
+
Here is a brief highlight of the base
- AbstractAggregatingMessageGroupProcessor (the responsibility of
- implementing the aggregateMessages method is left to the
- developer):
+ AbstractAggregatingMessageGroupProcessor (the
+ responsibility of implementing the aggregateMessages method is left to
+ the developer):
- public abstract class AbstractAggregatingMessageGroupProcessor
+
+
+ aggregateHeaders(MessageGroup group) {
....
}
protected abstract Object aggregatePayloads(MessageGroup group);
-}
- The CorrelationStrategy is owned by the CorrelatingMessageHandler and it has
- a default value based on the correlation ID message header:
- private volatile CorrelationStrategy correlationStrategy =
- new HeaderAttributeCorrelationStrategy(MessageHeaders.CORRELATION_ID);
+}]]>
- When appropriate, the simplest option is the DefaultAggregatingMessageGroupProcessor.
- It creates a single Message whose payload is a List of the payloads received
- for a given group. It uses the default CorrelationStrategy and
- CompletionStrategy as shown above. This works well for simple
- Scatter Gather implementations with either a Splitter, Publish Subscribe Channel,
- or Recipient List Router upstream.
+ The CorrelationStrategy is owned by the
+
+ CorrelatingMessageHandler
+
+ and it has a default value based on the correlation ID message header:
+
+
+
+
+
+ When appropriate, the simplest option is the
+ DefaultAggregatingMessageGroupProcessor. It creates a
+ single Message whose payload is a List of the payloads received for a
+ given group. It uses the default CorrelationStrategy and
+ CompletionStrategy as shown above. This works well for
+ simple Scatter Gather implementations with either a Splitter, Publish
+ Subscribe Channel, or Recipient List Router upstream.
+
+
- When using a Publish Subscribe Channel or Recipient List Router in this
- type of scenario, be sure to enable the flag to apply-sequence.
- That will add the necessary headers (correlation id, sequence number and sequence
- size). That behavior is enabled by default for Splitters in Spring Integration,
- but it is not enabled for the Publish Subscribe Channel or Recipient List
- Router because those components may be used in a variety of contexts where
- those headers are not necessary.
+ When using a Publish Subscribe Channel or Recipient List Router
+ in this type of scenario, be sure to enable the flag to
+ apply-sequence. That will add the necessary
+ headers (correlation id, sequence number and sequence size). That
+ behavior is enabled by default for Splitters in Spring Integration,
+ but it is not enabled for the Publish Subscribe Channel or Recipient
+ List Router because those components may be used in a variety of
+ contexts where those headers are not necessary.
-
+
+
+
When implementing a specific aggregator object for an application,
- a developer can extend AbstractAggregatingMessageGroupProcessor and
- implement the aggregatePayloads method. However, there are
- better suited (which reads, less coupled to the API) solutions for
- implementing the aggregation logic, which can be configured easily
- either through XML or through annotations.
+ a developer can extend
+ AbstractAggregatingMessageGroupProcessor and implement the
+ aggregatePayloads method. However, there are better suited
+ (which reads, less coupled to the API) solutions for implementing the
+ aggregation logic, which can be configured easily either through XML or
+ through annotations.
+
+
In general, any ordinary Java class (i.e. POJO) can implement the
aggregation algorithm. For doing so, it must provide a method that
@@ -146,6 +176,8 @@
supported as well). This method will be invoked for aggregating
messages, as follows:
+
+
if the argument is a parametrized java.util.List, and the
@@ -167,12 +199,16 @@
+
+
In the interest of code simplicity, and promoting best practices
such as low coupling, testability, etc., the preferred way of
implementing the aggregation logic is through a POJO, and using the
XML or annotation support for setting it up in the application.
+
+
@@ -181,11 +217,11 @@
The ReleaseStrategy interface is defined as
follows:
- public interface ReleaseStrategy {
+
+}]]>
In general, any ordinary Java class (i.e. POJO) can implement the
completion decision mechanism. For doing so, it must provide a method
@@ -208,27 +244,27 @@
- the method must return true if the message group is ready
- for aggregation, and false otherwise.
+ the method must return true if the message group is ready for
+ aggregation, and false otherwise.
- When the group is released for aggregation, all its
- unmarked messages are processed and then marked so they will not
- be processed again. If the group is also complete (i.e. if all
- messages from a sequence have arrived or if there is no sequence
- defined) then the group is removed from the message store.
- Partial sequences can be released, in which case the next time
- the ReleaseStrategy is called it will be presented
- with a group containing marked messages (already processed) and
- unmarked messages (a potential new partial sequence)
+ When the group is released for aggregation, all its unmarked
+ messages are processed and then marked so they will not be processed
+ again. If the group is also complete (i.e. if all messages from a
+ sequence have arrived or if there is no sequence defined) then the group
+ is removed from the message store. Partial sequences can be released, in
+ which case the next time the ReleaseStrategy is called it
+ will be presented with a group containing marked messages (already
+ processed) and unmarked messages (a potential new partial
+ sequence)
Spring Integration provides an out-of-the box implementation for
ReleaseStrategy, the
- SequenceSizerReleaseStrategy. This implementation uses
- the SEQUENCE_NUMBER and SEQUENCE_SIZE of the arriving messages for
- deciding when a message group is complete and ready to be
- aggregated. As shown above, it is also the default strategy.
+ SequenceSizerReleaseStrategy. This implementation uses the
+ SEQUENCE_NUMBER and SEQUENCE_SIZE of the arriving messages for deciding
+ when a message group is complete and ready to be aggregated. As shown
+ above, it is also the default strategy.
@@ -237,11 +273,11 @@
The CorrelationStrategy interface is defined as
follows:
- public interface CorrelationStrategy {
+ message);
-}
+}]]>
The method shall return an Object which represents the correlation
key used for grouping messages together. The key must satisfy the
@@ -249,8 +285,8 @@
equals() and hashCode().
In general, any ordinary Java class (i.e. POJO) can implement the
- correlation decision mechanism, and the rules for mapping a message to
- a method's argument (or arguments) are the same as for a
+ correlation decision mechanism, and the rules for mapping a message to a
+ method's argument (or arguments) are the same as for a
ServiceActivator (including support for @Header
annotations). The method must return a value, and the value must not be
null.
@@ -269,34 +305,34 @@
Configuring an Aggregator with XML
Spring Integration supports the configuration of an aggregator via
- XML through the <aggregator/> element. Below you can see an example of
- an aggregator with all optional parameters defined.
+ XML through the <aggregator/> element. Below you can see an example
+ of an aggregator with all optional parameters defined.
- <channel id="inputChannel"/>
+
-<aggregator id="completelyDefinedAggregator"
- input-channel="inputChannel"
- output-channel="outputChannel"
- discard-channel="discardChannel"
- ref="aggregatorBean"
- method="add"
- release-strategy="releaseStrategyBean"
- release-strategy-method="canRelease"
- correlation-strategy="correlationStrategyBean"
- correlation-strategy-method="groupNumbersByLastDigit"
- message-store="messageStore"
- send-partial-result-on-expiry="true"
- send-timeout="86420000" />
+
-<channel id="outputChannel"/>
+
-<bean id="aggregatorBean" class="sample.PojoAggregator"/>
+
-<bean id="releaseStrategyBean" class="sample.PojoReleaseStrategy"/>
+
-<bean id="correlationStrategyBean" class="sample.PojoCorrelationStrategy"/>
+]]>
@@ -368,49 +404,49 @@
present).
-
- A reference to a MessageGroupStore that
- can be used to store groups of messages under their
- correlation key until they are
- complete. Optional with default a
- volatile in-memory store.
-
+
+ A reference to a MessageGroupStore that can be used
+ to store groups of messages under their correlation key until they are
+ complete. Optional with default a volatile
+ in-memory store.
+
- Whether upon the expiration of the message group, the aggregator will
- try to aggregate the messages that have already arrived. Optional
- (false by default).
+ Whether upon the expiration of the message group, the aggregator
+ will try to aggregate the messages that have already arrived.
+ Optional (false by default).
- The timeout for sending the aggregated messages to the
- output or reply channel. Optional.
+ The timeout for sending the aggregated messages to the output or
+ reply channel. Optional.
- Using a "ref" attribute is generally recommended if a custom aggregator handler
- implementation can be reused in other <aggregator> definitions.
- However if a custom aggregator handler implementation should be scoped to a concrete
- definition of the <aggregator>, you can use an inner bean definition
- (starting with version 1.0.3) for custom aggregator handlers within the
- <aggregator> element:
-
+ Using a "ref" attribute is generally recommended if a custom
+ aggregator handler implementation can be reused in other
+ <aggregator> definitions. However if a custom
+ aggregator handler implementation should be scoped to a concrete
+ definition of the <aggregator>, you can use an inner
+ bean definition (starting with version 1.0.3) for custom aggregator
+ handlers within the <aggregator> element:
+
-]]>
-
+]]>
- Using both a "ref" attribute and an inner bean definition in the same
- <aggregator> configuration is not allowed, as it creates an
- ambiguous condition. In such cases, an Exception will be thrown.
-
+ Using both a "ref" attribute and an inner bean definition in the
+ same <aggregator> configuration is not allowed, as it
+ creates an ambiguous condition. In such cases, an Exception will be
+ thrown.
- An example implementation of the aggregator bean looks as follows:
+ An example implementation of the aggregator bean looks as
+ follows:
- public class PojoAggregator {
+ results) {
long total = 0l;
for (long partialResult: results) {
total += partialResult;
@@ -418,37 +454,34 @@
return total;
}
-}
+}]]>
An implementation of the completion strategy bean for the example
above may be as follows:
- public class PojoReleaseStrategy {
+ numbers) {
int sum = 0;
for (long number: numbers) {
sum += number;
}
- return sum >= maxValue;
+ return sum >= maxValue;
}
-}
-
-
- Wherever it makes sense, the release strategy method and
- the aggregator method can be combined in a single bean.
-
-
+}]]>
+ Wherever it makes sense, the release strategy method and the
+ aggregator method can be combined in a single bean.
+
An implementation of the correlation strategy bean for the example
above may be as follows:
- public class PojoCorrelationStrategy {
+
+}]]>
For example, this aggregator would group numbers by some criterion
(in our case the remainder after dividing by 10) and will hold the group
@@ -457,36 +490,109 @@
Wherever it makes sense, the release strategy method, correlation
- strategy method and the aggregator method can be combined in a single bean
- (all of them or any two).
+ strategy method and the aggregator method can be combined in a single
+ bean (all of them or any two).
+
+ Managing State in an Aggregator:
+ MessageGroupStore
+
+ Aggregator (and some other patterns in Spring Integration) is a
+ stateful pattern that requires decisions to be made based on a group of
+ messages that have arrived over a period of time, all with the same
+ correlation key. The design of the interfaces in the stateful patterns
+ (e.g. ReleaseStrategy) is driven by the principle
+ that the components (framework and user) should be to remain stateless.
+ All state is carried by the MessageGroup and its
+ management is delegated to the
+ MessageGroupStore.
+
+ The MessageGroupStore accumulates state
+ information in MessageGroups, potentially forever.
+ So to prevent stale state from hanging around, and for volatile stores to
+ provide a hook for cleaning up when the application shots down, the
+ MessageGroupStore allows the user to register
+ callbacks to apply to MessageGroups when they
+ expire. The interface is very straighforward:
+
+
+
+ The callback has access directly to the store and the message group
+ so it can manage the persistent state (e.g. by removing the group from the
+ store entirely).
+
+ The MessageGroupStore maintains a list of these callbacks which it
+ applies when asked to all messages whose timestamp is earlier than a time
+ supplied as a parameter:
+
+
+
+ The expireMessageGroups method can be called with a timeout value:
+ any message older than the current time minus this value wiull be expired,
+ and have the callbacks applied. Thus it is the user of the store that
+ defines what is meant by message group "expiry".
+
+ As a convenience for users, Spring Integration provides a wrapper
+ for the message expiry in the form of a
+ MessageGroupStoreReaper:
+
+
+
+
+
+
+
+
+]]>
+
+ The reaper is a Runnable, and all that is happening is that the
+ message group store's expire method is being called in the sample above
+ once every 10 seconds. In addition to the reaper, the expiry callbacks are
+ invoked when the application shuts down via a lifecycle callback in the
+ CorrelatingMessageHandler.
+
+ The CorrelatingMessageHandler registers its
+ own expiry callback, and this is the link with the boolean flag
+ send-partial-result-on-expiry in the XML configuration of the
+ aggregator. If the flag is set to true, then when the expiry callback is
+ invoked then any unmarked messages in groups that are not yet released can
+ be sent on to the downstream channel.
+
+
Configuring an Aggregator with Annotations
An aggregator configured using annotations can look like
this.
- public class Waiter {
+
- public Delivery aggregatingMethod(List<OrderItem> items) {
+ @Aggregator ]]> items) {
...
}
- @ReleaseStrategy
- public boolean releaseChecker(List<Message<?>> messages) {
+ @ReleaseStrategy ]]>> messages) {
...
}
- @CorrelationStrategy
+ @CorrelationStrategy ]]>
+}]]>
@@ -497,8 +603,8 @@
An annotation indicating that this method shall be
- used as the release strategy of an aggregator. If not present on
- any method, the aggregator will use the
+ used as the release strategy of an aggregator. If not present on any
+ method, the aggregator will use the
SequenceSizeCompletionStrategy.
@@ -510,12 +616,11 @@
- All of the configuration options provided by the xml element are also
- available for the @Aggregator annotation.
+ All of the configuration options provided by the xml element are
+ also available for the @Aggregator annotation.
The aggregator can be either referenced explicitly from XML or, if
the @MessageEndpoint is defined on the class, detected automatically
through classpath scanning.
-
diff --git a/src/docbkx/resequencer.xml b/src/docbkx/resequencer.xml
index 8c5034723f..acbdedf8ef 100644
--- a/src/docbkx/resequencer.xml
+++ b/src/docbkx/resequencer.xml
@@ -8,7 +8,7 @@
Introduction
Related to the Aggregator, albeit different from a functional
- standpoint, is the Resequencer.
+ standpoint, is the Resequencer.
@@ -16,9 +16,9 @@
The Resequencer works in a similar way to the Aggregator, in the
sense that it uses the CORRELATION_ID to store messages in groups, the
- difference being that the Resequencer does not process the messages in
- any way. It simply releases them in the order of their SEQUENCE_NUMBER
- header values.
+ difference being that the Resequencer does not process the messages in any
+ way. It simply releases them in the order of their SEQUENCE_NUMBER header
+ values.
With respect to that, the user might opt to release all messages at
once (after the whole sequence, according to the SEQUENCE_SIZE, has been
@@ -33,18 +33,20 @@
A sample resequencer configuration is shown below.
- <channel id="inputChannel"/>
+
-<channel id="outputChannel"/>
+
-<resequencer id="completelyDefinedResequencer"
- input-channel="inputChannel"
- output-channel="outputChannel"
- discard-channel="discardChannel"
- release-partial-sequences="true"
- message-store="messageStore"
- send-partial-result-on-expiry="true"
- send-timeout="86420000" />
+ ]]>
@@ -69,31 +71,47 @@
+
+
Whether to send out ordered sequences as soon as they are
available, or only after the whole message group arrives.
- Optional (false by default). If this
- flag is not specified (so a complete sequence is defined by the sequence
- headers) then it can make sense to provide a custom Comparator
- to be used to order the messages when sending
- (use the XML attribute comparator to point to a bean
- definition). If release-partial-sequences is true then
- there is no way with a custom comparator to define a partial sequence. To
- do that you would have to provide a release-strategy (also
- a reference to another bean definition, either a POJO or a ReleaseStrategy).
+ Optional (false by default).
+
+ If this flag is not specified (so a complete sequence is defined by the sequence headers) then it can make sense to provide a custom
+
+ Comparator
+
+ to be used to order the messages when sending (use the XML attribute
+
+ comparator
+
+ to point to a bean definition). If
+
+ release-partial-sequences
+
+ is true then there is no way with a custom comparator to define a partial sequence. To do that you would have to provide a
+
+ release-strategy
+
+ (also a reference to another bean definition, either a POJO or a
+
+ ReleaseStrategy
+
+ ).
- A reference to a MessageGroupStore that
- can be used to store groups of messages under their
- correlation key until they are
- complete. Optional with default a
+ A reference to a MessageGroupStore that can be
+ used to store groups of messages under their correlation key until
+ they are complete. Optional with default a
volatile in-memory store.
Whether, upon the expiration of the group, the ordered group
should be sent out (even if some of the messages are missing).
- Optional (false by default).
+ Optional (false by default). See .