From 91452275ad3d594d5041e6e7c420ee4228b3fad1 Mon Sep 17 00:00:00 2001 From: David Syer Date: Wed, 23 Jun 2010 16:54:50 +0000 Subject: [PATCH] INT-1128: Docs for MessageStoreReaper --- src/docbkx/aggregator.xml | 441 +++++++++++++++++++++++-------------- src/docbkx/resequencer.xml | 74 ++++--- 2 files changed, 319 insertions(+), 196 deletions(-) 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 .