A correlation key ties records together, so that when one of them fails, the rest of its group is turned away too. This page lets you break records and watch what reaches your output and what lands in the dead-letter file (the DLQ). It covers a plain pipeline, a pipeline with an Aggregate, and the case where the engine keeps part of a group instead of throwing it all away.
An orders file has one row per order line. Order A has three lines, B has two and C has one, and the last two lines have no order id. The pipeline reads the file and checks each line's quantity; a quantity that isn't a whole number fails the check.
- type: source
name: orders
config:
correlation_key: order_id
…
- type: transform
name: validate
input: orders
config:
cxl: |
emit q = qty.to_int() # "BAD" fails here
With correlation_key: order_id, a failing line takes the rest of its order with it:
The engine copies the key value into a hidden column as soon as the Source reads the row. Later steps can rename, rewrite or anonymise order_id and the row stays in its original group.
The hidden columns never appear in your output. To see them while debugging, set include_correlation_keys: true on a Sink. They show up as $ck.order_id.
To keep each group's rows together, the engine sorts the Source's rows by the key (empty keys first), keeping file order within a key. If the Source already declares a sort_order that starts with the key, nothing is re-sorted.
You can see this in section 01: with the key on, the output lists the empty-key line first, then B, then C. With the key off, it follows the file.
match: first picks.An Aggregate turns many lines into one total, so “turn the whole group away” needs a rule. The engine looks at your group_by. You never choose this setting yourself:
Each total belongs to exactly one order. If any line of the order fails, the order's total is dead-lettered as correlated, even though the engine did compute it from the clean lines.
Now one order spreads across several totals: order A has lines in both ENG and HR. Throwing away every total an order touched would void whole departments over one bad line. Instead the engine takes out only the failing line and recomputes the totals from the lines that remain, so the clean lines of a failed order still count.
A check after the Aggregate can also fail a total. Under retraction, that total is dead-lettered as a trigger and drops out of the output. The lines that fed it are not dead-lettered, and their orders are untouched. One exception today: those lines are also taken out of any other relaxed Aggregate on the same input, so a sibling total that did not fail loses them too (#1391).
strategy: streaming, because a streaming Aggregate writes each total before the run knows whether any line failed. That combination is rejected when the pipeline is checked (E15Y).A group can be split across branches. If any of its rows fails on any branch, the group is turned away from every branch.
With an inclusive Route, a row that reaches two branches but fails once is dead-lettered once.
Rows keep their group through a Merge, so rows from two Sources with the same key value end up in one group.
A failure only condemns rows from the Source that failed. If a row from src_b fails, the src_a row with the same id still reaches the output.
Every Combine must say which keys its output carries: propagate_ck: driver, all, or { named: [...] }.
If the Combine body fails for a driver row, the build row it matched is written next to it in the DLQ. The build row's own group is not condemned.
error_handling.max_group_buffer (default 100,000) caps how many entries a group may hold. A group over the cap is dead-lettered whole, under one group_size_exceeded row, and the run continues.
The cap does not limit the memory the group uses today (#1259).
Groups flow into a composition and back out unchanged. A composition can't declare its own key; keys belong to Sources.
E153: the key names a column that isn't in the Source's schema.
E15Y: a retraction Aggregate set to strategy: streaming.
Run clinker explain --code <code> for the full text.
| column | what it tells you |
|---|---|
_cxl_dlq_trigger | true when the row's own failure sent it there; false when another row's failure took it along. |
_cxl_dlq_trigger_id | The id of the trigger row behind this one. Group by it to see everything one failure took with it. A trigger points to itself. When an order has several failing lines, its correlated rows point to the first. |
_cxl_dlq_error_category | The real failure on a trigger (for example type_coercion_failure), or correlated on a row condemned by its group. Needs include_reason: true. |
_cxl_dlq_stage | Where it happened: transform:validate for the failing line, correlation_commit for the rows its group condemned. |