distribute

Processor that forwards a message to one or more target flows — all listed targets by default, or a content-filtered subset when targetFilterExpression is used.

Each flowId must point to either a Headless flow or a flow whose source is a handoff-family processor ( handoff, scheduledHandoff, or windowedAggregationHandoff). Flows with outward-facing sources such as restApi cannot be used as target flows.

The two most common approaches are:

  • Headless flow: The target pipeline is inlined into the main flow. The main flow waits for all target flows to complete and shares the same delivery guarantee and exchange pattern. No additional message persistence occurs, making this the higher-performance option. Use this approach when the main flow must wait for all target processing to complete or when minimizing the overhead of additional message persistence is a priority.

  • OneWay + handoff flow: The target flow uses exchangePattern = OneWay with a handoff source. The main flow completes after each message is successfully handed off to the target source. Each target flow independently manages persistence, delivery guarantees, redelivery, and dead letter handling. When distributing to multiple handoff target flows, each target independently persists the message in parallel. This additional persistence is the tradeoff for independent delivery guarantees and is the preferred approach when target systems must be decoupled from each other and from the main flow.

Distribution is not considered complete until all target flows have confirmed message receipt. The distribute processor maintains non-persistent state to track target confirmations. If redelivery is triggered from the source, the processor skips targets that have already confirmed receipt. However, because this state is non-persistent, a flow server restart between the original delivery and a redelivery attempt causes the message to be delivered to all targets again.

This has the following connotations:

  • Message processing is considered successful only when all target flows report success.

  • Any delivery guarantee must be implemented upstream in the main flow unless handoff target flows are used, in which case each target flow provides its own delivery guarantees.

  • Because state is non-persistent, target flows that are not idempotent must implement any required guarantees themselves.

Properties

Name Summary

flowId()

Adds a flow ID to a list of flows that the messages will be distributed to.

targetFilterExpression()

Adds a boolean expression to a list of filters used to select which flows to distribute individual messages to. If no filters are defined, all messages will be distributed to all target flows. If at least one filter is defined, each flow ID will be evaluated against the filter expression. If at least one expression returns true, the message will be distributed to that flow ID. Otherwise, it will not be distributed to that flow ID.

distributionExpiryMillis

The amount of time after distribution is initiated that the distributor maintains state of delivery.

distributionExpiryCheckMillis

The frequency of which the distributor checks for and evicts its internal state for abandoned/failed distributions.

distributionMaxMessages

The maximum number of active messages that the distribution processor can handle simultaneously. If this threshold is reached, the distributor will replace an existing entry by using an undefined strategy.

retainPayloadOnFailure

Whether the incoming payload is available for error processing on failure. Defaults to false.

name

Optional, descriptive name for the processor.

id

Required identifier of the processor, unique across all processors within the flow. Must be between 3 and 30 characters long; contain only lower and uppercase alphabetical characters (a-z and A-Z), numbers, dashes ("-"), and underscores ("_"); and start with an alphabetical character. In other words, it adheres to the regex pattern [a-zA-Z][a-zA-Z0-9_-]{2,29}.

exchangeProperties

Optional set of custom properties in a simple jdk-format, that are added to the message exchange properties before processing the incoming payload. Any existing properties with the same name will be replaced by properties defined here.

Sub-builders

Name Summary

inboundTransformationStrategy

Strategy that customizes the conversion of an incoming payload by a processor (e.g., string to object). Should be used when the processor’s default conversion logic cannot be used.

messageLoggingStrategy

Strategy for describing how a processor’s message is logged on the server.

payloadArchivingStrategy

Strategy for archiving payloads.

Details

Filtering

The targetFilterExpression() function evaluates expressions in a context that expands the standard expression context with the targetFlowId property. targetFlowId contains the current flow ID being evaluated.

You can add multiple expressions to the same processor to create advanced filtering strategies. For example, the following expressions:

targetFilterExpression("exchangePropertiesMap['targetSystem'] == null")
targetFilterExpression("exchangePropertiesMap['targetSystem'] == 'systemAbc' and targetFlowId.contains('-abc-')")
targetFilterExpression("exchangePropertiesMap['targetSystem'] == 'systemXyz' and targetFlowId.contains('-xyz-')")

The processor distributes messages without a targetSystem exchange property (that is, targetSystem == null) to all targets. If targetSystem is set, the processor distributes the message only to that target.

The processor treats messages for filtered targets the same way as messages filtered by the filter processor. If a filter processor on the target flow or targetFilterExpression() filters out all destinations, the main flow returns a filtered message.

Results

The distribute processor produces a result only after it receives a response from every target flow. The processor does not aggregate the responses. Instead, it returns the response from one target flow. The processor selects the response in the following order of precedence:

  1. Message failure from a non-transient error

  2. Message failure from a transient error

  3. Successful message

  4. Filtered message

Byte Streaming

In a streaming flow, the distribute processor acts as a Stream Agnostic Processor.

In streaming mode, the distribute processor can send the stream to exactly one downstream flow or to none. You can configure multiple targetFilterExpression entries, but the processor can select only one destination for a given stream.

For general design guidance and a pattern decision table for the distribute processor, see Distribution and Choosing the Right Pattern.