Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

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_exceeded trigger, the first of them; the rest are correlated rows that carry its id as their _cxl_dlq_trigger_id.
  • A group whose rows all failed writes only those failures, with no group_size_exceeded row.

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_buffer overflow 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