Octov0.11.7
Guides

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

The 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: 30s

Limits

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.

On this page