Splitting and Aggregating
Process a collection element by element, then re-join the results.
Break one message into many with split, process each element independently, and
re-join them with aggregate. This also covers what changes when a flow becomes
asynchronous and how to answer the caller once it has. It follows
samples/split-aggregate.yaml.
Why not foreach
foreach in map mode is a splitter
and an aggregator fused into one block: it materializes the whole collection,
runs the body once per element in order, and writes each result back
positionally.
That fusion costs three things. The collection must fit in memory, because
items evaluates to an array before the loop starts. One element's error aborts
the loop, so a single bad record in fifty thousand fails the other 49,999. The
join is immediate and positional, so there is no re-joining on a condition and no
combining messages that were never one collection.
split and aggregate are the same two halves, separated. Use foreach for a
small collection whose result you want back on the same message.
Splitting an order into lines
flows:
- name: order
process:
- type: split
name: lines
items: "body.lines"
buildResponse:
process:
- type: set-payload
settings:
value: '{"accepted": vars.groupSize}'
# Everything below here runs once per line.
- type: set-payload
name: price-line
settings:
value: '{"sku": body.sku, "total": body.qty * body.price}'The blocks after split become the rest of each element's journey. Every
element runs them as its own invocation, concurrently, with its own event ID, so
one element failing does not affect the others: a line that fails to price hits
the flow's error: chain alone and the other lines carry on.
split needs no concurrency setting. Elements are scheduled onto the flow's
shared pool while it has room and run on the splitting worker when it does not,
so a split can never outrun what the flow can absorb. Size it with the flow's
pool setting.
Answering the caller
The real work has moved off the caller's thread, so without buildResponse the
caller gets whatever the message looked like when it was handed over.
buildResponse runs once, on the original message, after every element has been
dispatched, and what it produces is what the caller receives. It sees
vars.groupSize, the number of elements dispatched, so the receipt can state how
much work was accepted:
{ "accepted": 2 }Re-joining them
- type: aggregate
name: rejoin
completionTimeout: 10s
# And from here down, once per completed group.
- type: log
settings:
message: '"re-joined " + string(vars.aggregateCount) + " lines"'That is the whole configuration. aggregate pairs with split out of the box
because split stamps every element with groupId and groupSize, and those
are exactly what aggregate's correlation and completionSize default to.
The group completes when its count reaches its size. Elements are placed by the position the split gave them, so the re-joined array is in the order the lines were sent rather than the order they finished, and the same input produces the same array every run.
Set a completionTimeout on any re-join. If an element fails, the count never
reaches the size and nothing else will ever close the group: the timeout is the
only condition that fires without a message arriving. vars.aggregateReason
tells the rest of the flow whether it got a whole group (size) or a partial one
(timeout), and a missing element leaves a null in its position rather than
renumbering the survivors.
Aggregating without a split
aggregate holds state between messages, so it can also combine messages that
were never one collection: a hundred webhook events, or every event in a
five-second window.
- type: aggregate
correlation: "body.customerId"
completionSize: "100"
completionTimeout: "5s"All three completion conditions can be combined and the first to fire wins: size, timeout, and a predicate over the group so far.
completionExpression: "vars.group.count >= 5 && vars.group.idleMs > 2000"vars.group is the accumulated state (count, size, acc, ageMs,
idleMs), so a predicate can ask about the group rather than only the message in
hand.
Large groups
The default append strategy keeps every body and rewrites the whole group on
each message, so its stored size grows quadratically. For a big re-join, fold as
you go:
- type: aggregate
strategy: expression
expression: >
vars.group.acc == null ? body.total : vars.group.acc + body.totalThe fold sees the group before the current message, so acc is the previous
accumulator and null on the first. What it returns becomes the new one.
Running it
octo invoke --config samples/split-aggregate.yaml --flow order \
--data '{"orderId":"A-1","lines":[{"sku":"x","qty":2,"price":10.0},{"sku":"y","qty":1,"price":5.5}]}'The receipt comes back immediately, and the re-joined order appears in the logs once the lines have been priced and combined.
In a cluster
Group state lives in the runtime KV store under optimistic concurrency, so any replica can fold into a group and two replicas never each hold half of one. The timeout sweep is gated on leader election, so a group is reaped once per cluster. None of this is configurable, and standalone uses an in-process store and a permanent leader, so one configuration behaves the same way in both.
Group state is namespaced by storeKey, which defaults to the block's address in
the flow. If you rename or move the block, in-flight groups are stranded until
their deadline. Set it explicitly when you expect to edit a flow that runs long
windows:
- type: aggregate
storeKey: order-totals
completionTimeout: 30sLimits
split and aggregate must be top-level blocks in a flow's chain. Inside
another block's slot there is no "rest of the flow" to continue into, so nesting
one fails at startup.
There is no synchronous split-and-join yet: the flow after a split is
asynchronous, and the caller is answered from buildResponse. Use foreach in
map mode when you need the joined result on the same message.