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+handoffbecause 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:
distributereturns one representative response (not an aggregate), in this order:-
Non-transient failure (any target)
-
Transient failure (any target)
-
Successful response
-
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+handoffbecause 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,splitreturns one representative response.aggregatereturns a composite result only in aHeadlesstarget 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
RequestResponseflow, 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
OneWayflow, 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 |
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 |
|
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 |
|
OneWay + |
Each target operates independently; failure in one does not affect others. |
Multiple source flows to one shared destination |
|
OneWay + |
Fan-in pattern; shared handoff flow consolidates from multiple sources. |
Process each collection item; main flow needs all results |
|
Headless + |
|
Process each collection item; decoupled from the main flow |
|
OneWay + |
Each fragment is persisted and processed independently. |
Process each collection item, then run post-split logic |
|
OneWay + |
Each fragment is persisted and processed independently; |
Process each collection item, then run post-split logic |
|
Headless, no |
Place further processing in the main flow after |
Route messages to different destinations by content |
|
Headless or OneWay |
Reference payload fields directly, or use |