Multi-Flow Design

Multi-flow design combines multiple flows to handle scenarios that a single flow cannot handle effectively, such as fan-out to multiple destinations, batch splitting, and content-based routing.

This page introduces key patterns and the processors that implement them. For a full reference of processor properties, follow the links to the relevant processor documentation.

Distribution

Fan-Out to Multiple Destinations

Consider a Sales system that receives an order event and must forward it to Finance, Delivery, and Reporting systems. A single-flow implementation places all three restRequest processors sequentially in one pipeline:

flowConfig {
    id = "order-flow"
    ownerId = OWNER_ID
    exchangePattern = RequestResponse
    restApi { id = "order-api"; ... }
    restRequest { id = "finance";   defaultMethod = POST; address = URL("https://OVERRIDE_ME/finance") }
    restRequest { id = "delivery";  defaultMethod = POST; address = URL("https://OVERRIDE_ME/delivery") }
    restRequest { id = "reporting"; defaultMethod = POST; address = URL("https://OVERRIDE_ME/reporting") }
}

This design tightly couples the downstream systems and can unintentionally pass downstream responses between processors.

The distribute Processor

By default, the distribute processor can forward the same message to multiple target flows simultaneously. Distribution is not complete until all target flows confirm reception.

flowConfig {
    id = "order-distribution-api"
    ownerId = OWNER_ID
    exchangePattern = RequestResponse
    restApi { id = "order-api"; ... }
    distribute {
        id = "distribute-order"
        flowId(FLOW_ID_FINANCE)
        flowId(FLOW_ID_DELIVERY)
        flowId(FLOW_ID_REPORTING)
    }
}

Valid Target Flows

Each flow referenced by distribute must be either a Headless flow or a flow whose source is a handoff-family processor (handoff, scheduledHandoff, or windowedAggregationHandoff).

Headless flow

Use this when the main flow must wait for all target processing to complete or high performance is required. The target pipeline is inlined in the main flow and shares the main flow exchange pattern and reliability configuration.

val financeTargetSpec = flowConfig {
    id = FLOW_ID_FINANCE_TARGET
    ownerId = OWNER_ID
    exchangePattern = Headless
    restRequest {
        id = "send-to-finance"
        defaultMethod = POST
        address = URL("https://OVERRIDE_ME/finance")
    }
}

OneWay + handoff Flow

Use this when targets should be decoupled from the main flow and from each other. The main flow completes once each message is handed off to the target source. Each target has its own persistence, redelivery, and dead letter handling.

val financeTargetSpec = flowConfig {
    id = FLOW_ID_FINANCE_TARGET
    ownerId = OWNER_ID
    exchangePattern = OneWay
    handoff { id = "order-source" }
    restRequest {
        id = "send-to-finance"
        defaultMethod = POST
        address = URL("https://OVERRIDE_ME/finance")
    }
}
Characteristic Headless OneWay + handoff

Main flow waits until the target…​

Completes its pipeline

Accepts the message at its source

Independent delivery guarantee

No; the main flow provides the guarantee

Yes; includes redelivery and dead-letter handling

Independent error handling

No

Yes

Persistence

No additional persistence; the main flow handles it

The handoff source persists the message, adding a second write to reliable storage for each distributed message

Use case

Synchronous fan-out, when the main flow must wait for all targets or high performance is required

Independent downstream systems that require decoupled processing

Each OneWay + handoff target flow must persist every incoming message before processing. That persistence is what gives each target flow its own delivery guarantee, redelivery, and dead letter handling. If you do not need that extra guarantee and decoupling, Headless target flows avoid this overhead and are the faster option.

Fan-In from Multiple Sources

The inverse pattern is also supported: multiple source flows distribute to a shared handoff flow. Consider a scenario where you need to consolidate data from several weather services. Some of these services provide XML, while others provide CSV files. In both cases, the data must be converted to JSON before reaching the final destination, which is a database.

The source flows can use different connectors to retrieve data and perform the necessary format conversions before distributing to the shared handoff flow. The shared handoff flow then performs the final transformation and writes to the database target. A dead letter strategy on the shared handoff flow can capture database write failures.

For fan-in, the shared target flow should be OneWay + handoff. Headless is inlined into each upstream flow, so it cannot provide a single shared consolidation point.

Design Considerations for Distribution

  • Performance: Headless target flows are faster than OneWay + handoff because they avoid the extra persistence and handoff overhead. Avoid distribution if the scenario can be implemented with a single flow.

  • Decoupling: With OneWay + handoff, slow or failed target processing does not block the main flow or other targets.

  • Non-persistence and idempotency: Distribution state is non-persistent, and redelivery can produce duplicates.

  • Main flow result: distribute returns one representative response (not an aggregate), in this order:

    1. Non-transient failure (any target)

    2. Transient failure (any target)

    3. Successful response

    4. Filtered message

Using RequestResponse on a handoff flow is a rare pattern that increases latency and resilience complexity. See handoff for full implications.

Split Collections

Scenario

The Sales system sends a batch (for example ['order1', 'order2']), but Finance only accepts one order per request.

The split Processor

The split processor divides a collection payload into individual fragments and sends each fragment as a separate message to one target flow. Processing completes only when every fragment has a response. Like distribute, split state is non-persistent.

By default, each element in the collection becomes one fragment — a collection of 10 orders produces 10 fragments. Set chunkSize to N to group elements into list fragments instead: each fragment contains a list of at most N elements, and the last fragment may contain fewer. For example, chunkSize = 3 on a collection of 10 orders produces 4 fragments: three containing 3 orders each and one containing the remaining 1. Use this when the target system accepts small batches but not individual records.

Valid Target Flows

The flow referenced by split must be either a Headless flow or a flow whose source is a handoff-family processor (handoff, scheduledHandoff, or windowedAggregationHandoff).

Headless Flow

The main flow waits for all fragments to finish and for all responses to return. If you want the split processor to receive a composite result, add aggregate to the end of the target pipeline.

val financeTargetSpec = flowConfig {
    id = FLOW_ID_FINANCE_TARGET
    ownerId = OWNER_ID
    exchangePattern = Headless
    restRequest {
        id = "send-to-finance"
        defaultMethod = POST
        address = URL("https://OVERRIDE_ME/finance")
    }
    aggregate {
        id = "collect-results"
        maxNoOfMessagesToKeep = 100
        maxAgeInSecondsForMessagesToKeep = 600
    }
}
// Result shape:
// { "numberOfFragments": 2, "fragments": ["result1", "result2"] }

OneWay + handoff Flow

The main flow receives a handoff receipt for each fragment as soon as the target source accepts it. The target flow processes the fragments independently.

val financeTargetSpec = flowConfig {
    id = FLOW_ID_FINANCE_TARGET
    ownerId = OWNER_ID
    exchangePattern = OneWay
    handoff { id = "order-source" }
    restRequest {
        id = "send-to-finance"
        defaultMethod = POST
        address = URL("https://OVERRIDE_ME/finance")
    }
}

The aggregate Processor

Use the aggregate processor only in the following cases:

Collect Fragment Results in the Main Flow

Use a Headless target + aggregate to return { "numberOfFragments": N, "fragments": […​] }.

Run Post-Split Processing After All Fragments Finish

Use a OneWay target + handoff and place aggregate as a synchronization point.

val financeTargetSpec = flowConfig {
    id = FLOW_ID_FINANCE_TARGET
    ownerId = OWNER_ID
    exchangePattern = OneWay
    handoff { id = "order-source" }
    restRequest {
        id = "send-to-finance"
        defaultMethod = POST
        address = URL("https://OVERRIDE_ME/finance")
    }
    aggregate {
        id = "sync-point"
        maxNoOfMessagesToKeep = 100
        maxAgeInSecondsForMessagesToKeep = 600
    }
    restRequest {
        id = "send-report"
        defaultMethod = POST
        address = URL("https://OVERRIDE_ME/reporting")
    }
}
If the target flow does not need to be OneWay, use a Headless target without aggregate and place post-split logic directly in the main flow after split.
flowConfig {
    id = "main-flow"
    ownerId = OWNER_ID
    exchangePattern = OneWay
    restApi { id = "order-api"; ... }
    setPayload { id = "create-list"; source = "envelope.payload.value" }
    split { id = "split-orders"; flowId = FLOW_ID_FINANCE_TARGET }
    restRequest { id = "send-report"; ... }
}

val financeTargetSpec = flowConfig {
    id = FLOW_ID_FINANCE_TARGET
    ownerId = OWNER_ID
    exchangePattern = Headless
    restRequest { id = "send-to-finance"; ... }
}

Key Design Implications

  • Performance: A Headless target flow is faster than OneWay + handoff because it avoids the extra persistence and handoff overhead.

  • Decoupling: With OneWay + handoff, slow or failed target processing does not block the main flow.

  • Non-persistence and idempotency: Split state is non-persistent, and duplicates can occur on redelivery.

  • Main flow result: Similar to distribute, split returns one representative response. aggregate returns a composite result only in a Headless target flow.

Content-Based Routing

Overview

Content-based routing directs messages to different target flows based on their content. In Connect, this is implemented using the distribute processor with targetFilterExpression.

Filter expressions run in the Standard Expression Context extended with one additional property: targetFlowId — the ID of the target flow currently being evaluated. This means expressions can filter on the payload directly, on exchange properties, on the target flow ID, or on any combination of these.

For example, to route messages to Shipping or Finance based on payload content, reference the payload field directly:

distribute {
    id = "distribute-message"
    flowId(FLOW_ID_SHIPPING)
    flowId(FLOW_ID_FINANCE)
    targetFilterExpression(
        "envelope.payload.value == 'address' and targetFlowId matches '.*shipping.*'"
    )
    targetFilterExpression(
        "envelope.payload.value == 'price' and targetFlowId matches '.*finance.*'"
    )
}

Alternatively, extract the routing key first with setExchangeProperty. This approach is useful when the same key is referenced in multiple expressions or the payload path is deeply nested:

setExchangeProperty {
    id = "set-exchange-property"
    propertyName = "filter"
    source = "envelope.payload.value"
}
distribute {
    id = "distribute-message"
    flowId(FLOW_ID_SHIPPING)
    flowId(FLOW_ID_FINANCE)
    targetFilterExpression(
        "exchangePropertiesMap['filter'] == 'address' and targetFlowId matches '.*shipping.*'"
    )
    targetFilterExpression(
        "exchangePropertiesMap['filter'] == 'price' and targetFlowId matches '.*finance.*'"
    )
}

Unmatched Messages

If all targetFilterExpression entries evaluate to false, the message is not forwarded to any target flow. The message is filtered out, and no error is raised.

To route unmatched messages to all target flows:

targetFilterExpression("envelope.payload.value != 'address' and envelope.payload.value != 'price'")

To route all messages to one specific target regardless of content:

targetFilterExpression("targetFlowId == 'my-target'")

The filter Processor

The filter processor cancels a message if its expression evaluates to false, stopping all further processing in the flow. Use it before distribute when messages that do not match any routing criteria should be rejected entirely — before distribution is attempted.

filter {
    id = "filter-message"
    expression = "envelope.payload.value == 'address' or envelope.payload.value == 'price'"
}

What happens on a false result depends on the flow’s exchange pattern:

  • In a RequestResponse flow, the client receives a cancellation response with a reason field:

    { "reason": "Message did not match filter expression: envelope.payload.value == 'address' or envelope.payload.value == 'price'" }
  • In a OneWay flow, the message is silently canceled. No error is raised, and no dead letter is produced.

filter vs targetFilterExpression

filter and targetFilterExpression serve different purposes and can be used together:

Characteristic filter targetFilterExpression

Scope

Cancels the message for the entire flow

Filters per target flow ID inside distribute

Use when

The message should not be processed at all if it does not match

The message should reach some targets but not others

Choose the Right Pattern

Scenario Processor Target flow type Notes

One message to multiple destinations; main flow waits for all

distribute

Headless

The main flow waits until all targets complete. Prefer this option when performance is critical and targets are idempotent.

One message to multiple destinations; decoupled processing

distribute

OneWay + handoff

Each target operates independently; failure in one does not affect others.

Multiple source flows to one shared destination

distribute (from each source)

OneWay + handoff

Fan-in pattern; shared handoff flow consolidates from multiple sources.

Process each collection item; main flow needs all results

split

Headless + aggregate

aggregate collects all fragment responses.

Process each collection item; decoupled from the main flow

split

OneWay + handoff

Each fragment is persisted and processed independently.

Process each collection item, then run post-split logic

split

OneWay + handoff + aggregate (sync point)

Each fragment is persisted and processed independently; aggregate synchronizes processing only.

Process each collection item, then run post-split logic

split

Headless, no aggregate

Place further processing in the main flow after split; this is cleaner than the synchronization-point pattern.

Route messages to different destinations by content

distribute + targetFilterExpression

Headless or OneWay

Reference payload fields directly, or use setExchangeProperty to extract a routing key when you reuse the same key or the payload path is deeply nested.