skip to content

Explain @Splitter and @Aggregator, and how correlation and release strategies tie a split flow back together.

level: seniorimportance: must knowfreq 50%

answer

  1. splitter: one→many, stamps CORRELATION_ID/SEQUENCE_SIZE/NUMBER
  2. aggregator = stateful barrier over MessageGroupStore
  3. three strategies: Correlation / Release / GroupProcessor
  4. defaults: HeaderAttributeCorrelationStrategy + SequenceSizeReleaseStrategy
  5. 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 s

A 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 lines
java
import 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

for a junior

Know splitter = one→many, aggregator = many→one, and they pair up.

for a middle

Know the sequence headers the splitter adds and the aggregator's default correlation/release strategies.

for a senior

Explain the three pluggable strategies, MessageGroupStore, and configure group-timeout / partial-result handling to avoid leaks.

for a principal

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).

context