Error Handling & DLQ
Clinker provides structured error handling with a dead-letter queue (DLQ) for records that fail processing. The error_handling: block at the top level of the pipeline YAML controls the behavior.
Configuration
error_handling:
strategy: continue
dlq:
path: "./output/errors.csv"
include_reason: true
include_source_row: true
Strategies
error_handling.strategy is pipeline-wide – it is set once at the top level, not per node. It controls what happens when a record fails:
| Strategy | Behavior | Exit code |
|---|---|---|
fail_fast | Default. Abort the run on the first record failure. | Non-zero, by the class of the aborting error (3 for an evaluation failure, 4 for an I/O failure – see Exit Codes) |
continue | Route the failing record to the DLQ and keep processing. | 2 if any record was dead-lettered, 0 otherwise |
There are exactly two, because the engine makes exactly one decision at each record failure: propagate it and stop, or dead-letter it and carry on.
fail_fast
The safest strategy. Any record-level error (type coercion failure, validation error, missing required field) halts the pipeline immediately, with a non-zero exit and no DLQ file. Use this when data quality is critical and you prefer to fix issues before reprocessing.
Some failures abort the run under either strategy, because they are not record-scoped: an unwritable output path, a config or CXL compile error, and the DLQ-rate ceiling (dlq.max_rate, E315/E316) all end the run regardless of the strategy.
CSV, JSON and XML resource failures are also fatal under either strategy. Memory or disk
admission refusal, allocation failure, descriptor exhaustion and temporary-storage
failure are not bad-record errors, so continue cannot turn them into successful
output. A typed resource diagnostic preserves the kind of failure rather than
reporting every case as a memory shortage. A failed destination can already have
accepted a prefix; see output preparation.
Malformed JSON/XML input encoding is a data failure, including when discovered
during schema discovery or envelope pre-scan. Under fail_fast, the CLI returns
exit 4 and machine code source.data.invalid. It does not report a compilation
error merely because no record has reached the pipeline. A late error can leave
an already delivered prefix; an envelope pre-scan may discover it before any
body records. Failed runs do not publish their staged normal output files.
Explicit cancellation ends an interrupted run with exit 130; it does not add a
Sink error. If a real I/O or resource failure occurs alongside a shutdown request,
the real failure retains its classification. Record and byte counters describe
established progress, not rows merely attempted or prepared.
An executor invariant failure also aborts under either strategy with exit code
1. In particular, if a planned materialized input is unavailable when its
consumer runs, Clinker stops instead of treating that input as a legitimate
zero-row result. The message names the consuming node and planned producer
(including the producer port when applicable) and says the input was not
treated as empty. A source or stage that really emits zero rows remains valid;
it carries an explicit empty buffer and completes normally. Report any missing-
input internal error as an engine defect rather than routing it to the DLQ.
continue
The production workhorse. Bad records are written to the DLQ file with diagnostic metadata, and the pipeline continues processing remaining records. After the run completes, inspect the DLQ to understand and correct failures.
A pipeline that completes with DLQ entries exits with code 2 – this signals “pipeline completed successfully but some records were rejected.” It is not a crash or internal error. A continue run that dead-letters nothing exits 0, exactly like a clean fail_fast run.
Migrating from
best_effort. The removedbest_effortspelling was a third name for thecontinuebehavior: it wrote the same DLQ entries and produced the same exit code, because the runtime never distinguished the two. Replace it withstrategy: continue. A pipeline still carryingbest_effortis rejected at config-validation time with a message naming the replacement.
Declared source-type failures are deliberately not lossy: under continue,
the complete original record is written to the
configured DLQ and no null, raw, or partially converted replacement enters the
pipeline. This strategy therefore requires an error_handling.dlq block when
such a failure occurs. fail_fast stops on the first failure without emitting a
replacement. This includes fields renamed by a source schema: rejection retains
the original decoded record and its values, even when conversion failed after
other fields had already been examined.
An evaluation error is never false
A condition that fails to evaluate, such as one that divides by zero, has not said whether it holds. The engine never reads that failure as “false”, on any node: the record the condition was about is dead-lettered, and no decision that depends on the condition being false is taken for it.
- A Transform
filterthat fails dead-letters the record; it is neither kept nor filtered out. - A Route branch condition that fails dead-letters the record, which takes no
branch and not the
default. Inexclusivemode only the conditions up to the first true one are evaluated, so a later condition cannot fail. - A Combine
where:that fails for a candidate build row dead-letters that pair. The driver is not unmatched, soon_missdoes not fire; undermatch: firsta failure on the deciding candidate is the driver’s only result, and undermatch: collectthe driver writes no row. See Combine.
A pipeline that wants a failing condition treated as false says so in CXL, for
example by guarding the division or coalescing the result with ?? false.
DLQ configuration
The DLQ is always written as CSV, regardless of the pipeline’s input/output formats.
dlq:
path: "./output/errors.csv"
include_reason: true
include_source_row: true
| Field | Required | Default | Description |
|---|---|---|---|
path | No | – | The pipeline-wide DLQ file. It receives the dead letters of every Source without its own per_source path. A dead letter with neither this path nor a per_source path for its Source is counted in the run’s dead-letter totals, sets exit code 2 and counts toward max_rate, but is written nowhere. The dlq: block itself is required to continue past a declared source-type failure. |
include_reason | No | true | Include _cxl_dlq_error_category and _cxl_dlq_error_detail columns. |
include_source_row | No | true | Include the failing record’s columns after the _cxl_dlq_* columns. Which record columns each DLQ file carries is fixed when the pipeline compiles; see How the DLQ columns are chosen. With false, only the _cxl_dlq_* columns are written. |
max_rate | No | none | Stop the run (E315, exit code 3) once the dead-lettered rows reach this fraction of the source rows read so far, both counted across the whole run. Must be greater than 0.0 and at most 1.0 (E318). Without it, the run is never stopped for its dead-letter rate. See Bounding how much can dead-letter. |
min_records | No | 100 | How many source rows must have been read before max_rate is checked, so the first failures of a run cannot trip it on a tiny denominator. Also the default for each per_source min_records. |
per_source | No | – | Settings for individual Sources, keyed by Source node name: a separate DLQ file, and a rate ceiling of their own. See Per-source DLQ settings. |
Per-source DLQ settings
per_source gives a Source its own DLQ file, its own rate ceiling, or both.
Each key is the name of a Source node:
error_handling:
strategy: continue
dlq:
path: ./output/errors.csv
max_rate: 0.05
per_source:
vendor_feed:
path: ./output/vendor_feed_errors.csv
max_rate: 0.20
min_records: 500
orders:
max_rate: 0.01
| Field | Default | Description |
|---|---|---|
path | – | A separate DLQ file for this Source’s dead letters. They are written only there and do not appear in the pipeline-wide file. Without it, the Source’s dead letters go to the pipeline-wide path. |
max_rate | none | Stop the run (E316, exit code 3) once this Source’s dead-lettered rows reach this fraction of the rows read from this Source so far. Must be greater than 0.0 and at most 1.0 (E318). |
min_records | the pipeline-wide min_records, else 100 | How many rows must have been read from this Source before its max_rate is checked. |
A Source’s own max_rate is checked first, so a breach names that Source.
The pipeline-wide max_rate, when set, still applies to the run as a whole.
In the example, vendor_feed may dead-letter up to 20% of its own rows, but
the run still stops when all dead letters together reach 5% of all rows read.
A key that does not name a declared Source is rejected at compile time (E317). Two DLQ paths that name one file are rejected too (E318), including paths that differ only in case on a case-insensitive filesystem, or in being written relatively and absolutely. A DLQ path that names the same file as a Sink’s path is rejected with E322.
Which record columns each file carries follows from the Sources routed to it; see How the DLQ columns are chosen.
How DLQ output is written
Dead-letter rows are written while the run executes, not collected until it
ends. Each row is formatted under its file’s header, which is fixed when the
pipeline compiles (see
How the DLQ columns are chosen), and written into
a staged copy of that file in the run’s publication attempt: in quarantine
next to the destination by default, or under local_spool_dir with
mode = "local_then_publish" (see
Output publication).
Each open DLQ file writes through one fixed 64 KiB buffer, so the memory the
DLQ files use does not grow with the number of failures.
-
A DLQ file is created when its first row arrives. A file no row reaches is not created, and no empty file is published.
-
DLQ files are published only if the run succeeds, by the same publication step as the pipeline’s other outputs. A failed or interrupted run publishes no DLQ file.
-
The failures that are counted but have no destination (see
pathabove) are never formatted or written. -
Three kinds of dead letter are held in memory until the stage that found them finishes, and written then:
- a
join_valuescollision at a Sink that writes on its own thread; - an Aggregate
add_recordfailure found while the Aggregate reads its input on its own thread; - a Combine output-row failure found while the Combine streams its driver on its own thread, or inside a grace-hash, sort-merge or IEJoin join.
A Sink that writes on its own thread stops the run with an internal error once 65,536 collisions are waiting this way.
- a
-
Under a correlation key or
dlq_granularity: document, records are held until their group or document is decided. That is those features’ own state, described in their sections, not DLQ output; the rows they dead-letter are then written like any other. Underdlq_granularity: documentthe engine also keeps, for each rejected document, a compressed record of which rows it has already written, so a row held by several Sinks is written once. That record is charged to the memory budget. While a Sink writes a rejected document’s rows, the record grows by:- up to 16 bytes per row when the Sink receives the rows in the order their Source read them;
- up to 96 bytes per row when it receives them in any other order, for example after a Sort on a data column;
- up to 672 bytes per row for a document whose rows span its Source’s 4,294,967,296th row, in any order. A Source counts its rows across every file it reads.
The first row of each document a Sink starts costs up to 920 bytes. The growth is released when the Sink’s pass ends, or every 65,536 rows of a document. Between releases one document’s record therefore grows by at most about 1 MiB when its rows arrive in the order they were read, about 6 MiB in any other order, and about 42 MiB when its rows span the 4,294,967,296th row. For example, 65,535 rows that alternate between a document’s first 65,536 rows and its next 65,536 grow the record to about 5.8 MB, and the release at the next row brings it back to 880 bytes. Once released, the record costs at most 4 bytes per written row, plus 96 bytes for every 65,536 rows of the document, written or not, plus under 1 KiB. That is about 2.3 KiB for a million-row document whose rows are contiguous. When only one row in each 65,536 is written, it is 88 bytes per written row plus 672 bytes, within the cap of about 100 bytes per written row. The record cannot spill: if one more row would not fit once every held row (below) has moved to disk, the run fails with E310. A failed document’s failing records are formatted as DLQ rows when they fail and held until the document is rejected. They are held in memory, charged to the memory budget, and move to one file in the spill directory only when the budget needs the memory, counting toward
storage.spill.disk_cap_bytes(E320). With memory to spare nothing is written to disk. If even one more held row would not fit with every held row on disk, the run fails with E310.
Disk bounds how much DLQ output a run can produce: the free space at the
staging location, and the publication attempt’s byte ceiling
(storage.publication.max_attempt_bytes, and no more than
retained_byte_limit), which every staged file of the run counts toward,
DLQ files included. An attempt larger than that ceiling is refused at
publication and nothing is published.
When the staging location fills, the run fails with an I/O error (exit code 4) and publishes nothing. After the destination’s own error text, the message names the DLQ file, the number of rows dead-lettered so far, and the stage and category with the most rows, and suggests a breaker:
<destination error>: dead-letter output ./output/errors.csv could not be written after 1048576 dead-lettered rows (most from stage transform:validate_orders, category validation_failure)
help: stop the run before dead letters fill the destination, for example:
error_handling:
type_error_threshold: 0.05
dlq:
max_rate: 0.05
These breakers bound how many rows can dead-letter; disk at the staging destination bounds their volume.
Bounding how much can dead-letter
Disk bounds the volume of dead-letter output; the breakers bound how many rows can dead-letter in the first place. Set one wherever a wrong schema could make most rows fail, for example when a feed can change its columns or types without notice:
error_handling:
strategy: continue
type_error_threshold: 0.05
dlq:
path: ./output/errors.csv
max_rate: 0.05
dlq.max_rate(E315) anddlq.per_source.<name>.max_rate(E316) stop the run when the fraction of dead-lettered rows crosses the ceiling, oncemin_recordsrows have been read. The numerator counts dead-letter rows, collateral rows included, whether or not they have a DLQ file to go to. A source row counts once for each failure it took part in, so a row that failed twice counts twice, with or without a correlation key. Underdlq_granularity: document, a row that failed twice at the Source, in a Transform or in a Route is written once and counts once (see Document-level DLQ).type_error_threshold(E368) stops the run when the fraction of declared source-type failures crosses the threshold. It catches a schema mismatch at the Source, before the failing rows reach later stages.
A rate ceiling is checked each time a dead letter is counted, so the row that crosses it is itself counted and written before the run stops. A stopped run exits with code 3 and publishes nothing. No breaker is set by default.
DLQ columns
Every DLQ record includes these metadata columns:
| Column | Description |
|---|---|
_cxl_dlq_id | UUID v7 (time-ordered unique identifier), unique to the row. It is taken together with _cxl_dlq_timestamp, so ids order the same way as timestamps. |
_cxl_dlq_trigger_id | The _cxl_dlq_id of the trigger row whose failure produced this row. That trigger row is always written in the same run. A trigger row points to itself, so _cxl_dlq_trigger is true exactly when this value equals _cxl_dlq_id; every row one failure produced carries the same value. See Pairing the rows one failure produced. |
_cxl_dlq_timestamp | RFC 3339 timestamp of when the failure was observed, not of when the row was written. A collateral row (correlated, document_rejected) carries the time its correlation group or document was condemned, except that a rejected document’s other failing records carry the time each one failed. A group_size_exceeded row carries the time its group went over max_group_buffer. |
_cxl_dlq_source_file | Input filename carried by that failing record’s $source.file provenance (or <merged> when no source-file provenance exists) |
_cxl_dlq_source_name | Name of the Source the failing record came from (or <merged> when the record carries no Source identity) |
_cxl_dlq_source_row | 1-based row number in the source file |
_cxl_dlq_triggering_field | The field whose evaluation failed, when the failure names one; empty for collateral rejections |
_cxl_dlq_triggering_value | The value the failure reported, when it carries one (for example the text that failed to convert) |
_cxl_dlq_stage | Name of the transform or aggregate node where the error occurred |
_cxl_dlq_route | Route branch name (if the error occurred after routing) |
_cxl_dlq_trigger | true when the row’s own failure dead-lettered it; false when another row’s failure took it along (a correlated or document_rejected row, or a Combine build row) |
_cxl_dlq_source_record | One of the record columns rather than a metadata column: present in any file a Source rejection can reach under strategy: continue, and filled only for a record-grained E345 rejection. Contains the fixed-width line text or a JSON array of decoded CSV cells, preserving the physical row without assigning it a declared record shape. |
Timestamps need not increase down a file. A failure can be held before its
row is written, for example in a correlation group that commits later, or by
a stage listed under How DLQ output is written,
so its row can follow rows observed after it. Sort on _cxl_dlq_timestamp or
_cxl_dlq_id to read failures in the order they were observed.
When include_reason: true is set, two additional columns appear:
| Column | Description |
|---|---|
_cxl_dlq_error_category | Machine-readable error classification |
_cxl_dlq_error_detail | Human-readable error description |
Pairing the rows one failure produced
Two rules decide which rows a DLQ file holds:
- One row per failure. Every failure writes its own trigger row, with its
own category, detail, triggering field and value, stage and route. A row
that fails twice, on two Route branches or against two Combine build rows,
is written twice. A Combine failure also writes the build row that
contributed to it, right after its trigger. Under
dlq_granularity: documentthis rule has an exception: a rejected document’s later failing record is written as adocument_rejectedrow paired with the document’s first failure, and a record that failed twice at the Source, in a Transform or in a Route is written once (#1316; see Document-level DLQ). - Each condemned row once. A correlation group or document that fails
adds each of its other rows once, as a collateral of its first failure.
Under
dlq_granularity: document, a record that an Aggregate, a Combine or a Reshape failure dead-lettered is written again, as adocument_rejectedrow, only when a Sink on another branch also received it and its document is rejected. If its document is not rejected, that Sink can publish it (see Not covered).
A correlation key never removes, merges or relabels a failure row: the same failures are written with and without a key, and the key only adds the rows a failing group condemns and decides when rows are written.
One failure can dead-letter several rows:
- a Combine body that fails writes the driver row and its matched build row,
each under its own Source’s
_cxl_dlq_source_nameand_cxl_dlq_source_row, whichever join strategy ran. The build row is written once per failure, so a build record that two failing drivers matched is written twice, once after each driver’s row, and a driver that fails against three build rows is written three times, each copy followed by its own build row; - a failing row in a correlation group takes the rest of its group with it as
correlatedrows; - a group larger than
max_group_bufferwrites agroup_size_exceededrow and the rest of its group ascorrelatedrows (the group’s own failures keep their own trigger rows, see below); - under
dlq_granularity: document, a failing record rejects the rest of its document asdocument_rejectedrows.
Every row carries _cxl_dlq_trigger_id, the _cxl_dlq_id of the trigger row
whose failure produced it, and that trigger row is always in the run’s
output. A trigger row points to itself, so _cxl_dlq_trigger
is true exactly when _cxl_dlq_trigger_id equals _cxl_dlq_id. Group by
_cxl_dlq_trigger_id to see everything one failure took with it.
In this excerpt (other columns omitted), the first two rows are a Combine
driver row and its build row, and the last two are a correlation trigger and
one of its collaterals:
_cxl_dlq_id,_cxl_dlq_trigger_id,_cxl_dlq_source_name,_cxl_dlq_error_category,_cxl_dlq_trigger
01928f3a-6c10-7b21-8a4e-3f1c2d9e0a01,01928f3a-6c10-7b21-8a4e-3f1c2d9e0a01,orders,combine_output_row,true
01928f3a-6c10-7b22-9f07-51e6a8b4c302,01928f3a-6c10-7b21-8a4e-3f1c2d9e0a01,rates,combine_output_row,false
01928f3a-6c14-7c03-b2d8-0a9e7f615203,01928f3a-6c14-7c03-b2d8-0a9e7f615203,employees,type_coercion_failure,true
01928f3a-6c19-7d40-8c11-6e2b90d3f404,01928f3a-6c14-7c03-b2d8-0a9e7f615203,employees,correlated,false
A trigger whose failure wrote nothing else carries its own id and shares it
with no other row. Every row has its own _cxl_dlq_id, so the value only
repeats across the rows of one failure.
When a correlation group holds several failing rows, each failing row is a
trigger and keeps its own id as its trigger id. The group’s correlated rows
carry the trigger id of the group’s first failing row, the one whose error
their _cxl_dlq_error_detail quotes. A rejected document has one trigger, its
first failing record; its other records, including any that failed after it,
carry that trigger’s id.
A group larger than max_group_buffer can hold failing rows too. Each of
them is still written as its own trigger, with its own category and id, and
before the rest of the group. The group’s other rows follow under one
group_size_exceeded trigger, the first of them, and the remaining ones are
correlated rows carrying its id. Each failure is written once, as its own
trigger; a row that failed on one Route branch and reached a Sink on another
is written as its own failure, and never again as correlated or as the
overflow trigger. A group whose rows all failed writes no
group_size_exceeded row, because the overflow took nothing with it that
had not already failed.
The rows of one failure can land in different DLQ files: with
per_source paths, a Combine build row goes to its own Source’s file while
its driver row goes to the driver’s.
How the DLQ columns are chosen
Each DLQ file’s header is fixed when the pipeline compiles, before any record is read. It does not depend on which records failed, or on which stages they failed in: every time a pipeline writes a given DLQ file, that file has the same columns in the same order.
A header starts with the _cxl_dlq_* metadata columns, always in this order:
_cxl_dlq_id, _cxl_dlq_trigger_id, _cxl_dlq_timestamp, _cxl_dlq_source_file,
_cxl_dlq_source_name, _cxl_dlq_source_row, _cxl_dlq_triggering_field,
_cxl_dlq_triggering_value, then _cxl_dlq_error_category and
_cxl_dlq_error_detail when include_reason is on, then _cxl_dlq_stage,
_cxl_dlq_route and _cxl_dlq_trigger.
With include_source_row on, the record columns follow. They come from every
record shape that can reach that file:
- the declared columns of each Source, and
_cxl_dlq_source_record, when Source rejections are dead-lettered (strategy: continue); - the shape of the records entering each Transform, Route, Reshape, Aggregate, Combine and Sink that can dead-letter a record under the pipeline’s strategy (stages inside a composition count at the composition’s place in the pipeline);
- the output columns of an Aggregate without
group_byunderstrategy: continue, which go to the pipeline-wide file.
A shape is added to the file of each Source whose records can carry it: that
Source’s per_source.<name>.path when it has one, otherwise the pipeline-wide
path. A shape that carries no Source identity, such as the output of a
Combine using match: first or match: all, is added to the pipeline-wide
file.
The shapes then combine as follows:
- One shape keeps its natural column order.
- Several shapes give their first-seen union in plan order: each column
appears once, at the position where it was first seen. Plan order is the
node order
--explainprints, which need not match the order the nodes are written in the YAML. - A column a record does not carry is an empty cell in that record’s row.
- Engine-only sidecar columns (
$widenedand the$source.*stamps) are never written. Correlation-key columns ($ck.*) are.
For example, this pipeline has two Sources with different columns, one pipeline-wide DLQ file, and a Transform on each Source that can fail:
error_handling:
strategy: continue
dlq:
path: rejects.csv
nodes:
- type: source
name: orders
config:
schema:
- { name: order_id, type: int }
- { name: amount, type: int }
# ...
- type: source
name: refunds
config:
schema:
- { name: refund_id, type: int }
- { name: order_id, type: int }
- { name: amount, type: int }
# ...
# one Transform on each Source, then a Sink on each Transform
In this plan refunds comes before orders, so the record columns of
rejects.csv are:
refund_id,order_id,amount,_cxl_dlq_source_record
An orders row writes an empty refund_id cell. order_id and amount
appear once, although both Sources declare them. _cxl_dlq_source_record is
in the header because a Source rejection can reach the file under continue.
It is empty for these rows.
To see every DLQ file’s columns before a run, use
clinker run pipeline.yaml --explain: its === Dead-Letter Output ===
section lists each file, the Sources routed to it and its full header (see
Explain Plans).
Source-file provenance (_cxl_dlq_source_file) is read from each record, so
one file’s path is never reused for a later row.
Error categories
The _cxl_dlq_error_category column contains one of these values:
| Category | Description |
|---|---|
missing_required_field | A required field is absent from the record |
type_coercion_failure | A value could not be converted to the expected type |
required_field_conversion_failure | A required field exists but its value cannot be converted |
nan_in_output_field | A computation produced NaN |
aggregate_type_error | An aggregate function received an incompatible type |
validation_failure | A declarative validation check failed |
aggregate_finalize | An aggregate function failed during finalization: an integer sum outside the 64-bit range, a decimal total or quotient outside the decimal range, a weighted_avg whose weights total zero or whose row product is out of range, or a group holding both decimals and floats. The reason names the Aggregate and the emit, and gives the fix. Under continue the failed group goes to the dead-letter output. An aggregate in an Envelope footer: stops the run under every strategy. |
correlated | A non-failing record was DLQ’d as collateral because another record in its correlation group failed |
group_size_exceeded | A correlation-key group exceeded the configured max_group_buffer limit |
document_rejected | A non-failing record was DLQ’d as collateral because another record in its document failed under a source’s dlq_granularity: document policy |
late_record | A record arrived at a time-windowed aggregate after its event-time window had already closed |
expansion_limit_exceeded | Per-input fan-out exceeded its authored ceiling. Transform max_expansion rejects before body rows emit; Source max_output_rows_per_input emits exactly its ceiling, then DLQs the original input on the first attempted row above it. Neither is silent truncation. |
combine_output_row | A Combine output-stage eval failed for one driver row (probe-key or on_miss: null_fields body) or for one matched pair (residual or matched body). A failing residual is neither a match nor a miss: on_miss never fires for its driver, under match: all the driver’s other matches are still evaluated and emitted, under match: first a failure on the deciding candidate is the driver’s only result and a failure after it is never written, and under match: collect the driver writes no row; the entry carries the contributing-build lineage and rewinds both the driver and matched build source’s rollback cursor. The driver row and the matched build row each report their own Source in _cxl_dlq_source_name and their own row in _cxl_dlq_source_row, whichever join strategy ran. Routed to the DLQ under continue across every Combine join mode; fail_fast propagates the eval error |
structural_validation | A structural source rule failed: an envelope trailer’s declared count did not match its streamed body, a multi-record body appeared after its closing trailer, or a record type discriminator was unknown. Under dlq_granularity: document, the root cause has trigger: true and every already-streamed record of that file is document_rejected collateral. Under record-grained continue, E345 instead emits only the unknown row with _cxl_dlq_source_record. |
Advanced options
Type error threshold
Abort the pipeline if the fraction of declared source-type failures exceeds a threshold:
type_error_threshold: 0.05 # Abort if >5% of records fail
The cumulative ratio is:
declared source-type failures / decoded source rows observed
The rejected row appears once in both numerator and denominator. The same
typed-error event is used for strategy routing, DLQ accounting, and this
circuit breaker; unrelated validation, structural-document, and collateral
DLQ entries do not enter the numerator. Equality is allowed: a threshold of
0.05 stops only when the ratio is strictly greater than 5%. 0.0 stops on
the first type failure, while 1.0 never trips. Values must be finite and in
[0.0, 1.0].
Correlation key
Declare correlation_key on the contributing Source’s config: block, not on
error_handling:. Group DLQ rejections by a key field. When any record in a correlation group fails, records from the failing source’s contribution to that group are routed to the DLQ:
# Inside a Source's config:
correlation_key: order_id
For compound keys:
# Inside a Source's config:
correlation_key: [order_id, customer_id]
This is useful for transactional data where partial processing of a group is worse than rejecting the entire group. For example, if one line item in an order fails validation, you may want to reject the entire order.
Under multi-source ingest, the collateral fan-out narrows to the failing source: a src_b trigger does NOT DLQ records from src_a that share the same correlation key. Single-source pipelines see bit-identical behavior to today’s pipeline-wide collateral DLQ. See Per-source rollback narrowing for the full semantic and the two documented exceptions (max_group_buffer overflow and Combine output failures).
When a Combine output row fails under a correlation key, the failing driver
row is the trigger of the driver’s correlation group. The dead letter for the
matched build record is held with that group as a collateral
(_cxl_dlq_trigger: false, category combine_output_row), written right after
its driver’s row with its driver’s _cxl_dlq_trigger_id, and written or rolled
back exactly when that group is. It never condemns the build record’s own
correlation group, so another driver that matched the same build record keeps
its output unless its own group failed. Each failure gets its own copy of the
build row, paired with its own driver row: two failing drivers of one group
each get one, and a driver that fails against several build rows is written
once per failure, each copy followed by the build row of that failure. The
held failure keeps its triggering field and value, as it does without a key.
Because every failure is written, a correlation key does not lower
dlq_count, records_dlq or the max_rate numerators: they equal the
counts the same failures give without a key, plus the rows the failing
groups condemn.
For the full lifecycle and per-operator semantics (route, merge, aggregate, combine), see Correlation Keys.
Max group buffer
Limit the number of records buffered per correlation group:
max_group_buffer: 100000 # Default: 100,000
A group that goes over this limit is dead-lettered whole. Its failing rows
are written as their own triggers; its other rows are written under one
group_size_exceeded row, as correlated rows. The limit counts held
entries, not distinct rows: see
Correlation Keys.
Document-level DLQ
By default a record failure dead-letters only that record (dlq_granularity: record). A source can instead reject the entire document any record of which fails, by declaring the granularity per source:
nodes:
- type: source
name: claims
config:
name: claims
type: x12
glob: ./claims/*.edi
schema: [{ name: seg_id, type: string }]
dlq_granularity: document # record (default) | document
Under dlq_granularity: document and the continue strategy, a document is rejected when one of its records fails at the Source (against its declared type, or a structural rule), in a Transform, or in a Route. A record that fails in an Aggregate, a Combine or a Reshape is dead-lettered as its own row, as under record granularity, and does not reject its document (#1232; see Not covered). When a document is rejected:
- the failing record becomes the root-cause DLQ entry (
_cxl_dlq_trigger = true, carrying its original error category); - every other record of the document that reaches a Sink becomes a collateral entry (
_cxl_dlq_trigger = false, categorydocument_rejected); a record dropped before any Sink is not written; - every other record of the document that fails at the Source, in a Transform or in a Route is also a
document_rejectedcollateral, written right after the root cause in the order the records failed. These records are stamped (_cxl_dlq_id,_cxl_dlq_timestamp) when they fail; the records a Sink held are stamped when the Sink rejects the document. Every collateral carries the root cause’s id in_cxl_dlq_trigger_id; - no record read from that document is written by any Sink, however many Sinks read it. Rows derived from the document’s records can still reach a Sink; see Not covered.
Clean documents in the same run stream through untouched, and records from sibling sources still on the default record granularity keep per-record semantics — the policy is per source.
Sinks run last. Every Sink runs after every other node, so a document’s verdict is final before any Sink writes one of its records. --explain lists the Sinks last. Dead-letter rows that other nodes write directly come before the rows a Sink writes.
Several Sinks. Each source row of a rejected document appears in the DLQ once, however many Sinks reached it. A row that failed is written as it was when it failed, ahead of any Sink’s copy of it. Any other row is written as it was held by the first Sink, in run order, that held it; a row that reached only a later Sink (a Route sent it there, say) is written by that Sink. Rows are matched by source row, so of the records one emit each makes from a source row, one is written. Every collateral, whichever Sink writes it, names the document’s root-cause entry in _cxl_dlq_trigger_id.
There is one exception, described under Not covered: a record that an Aggregate, a Combine or a Reshape failure dead-lettered (for a Combine, the driver row and the build row it matched) is written for that failure. It is written again, as a document_rejected row, only when a Sink on another branch also received it and its document is rejected; if its document is not rejected, that Sink can publish it (#1232). The rows a Combine or an Aggregate writes for a rejected document, and the rows that pass through a Reshape, are not held back at all.
This is the document-shaped analogue of correlation keys: use it when partial processing of a document (an EDI interchange, a batch file with a header/trailer) is worse than rejecting the whole document. Unlike correlation keys, which group across files by a key value, document-level DLQ scopes rejection to a single document’s records.
Document grain. The document is the outermost level — the source file. For a flat format (CSV, JSON, plain XML) each input file is one document. For a nested-envelope format (an X12 ISA → GS → ST interchange, an EDIFACT UNB → UNG → UNH) the document is the whole interchange / file, not an inner functional group or transaction set: a failure anywhere in the interchange rejects the entire interchange, including the transaction sets that validated cleanly. Reject the inner-level grain instead by partitioning the input so each interchange is its own file is not currently offered — the grain is fixed at the file.
DLQ rate. Each dead-letter row — the trigger and each collateral — counts once toward the configured DLQ max_rate, matching the correlated-collateral precedent, however many Sinks held it; a record written twice under the exception in Several Sinks counts twice. A rejected 1000-record document contributes 1000 when every one of its records reaches a Sink. It does not affect type_error_threshold, whose numerator contains only declared source-type failures.
dlq_count never counts a row that ok_count also counts, because no Sink writes a record read from a rejected document, with these exceptions, each described under Not covered:
- a record that an Aggregate, a Combine or a Reshape failure dead-lettered, which a Sink on another branch can still write when its document is not rejected (#1232);
- a rejected document’s rows that a Combine or an Aggregate writes, or that pass through a Reshape, which reach a Sink while the document’s records are dead-lettered;
- a record whose joined or aggregated row fails in a node after the Combine or Aggregate: that failure is written while the record itself can still reach a Sink on another branch.
Memory. The engine buffers each open document’s records until its boundary, then flushes the document clean to the sink or rejects it and drops the buffer. Peak memory scales with the concurrently-open documents, not the total input; a single very large document spills its buffer to disk under the run’s memory budget rather than holding everything in RAM. See Streaming vs blocking for the spill model.
Sink restriction. Document-level DLQ flushes each whole document to a single output writer, so it cannot be combined with a per-source-file Sink (a {source_file} / {source_path} path template over a multi-file source). The two are rejected together at compile time (E343); use a single output path, or set dlq_granularity: record if per-file output is the requirement.
Strategy requirement. dlq_granularity: document requires error_handling.strategy: continue. It is incompatible with the default fail_fast: document-level dead-lettering keeps the run going past a bad document, which contradicts fail-fast’s abort-on-first-error. The combination is rejected at compile time (E344) — set strategy: continue to dead-letter bad documents, or keep fail_fast with the default dlq_granularity: record.
Correlation restriction. Document-level and correlation-key rejection are
alternative atomic-disposition models. A document is keyed by its source file;
a correlation group can span files and is keyed by authored field values. The
engine does not define precedence or a combined writer boundary for those two
populations, so a pipeline containing both dlq_granularity: document and any
correlation_key is rejected at compile time (E370). Remove every
correlation_key to keep document rejection, or set dlq_granularity: record
to keep correlation rejection.
Composition restriction. A composition body may not declare a Sink when any
Source uses dlq_granularity: document (E378): a body Sink runs inside its
composition, where it cannot be held back until every document’s verdict is
final. Declare the Sink at pipeline level instead. Where moving it through a
new composition output port works today, the diagnostic prints that move ready
to paste. Otherwise it names clinker explain --code E378, which shows how to
declare the Sink’s work at pipeline level.
Not covered. Document rejection does not yet reach these rows:
- A row failure inside an Aggregate, a Combine or a Reshape is dead-lettered
per record and does not reject its document
(#1232). When a Sink on
another branch also received that record, it is written again as a
document_rejectedrow if another failure rejects its document, and that Sink can publish it if none does. - The rows a Combine or an Aggregate writes are not held back by their document’s verdict, so a rejected document’s joined or aggregated rows still reach a Sink. A failure on such a row, in any node after the Combine or Aggregate, dead-letters only that row and does not reject its document, so the document’s other records, and that row’s own source record on another branch, can still be published (#1317).
- A Reshape’s output rows do not keep their input row’s document, so a rejected document’s rows that pass through a Reshape still reach a Sink.
- In a pipeline where any Source declares
dlq_granularity: document, a CSVjoin_valuescollision withon_conflict: errorat a Sink fails the run rather than dead-lettering the record (#933).
A document is identified by its file path, so two Sources reading the same file share one verdict.
Spilling stages. Document identity survives memory pressure end to end. The per-document buffer identifies each document before buffering and spills under the memory budget, and a Sort between the source and the output preserves each record’s document context — including the source file the grain keys on — across its own spill round-trip. A document whose records pass through a spilling Sort is therefore still grouped and rejected as one document under memory pressure, exactly as it would be in memory. The rows a hash Aggregate or a grace-hash Combine writes are not held back by their document’s verdict, spilled or not; see Not covered.
Malformed envelopes (structural validation)
Envelope formats carry their own structural-integrity claims: an X12 interchange declares a segment count in each SE/GE/IEA trailer, EDIFACT in each UNT/UNZ, HL7 batch/file in each BTS/FTS, and a multi-record flat file’s trailer record declares a body count via its structure: constraint. When the declared count does not match the body the reader actually streamed, the file is structurally invalid. A multi-record flat file can also break a non-count structural rule — a line whose record-type discriminator matches no declared records: entry (E345), or a body record appearing after the trailer that closes the document; these are classified separately from a count mismatch but carry the same disposition.
Under dlq_granularity: document, such a structural failure dead-letters the whole source file to the DLQ rather than aborting the run:
- the file’s records dead-letter as one
structural_validationroot-cause entry (_cxl_dlq_trigger = true) plus adocument_rejectedcollateral for every other already-streamed record of the file; - no record of the malformed file reaches the success sink.
nodes:
- type: source
name: claims
config:
name: claims
type: x12
glob: ./claims/*.edi
schema: [{ name: seg_id, type: string }]
dlq_granularity: document # reuse the document opt-in — no separate config
The opt-in is the same dlq_granularity: document knob that governs per-record document rejection above; there is no separate validation: block. A malformed envelope is simply one more reason a source under the document policy condemns a whole document. E345 is the one structural class that also has a record-grained recovery: under strategy: continue and the default dlq_granularity: record, only the unknown-tag row is dead-lettered and the reader continues at the next physical row. The DLQ row exposes the unknown tag as record_type and the unguessed decoded input as _cxl_dlq_source_record (fixed-width line text, or a JSON array of decoded CSV cells).
Honest timing — rejected at the sink boundary, not before the first record. The trailer that carries the count arrives at the end of the file, after every body record it counts has already streamed through the DAG. Clinker is a bounded-memory streaming engine — it does not buffer the whole file up front to pre-validate it (that would defeat the streaming model). So the count mismatch is detected mid-stream, the file is marked failed, and the document-level DLQ buffer rejects every already-streamed record of the file at its close. The user-visible outcome is the same — no record of a malformed envelope is ever written to the output — but the rejection lands at the sink boundary, not literally before the file’s first record streams.
Grain — the whole file. An SE-level mismatch (one transaction set inside a larger interchange) rejects the entire interchange / file, not just that one transaction set, because the document grain is the outermost source file (see Document grain above). Split the input so each interchange is its own file if you need finer rejection.
Multiple files keep flowing. When a glob / paths source matches several files and one is malformed, only that file dead-letters — ingestion continues to the remaining files, so the clean files after a bad one still reach the sink. (This is unlike the default record granularity, where a count mismatch aborts the whole run and no file’s records are written.) Dead-lettering one malformed file never silently drops the good files around it.
Record-grained E345 is narrow. Under the default dlq_granularity: record, strategy: continue can recover only from an unknown multi-record discriminator because the reader has consumed exactly one bounded physical row and can resume unambiguously. Trailer-count mismatches and a body record after a document-closing trailer still abort at record granularity: neither belongs to one independently recoverable row. Genuine corruption (a truncated stream, a bad delimiter, a control-number echo mismatch, a segment after an X12/EDIFACT/HL7 envelope trailer) always aborts, even under the document opt-in. Under fail_fast, E345 also aborts at the offending line.
Cryptographic integrity (checksums / signatures) is not yet validated. Envelope formats can also carry a SHA-256 body hash, a JWS-signed JSON payload, or an XML Signature. Clinker extracts these envelope sections but does not yet verify them. Tracked for a future release.
Exit codes
| Code | Meaning |
|---|---|
| 0 | Pipeline completed successfully, no errors |
| 1 | Configuration error – the pipeline never started |
| 2 | Pipeline completed, but DLQ entries were produced |
| 3 | Data error halted the run: a fail_fast evaluation/accumulator failure, or the DLQ-rate ceiling |
| 4 | I/O, format, or spill failure |
Exit code 2 is not a failure – it means the pipeline ran to completion and handled errors according to the configured strategy. Check the DLQ file for details. See Exit Codes & Error Diagnosis for the full reference and the orchestrator retry policy.
Complete example
pipeline:
name: order_processing
memory: { limit: "512M" }
nodes:
- type: source
name: orders
config:
name: orders
type: csv
path: "./data/orders.csv"
correlation_key: order_id
schema:
- { name: order_id, type: int }
- { name: customer_id, type: int }
- { name: amount, type: float }
- { name: email, type: string }
- type: transform
name: validate_orders
input: orders
config:
cxl: |
emit order_id = order_id
emit customer_id = customer_id
emit amount = amount
emit email = email
validations:
- field: email
check: "not_empty"
severity: error
message: "Customer email is required"
- check: "amount > 0"
severity: error
message: "Order amount must be positive"
- type: sink
name: valid_orders
input: validate_orders
config:
name: valid_orders
type: csv
path: "./output/valid_orders.csv"
error_handling:
strategy: continue
dlq:
path: "./output/rejected_orders.csv"
include_reason: true
include_source_row: true
type_error_threshold: 0.10