Explain @Splitter and @Aggregator, and how correlation and release strategies tie a split flow back together.
answer
- splitter: one→many, stamps CORRELATION_ID/SEQUENCE_SIZE/NUMBER
- aggregator = stateful barrier over MessageGroupStore
- three strategies: Correlation / Release / GroupProcessor
- defaults: HeaderAttributeCorrelationStrategy + SequenceSizeReleaseStrategy
- group-timeout to avoid leaked partial groups
basics
~20 s@Splitter breaks one message into many (e.g. a list into per-item messages). @Aggregator collects related messages back into one. They pair up: the aggregator uses a correlation strategy to group messages and a release strategy to decide when a group is complete and ready to emit.
solid answer
~40 sA splitter (@Splitter) takes one inbound message and emits multiple — typically returning a Collection/array/Iterator whose elements each become a message. Spring Integration automatically stamps sequence headers on the split messages: CORRELATION_ID (the original message id), SEQUENCE_SIZE, and SEQUENCE_NUMBER. An aggregator (@Aggregator) is the inverse: a stateful barrier that buffers messages in a MessageGroupStore keyed by a correlation key, then releases the group as a single combined message. It has three pluggable concerns: the CorrelationStrategy (which group a message belongs to — default HeaderAttributeCorrelationStrategy on CORRELATION_ID), the ReleaseStrategy (when the group is complete — default SequenceSizeReleaseStrategy, releasing when count == SEQUENCE_SIZE), and the aggregating method/MessageGroupProcessor (how to combine payloads). Because splitters set these headers, a downstream aggregator can reassemble the exact original set with zero config. Watch group-timeout, partial groups, and store persistence.
code
java · 38 linesimport org.springframework.integration.annotation.*;
import org.springframework.messaging.Message;
import java.util.List;
public class ScatterGather {
// SPLIT: each LineItem becomes its own message; framework stamps
// CORRELATION_ID (= original msg id), SEQUENCE_SIZE, SEQUENCE_NUMBER.
@Splitter(inputChannel = "orders", outputChannel = "lineItems")
public List<LineItem> split(Order order) {
return order.items();
}
// Custom correlation (optional): group by the original correlation id.
@CorrelationStrategy
public Object correlate(@Header("correlationId") Object id) {
return id;
}
// RELEASE: complete when we've collected SEQUENCE_SIZE messages.
@ReleaseStrategy
public boolean released(List<Message<LineItemResult>> group) {
Message<?> first = group.get(0);
int expected = (int) first.getHeaders().get("sequenceSize");
return group.size() == expected;
}
// AGGREGATE: combine the grouped payloads into one result.
@Aggregator(inputChannel = "lineItemResults", outputChannel = "orderResults")
public OrderResult combine(List<LineItemResult> results) {
return new OrderResult(results);
}
}
class Order { List<LineItem> items() { return List.of(); } }
class LineItem {}
class LineItemResult {}
class OrderResult { OrderResult(List<LineItemResult> r) {} }go deeper
Know splitter = one→many, aggregator = many→one, and they pair up.
Know the sequence headers the splitter adds and the aggregator's default correlation/release strategies.
Explain the three pluggable strategies, MessageGroupStore, and configure group-timeout / partial-result handling to avoid leaks.
Design reliable scatter-gather: persistent stores, correlation key lifecycle (expire-on-completion), straggler/duplicate handling, ordering under parallelism, and back-pressure.
## The Splitter/Aggregator pair These two EIP endpoints are complementary. A **Splitter** fans **one message out into many**; an **Aggregator** fans **many messages back into one**. Together they implement scatter/gather-style processing (split → process each in parallel → recombine). ## @Splitter ```java @Splitter(inputChannel = "orders", outputChannel = "lineItems") public List<LineItem> split(Order order) { return order.getLineItems(); } ``` - The method typically returns a **`Collection`, array, `Iterator`, `Iterable`, or `Stream`**; **each element becomes its own message** sent to the output channel. Returning a single object sends one message; returning a `Message<?>` list lets you control headers. - **Automatic sequence headers** — for every split message the framework sets: - `IntegrationMessageHeaderAccessor.CORRELATION_ID` = the *original* message's id, - `SEQUENCE_SIZE` = total number of split messages, - `SEQUENCE_NUMBER` = 1-based position. These are exactly what the default aggregator needs to reassemble the group. (You can disable via `applySequence=false` — then you must supply your own correlation.) ## @Aggregator — a stateful barrier ```java @Aggregator(inputChannel = "lineItems", outputChannel = "orderResults") public OrderResult combine(List<LineItemResult> results) { ... } ``` The aggregator **buffers** incoming messages until a group is complete, then invokes your method with the **list of grouped payloads** (or messages) and sends the single result. It composes **three strategies**: ### 1. CorrelationStrategy — "which group?" Computes a **correlation key** per message; messages with the same key form a `MessageGroup`. Default: **`HeaderAttributeCorrelationStrategy`** reading `CORRELATION_ID` (set by the splitter). Customize with `@CorrelationStrategy` on a method, a `CorrelationStrategy` bean, or a SpEL `correlation-strategy-expression`. ### 2. ReleaseStrategy — "is the group complete?" Decides when a buffered group is ready to emit. Default: **`SequenceSizeReleaseStrategy`** — releases when the number of collected messages equals `SEQUENCE_SIZE`. Alternatives: **`MessageCountReleaseStrategy`**, a `@ReleaseStrategy` method returning boolean over the `List`/`MessageGroup`, or a SpEL `release-strategy-expression` (e.g. `size() == 5`). ### 3. MessageGroupProcessor — "how to combine?" Your `@Aggregator` method (or a `MessageGroupProcessor`) turns the group into the output. Default processor concatenates payloads into a `List`. ## The MessageGroupStore Buffered groups live in a **`MessageGroupStore`** — by default **`SimpleMessageStore`** (in-memory). For durability/restart-safety you plug in `JdbcMessageStore`, `MongoDbMessageStore`, Redis, etc. This is where partially-complete groups sit waiting. ## Critical operational concerns (senior signal) - **Incomplete groups / timeouts**: if a release condition never becomes true (a message is lost), a group leaks forever in the store. Configure **`group-timeout`** (or `groupTimeoutExpression`) so a partial group is force-released or discarded after a delay; `send-partial-result-on-expiry` controls whether the partial group is emitted or discarded to `discard-channel`. - **`expire-groups-upon-completion` / `expireGroupsUponTimeout`**: whether a correlation key can be reused after release (else late messages for that key start a new group vs. are ignored). - **Ordering & parallelism**: splitting then processing on an executor channel means aggregation order isn't guaranteed; the aggregator doesn't require order but your combine logic must not assume it. - **Memory/persistence**: in-memory store loses partial groups on restart; use a persistent store for reliability. - **Late/duplicate messages**: after a group is released and the key expired, a straggler may create a new one-element group — guard with release/timeout config. ## When to use Split large batch payloads for per-item processing/back-pressure; aggregate responses from parallel calls (scatter-gather), collect a fixed batch, or window a stream. If you only need to *drop* or *reshape*, use filter/transformer instead.
- What headers does a splitter add, and why do they matter for the aggregator?CORRELATION_ID (original message id), SEQUENCE_SIZE, and SEQUENCE_NUMBER. The default aggregator uses CORRELATION_ID (HeaderAttributeCorrelationStrategy) to group and SEQUENCE_SIZE (SequenceSizeReleaseStrategy) to know when the group is complete — so a splitter and aggregator reassemble the original set with no custom config.
- A message in a group is lost so the release condition never fires. What happens and how do you handle it?The partial group sits in the MessageGroupStore forever (a leak). Configure group-timeout / groupTimeoutExpression to force-release or discard after a delay, and use send-partial-result-on-expiry + a discard-channel to decide whether to emit or drop the partial group.
- How do you make aggregation survive an application restart?Replace the default in-memory SimpleMessageStore with a persistent MessageGroupStore (JdbcMessageStore, MongoDbMessageStore, Redis) so partially-collected groups aren't lost on restart.
saying these in an interview costs you the question
- Thinking the aggregator is stateless — it buffers groups in a MessageGroupStore.
- Assuming groups always complete — without a timeout a lost message leaks the group forever.
- Believing you must always hand-code correlation — the splitter's CORRELATION_ID + default strategies reassemble automatically.
- Confusing SEQUENCE_NUMBER (position) with SEQUENCE_SIZE (total).