Correlation Keys
A correlation key declares a set of records from a single source as an atomic group: if any record in the group fails validation or processing, the whole group is sent to the DLQ. This is the right shape for transactional data where partial processing is worse than total rejection – the canonical example is an order with multiple line items where one bad line should reject the entire order.
This page describes how to declare a correlation key and how it behaves through each node that can fan out, fan in, group, or join records.
Interactive companion: the correlation keys explainer lets you make lines fail and see what reaches the output and the DLQ, with and without an Aggregate.
Declaration
Correlation keys are declared per source. Each source’s config: block carries an optional correlation_key: field naming the column (or list of columns) whose value identifies a record’s correlation group within that source.
nodes:
- type: source
name: orders
config:
name: orders
type: csv
path: ./data/orders.csv
correlation_key: order_id
schema:
- { name: order_id, type: string }
- { name: amount, type: int }
- type: source
name: customers
config:
name: customers
type: csv
path: ./data/customers.csv
correlation_key: [customer_id, region] # multi-column key
schema:
- { name: customer_id, type: string }
- { name: region, type: string }
- { name: name, type: string }
- type: source
name: sensor_readings
config:
name: sensor_readings
type: csv
# No correlation_key: record-level errors land in the DLQ as
# standalone entries with no group atomicity.
schema:
- { name: ts, type: date_time }
- { name: value, type: float }
A record’s correlation group is identified by the tuple of values for that source’s listed fields. Records sharing the same tuple within the same source belong to the same group. There is no pipeline-level correlation key — declare it on each contributing source.
The group identity is captured at ingest, so rewriting the key column in a later Transform does not change a record’s group — anonymizing or transforming order_id downstream still keeps the original grouping intact.
A source whose declared correlation_key: field names a column not present in its own schema: block is rejected at compile time with diagnostic E153. The fix is to add the field to the schema or remove it from correlation_key:.
Row order
Declaring correlation_key on a Source changes the order its rows reach the rest of the pipeline. When the Source’s declared sort_order already begins with the correlation-key fields, in either direction, that declared order is kept as it is. Otherwise the planner inserts a sort after the Source, ascending on the correlation key and then on the remaining declared sort_order fields, so that each group’s rows are adjacent. Rows with equal keys keep their order in the file.
Every consumer that depends on row order sees this order: Sink output order, an Aggregate’s first-arriving value, Cull and Reshape ties, and which build record a Combine’s match: first picks. Adding a correlation key to an existing pipeline can therefore change those results even when no record fails.
DLQ semantics
When a record fails inside a correlation group:
- The failing record produces a trigger DLQ entry. Its category reflects the actual failure (e.g.
type_error,validation_failed). - Every other record from the same source in that group produces a collateral DLQ entry, carrying the category
correlated. - Records belonging to other (clean) groups proceed normally.
A record with a null value for the correlation-key field is treated as its own group: it has no peers, so DLQ atomicity does not span multiple records.
The dlq_count counter sums triggers and collaterals. It counts one trigger per failure, so a row that fails twice in a group counts twice, as it does without a key.
Group buffering
The engine buffers records per correlation group until either the group completes or a failure triggers a flush. The max_group_buffer: field on the pipeline-level error_handling: block caps per-group buffering across every source’s groups:
error_handling:
max_group_buffer: 100000 # Default: 100,000
A group that goes over the cap is dead-lettered whole when the run commits it. It is not a hard error: the run continues.
- Each row of the group that failed on its own is written as its own trigger, with its own category and
_cxl_dlq_id. - The group’s other rows are written under one
group_size_exceededtrigger, the first of them; the rest arecorrelatedrows that carry its id as their_cxl_dlq_trigger_id. - A group whose rows all failed writes only those failures, with no
group_size_exceededrow.
Each failure is written once, with its Combine build row if it has one (see Combine interaction), each condemned row once, and every row counts toward dlq_count and the DLQ rate limits. The group_size_exceeded row’s _cxl_dlq_timestamp is when the group went over the cap, and its error detail states the cap and how many entries the group held.
The cap counts the entries a group holds, not its distinct rows: a row counts once for each Sink it reaches, and each failure counts once. A row that an inclusive Route sends to two Sinks counts twice, and so does a row that fails on one branch and reaches a Sink on another. A failing Combine match counts once, although it holds both the driver row and the matched build row.
Going over the cap does not stop a group from buffering. The group keeps buffering its rows until the run commits it, so today the cap decides how a large group is dead-lettered but does not bound the memory it uses.
Per-operator interactions
Route interaction (fan-out)
A correlation group can span multiple route branches. Group atomicity is preserved across branches: if any record in the group fails (in any branch’s transform, or in the route predicate itself), the entire group is rejected from every branch.
For an inclusive route where one record reaches both branches, a single failure DLQ’s that source row exactly once — not once per branch. A record that fails on both branches has two failures and is written twice, once as each branch’s trigger with its own stage, as it is without a key.
Merge interaction (fan-in)
Merge concatenates upstream branches that share a schema. Records keep their correlation identity through the merge, so rows from different sources that share the same key value become one correlation group downstream: a failure on any one of them DLQ’s the whole group across both sources.
Per-source rollback narrowing
When two sources contribute records to the same correlation group, a failure originating from one source does not collaterally DLQ records from the other source. The collateral fan-out is scoped to the failing source’s records only.
For example, with [src_a, src_b] → merge → transform → out where both declare correlation_key: id, an error that fires on a src_b row produces a trigger for that row while the src_a row sharing the same id is spared and reaches the output. Single-source pipelines behave exactly as a pipeline-wide collateral DLQ would, since every co-grouped record shares the one source.
Two cases stay group-wide rather than narrowing per source:
max_group_bufferoverflow DLQ’s every record in the overflowing group — no single source is to blame for the overflow.- Combine output failures DLQ the synthesized output row, which has no single-source attribution. The exception concerns the output row only: the matched build record’s dead letter follows the failing driver’s group, as described under Combine interaction, and never widens the narrowing to the build record’s source.
Aggregate interaction
When an aggregate’s group_by covers every correlation-key field, the aggregate stays on the strict path: each emitted row inherits the correlation identity of its inputs, and any DLQ trigger in the group rolls back every record in the group, including the aggregate output row.
- type: aggregate
name: order_totals
input: orders # correlation_key: order_id
config:
group_by: [order_id] # covers the key
cxl: |
emit total = sum(amount)
When an aggregate’s group_by omits a correlation-key field, the engine automatically retracts only the failing records and recomputes the affected groups, so the surviving contributions still produce a correct aggregate row. You do not configure this — the engine picks the path from the group_by content. (One restriction: this mode cannot be combined with strategy: streaming, which is rejected at compile time.)
- type: aggregate
name: dept_totals
input: orders # correlation_key: order_id
config:
group_by: [department] # omits the key — surviving rows recomputed
cxl: |
emit total = sum(amount)
Combine interaction
Every combine declares propagate_ck: to select which correlation-key fields its output rows carry:
propagate_ck: driver— output inherits only the driver input’s correlation identity. The common case; today’s strict-correlation pipelines stay on this setting.propagate_ck: all— output carries the union of correlation-key fields across every input. Use when the build side carries keys that downstream operators need to read.propagate_ck: { named: [<field>, ...] }— output carries exactly the named subset. Use to project a multi-field key down after a join.
- type: combine
name: enriched
input:
o: orders # driver (correlation_key: employee_id)
d: departments # build side
config:
where: "o.employee_id == d.employee_id"
match: first
on_miss: skip
cxl: |
emit employee_id = o.employee_id
emit amount = o.amount
emit dept = d.dept
propagate_ck: driver
How match mode fills the propagated key:
match: first— the single matched build’s key fills the slot.match: all— one output row per matched build, each carrying its own build’s key.match: collect— one row per driver; the first matched build’s key fills the (single-valued) slot, while every matched build’s full payload still rides inside the array column.
Driver wins on a name collision: if both the driver and a build input declare the same key field, the output keeps the driver’s value.
propagate_ck is a required field — every combine must spell out which mode it uses.
A failing match’s build-side dead letter follows the driver’s group. When the combine body fails for a driver row, that driver row is the trigger of the driver’s correlation group. The matched build record’s dead letter is held with the same group as a collateral (_cxl_dlq_trigger: false, category combine_output_row), written right after its driver’s row and carrying that driver’s _cxl_dlq_trigger_id, and it is written or rolled back exactly when the driver’s group is. It never condemns the build record’s own correlation group: another driver that matched the same build record keeps its output unless its own group failed. The build record is written once per failure: when several drivers fail against one build record, in one group or in several, each failing driver’s row is followed by its own copy of the build row, carrying that driver’s _cxl_dlq_trigger_id. A driver that fails against several build rows is written once per failure, each copy followed by the build row of that failure.
Composition interaction
A composition’s body operates on records flowing in from the parent pipeline; correlation identity flows into the composition inputs and back out the named ports unchanged. Compositions cannot declare their own correlation key — a key is a property of a source, not of the composition body that consumes a source’s records.
Debugging
Correlation grouping is tracked on internal columns you never write in YAML or CXL, and they are hidden from writer output by default. To surface them for debugging, set include_correlation_keys: true on a Sink node:
- type: sink
name: debug
input: any_node
config:
type: csv
path: "./debug.csv"
include_correlation_keys: true
The output then contains extra columns named $ck.<field> (literal prefix in the CSV header) for each declared correlation-key field.
To investigate DLQ collaterals: every collateral entry’s category is correlated, and the trigger entry in the same group carries the actual failure category and message.
See also
- Error Handling & DLQ – general DLQ configuration, fail-fast vs continue, type-error thresholds.
- Aggregate Nodes – group-by semantics and the strategy hint.
- Combine Nodes – driver selection and match modes.
- Sink Nodes –
include_correlation_keysand other field-control flags.