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

Source Nodes

Source nodes read data from files and are the entry points of every pipeline. They have no input: field – they produce records, they do not consume them.

Basic structure

- type: source
  name: customers
  config:
    name: customers
    type: csv
    path: "./data/customers.csv"
    schema:
      - { name: customer_id, type: int }
      - { name: name, type: string }
      - { name: email, type: string }
      - { name: status, type: string }
      - { name: amount, type: float }

Schema declaration

The schema: field is required on every source node. Clinker does not infer types from data – you must declare each column’s name and CXL type explicitly. This schema drives compile-time type checking across the entire pipeline.

Each entry is a { name, type } pair:

schema:
  - { name: employee_id, type: string }
  - { name: salary, type: int }
  - { name: hired_at, type: date_time }
  - { name: is_active, type: bool }
  - { name: notes, type: nullable(string) }

Available types

TypeDescription
stringUTF-8 text
int64-bit signed integer
float64-bit IEEE 754 floating point
boolBoolean (true / false)
dateCalendar date
date_timeDate with time component
arrayOrdered sequence of values
anyUnknown type – field used in type-agnostic contexts
nullable(T)Nullable wrapper around any inner type (e.g. nullable(int))

A source column’s declared type must be concrete: numeric — the inference-only int | float union CXL resolves during type unification — is not a valid source column type. Declaring one is rejected at compile with E158; declare int or float explicitly.

long_unique — storage hint for high-cardinality text

A string column may carry an optional long_unique: true flag. It is an advisory, opt-in hint, not a type change: it tells Clinker the column’s values are long and effectively unique — never repeated across records — so the run uses less memory for that column. Typical candidates are UUIDs rendered as text, street addresses, and free-text comment or note fields.

schema:
  - { name: ticket_id,  type: string, long_unique: true }   # 36-char UUID
  - { name: notes,      type: string, long_unique: true }   # free text
  - { name: department, type: string }                      # low-cardinality, default

The flag lowers memory use only. A value’s content and its comparison, grouping, join, sort, and output behavior are all unchanged — a long_unique value behaves identically to the same text in any other column. Omitting the flag (the common case) leaves the default behavior untouched. Set it only when you know a column is genuinely high-cardinality free text; on a column whose values repeat, leave it off.

source_name — read a differently-named physical column

A column may carry an optional source_name naming the physical input column it reads from, when that differs from the exposed name. The reader matches input fields by physical name and re-labels the value under name, so downstream CXL and the output see name carrying the physical column’s data.

schema:
  # read the physical `cust_id` column, expose it downstream as `customer_id`
  - { name: customer_id, type: string, source_name: cust_id }

Omitting source_name (the common case) reads the input field whose key equals name, unchanged from before. A channel schema patch’s rename op sets this alias automatically (see Channels).

Transport vs format

A source declaration has two independent layers:

  • Transport (transport:) selects where the records come from. Two transports exist: file — read bytes from the filesystem, resolved through one of the file matchers (path / glob / regex / paths) — and rest — pull records from a paginated HTTP endpoint under a hard page/record cap (see Network Sources (REST)). transport: is optional and defaults to file, so a source that omits it reads from disk exactly as before.
  • Format (type:) selects how the bytes decode into records: csv, json, xml, fixed_width, edifact, x12, hl7, swift.
- type: source
  name: orders
  config:
    name: orders
    transport: file        # optional; this is the default
    type: csv              # the on-disk format
    path: "./data/orders.csv"
    schema:
      - { name: order_id, type: int }

A file transport requires exactly one file matcher (path, glob, regex, or paths). Declaring none fails validation with E211; declaring more than one fails with E210. Both are reported at config-load time, before any file is opened.

Format types

The type: field inside config: selects the on-disk format. Each format has its own reference page covering its options and decoding model:

type:FormatReference
csvDelimited text (RFC 4180)CSV Format
jsonArray / NDJSON / wrapper objectJSON Format
xmlElement-path-selected record elementsXML Format
fixed_widthColumn-positioned legacy extractsFixed-Width Format
edifactUN/EDIFACT interchangesEDIFACT Format
x12ANSI ASC X12 interchangesX12 Format
hl7HL7 v2.x pipe-and-hat messagesHL7 v2 Format
swiftSWIFT MT (FIN) messagesSWIFT MT Format

The same schema: rules apply regardless of format: the reader maps each decoded record onto the declared schema, and undeclared input fields fall under the on_unmapped policy below.

on_unmapped — undeclared input fields

The per-source on_unmapped policy decides what to do with input fields the source’s schema: block does not name. Three modes — auto_widen (default), drop, reject:

- type: source
  name: orders
  config:
    name: orders
    type: csv
    path: "./data/orders.csv"
    on_unmapped:
      mode: auto_widen     # default; other values: drop, reject
    schema:
      - { name: order_id, type: string }
      - { name: amount, type: float }

See Auto-Widen & Schema Drift for the full specification: how undeclared columns flow through each downstream node type, the include_unmapped Output flag, E315 merge-policy mismatch, and fixed-width behavior.

Sort order

If your source data is pre-sorted, declare the sort order so the optimizer can use streaming aggregation instead of hash aggregation:

- type: source
  name: sorted_transactions
  config:
    name: sorted_transactions
    type: csv
    path: "./data/transactions_sorted.csv"
    schema:
      - { name: account_id, type: string }
      - { name: txn_date, type: date }
      - { name: amount, type: float }
    sort_order:
      - { field: "account_id", order: asc }
      - { field: "txn_date", order: asc }

Sort order declarations are trusted – Clinker does not verify that the data is actually sorted. If the data violates the declared order, downstream streaming aggregation may produce incorrect results.

The shorthand form is also accepted – a bare string defaults to ascending:

    sort_order:
      - "account_id"
      - { field: "txn_date", order: desc }

Watermarks

An event-time watermark declares which column on the source carries each record’s event time — the wall-clock instant the event happened, distinct from when Clinker read the row. When set, Clinker takes the column on every record, subtracts the source’s delay, and uses the result to track event-time progress so downstream time windows know when to close. The delay-corrected value is also stamped on every record as $source.event_time, the column a downstream time-windowed aggregate uses to assign records to windows.

- type: source
  name: clicks
  config:
    name: clicks
    type: csv
    path: "./data/clicks.csv"
    options:
      has_header: true
    watermark:
      column: event_ts       # must be date_time or date
      delay: 5s              # bounded out-of-order tolerance
      idle_timeout: 30s      # flip partitions to idle if quiet
    schema:
      - { name: user_id, type: string }
      - { name: event_ts, type: date_time }
      - { name: amount, type: int }

Fields:

  • column (required) — the schema column whose value is each record’s event time. The column’s declared type must be date_time or date. A column: that names a field absent from schema: raises E154; a column: whose declared type is neither raises E155.

  • delay (optional duration, default unset) — bounded out-of-order tolerance. Each record’s event time is shifted earlier by delay before being folded into the watermark, so the source’s effective watermark trails its observed max event time by this amount. Mirrors Flink’s BoundedOutOfOrdernessWatermarks. Without delay, the watermark advances strictly to the observed max — a single late record routes to the DLQ.

  • idle_timeout (optional duration, default unset) — if a source stays quiet longer than this, it stops holding back downstream window-close progress, so windows keep closing when one source pauses. Unset means the source never goes idle.

Durations use the suffixes ms, s, m, h, d. ms is matched before the single-character s, so 500ms reads as 500 milliseconds, not 500 seconds with a stray m.

A pipeline whose aggregate declares time_window: must have a watermark.column on every upstream-reachable source. Without it, event-time progress can never advance and the window can never close — the planner rejects this with E156.

Multi-value fields

A field that holds more than one value is declared on the schema column, not on the pipeline:

schema:
  - { name: order_id, type: string }
  - { name: tags, type: string, multiple: true }

multiple: true says the column holds zero or more values of its declared type. Reading collects every occurrence of the field into one array — a single occurrence is still an array, so downstream code never has to branch on how many values happened to arrive. A field absent from a record has no column at all and resolves to null, exactly as any other absent column does. CXL sees the column as an array; the declared type: describes each element and drives coercion. The declaration describes the shape of the data, so it serves both directions: a writer that can encode repetition reads the same declaration.

The split_to_rows, split_values, and join_values blocks below accept a compact shorthand (a bare field name, or a mapping that omits defaults). To see the fully-materialized form the engine actually runs — every default spelled out — print the canonical config with clinker config --resolved. It rewrites only those shorthand blocks and leaves the rest of the file untouched.

Both ends of the declaration are checked at compile, so a shape the formats cannot carry fails before a run starts rather than mid-stream:

FormatAs a sourceAs an output
jsonnative — an arraynative — an array
xmlnative — repeated child elementsnot yet (issue 916)
csvdelimited cell via split_valuesdelimited cell via join_values
fixed_widthdelimited cell via split_valuesnot yet (issue 918)
edifact, x12, hl7, swiftno — repetition is positionalno — repetition is positional

A multiple: true column reaching an output that cannot encode it is E359; one on a source that cannot produce it is E361. E359 covers an output’s own schema: block too — the attribute is direction-neutral, but the remaining writers do not encode repetition yet, so declaring it on such a sink would be accepted and ignored. Run clinker explain --code E361 for the full remediation of either.

csv and fixed_width read a multiple: true column through a split_values entry. Neither wire format repeats a field, but a cell’s text may hold several values separated by a delimiter. Declare that delimiter with a split_values entry and the reader parses the cell into the array the column holds. A multiple: true column no entry covers is rejected by E361 — the reader would have no delimiter and deliver the raw cell; either add the entry or leave the column single-valued and split it in a transform (tags.split(";")). The entry is read only on a single-schema source: a multi-record source of either format runs a backend that does not consume it. On the output side, a CSV sink joins a multiple: field into one delimited cell with join_values (defaulting to ; / on_conflict: error); fixed_width output is still pending (#918).

A split_values entry also recovers a CSV cell a sink wrote under join_values on_conflict: escape or encode_json: add escape: "\\" to un-escape an escaped delimiter, or json: true to read the whole cell as an embedded JSON array.

The segment formats are a permanent no, not a pending one. Repetition there is a positional coordinate rather than a list: a repeated composite is written as two axes interleaved in one element (11:B:1^12:B:2), which a flat array cannot represent without losing the component axis. The faithful shape is one column per coordinate — for HL7, that is what options.split_fields produces, with a writer that reassembles the wire field byte-for-byte.

One record per value: split_to_rows

split_to_rows fans a record out to one record per occurrence of a repeated field. Each entry is either a bare field name or a full mapping, and the two forms mix freely in one list:

- type: source
  name: invoices
  config:
    name: invoices
    type: json
    path: "./data/invoices.json"
    schema:
      - { name: invoice_id, type: int }
      - { name: customer, type: string }
      - { name: line_item, type: string }
      - { name: line_amount, type: float }
      - { name: line_no, type: int }
    split_to_rows:
      - tags                      # shorthand: field name, all defaults
      - field: line_items         # full form
        keep_empty: true
        mode: extract
        position_column: line_no
KeyDefaultMeaning
fieldThe repeated field, as a flattened dotted name
keep_emptytrueWhether a record whose field is empty or absent survives
modeextractextract — the occurrence becomes the record; split — the record shape is kept
position_columnnoneColumn receiving each occurrence’s 1-based position

The field is named as it appears in the input document, not as the schema exposes it: a column declared source_name: is addressed by that source_name. The same rule applies to split_values below.

keep_empty defaults to true. A record whose field holds an empty array, or carries no such field at all, is emitted with that field unset rather than disappearing. Several widely used engines drop the record instead; a vanished row is the costliest failure mode there is, so dropping is opt-in here.

mode: extract (the default) makes the occurrence the record: its own fields are lifted out from under the field name and every field outside the group is merged onto each output. {"orders": [{"id": 1}]} yields a top-level id, and repeated <Item><name> children yield name. When lifting lands an occurrence’s field on a name an outside field already occupies, the occurrence wins — it is the record, so its own value is not shadowed by the parent it was merged with. A position_column wins over both: you named it, so a field of the same name inside or outside the occurrence gives way to the index.

mode: split preserves the record shape: the occurrence’s fields keep their dotted path (orders.id, Item.name) and each output carries exactly one occurrence.

Entries apply in declaration order, so two entries multiply. Declaring the same field twice is rejected at compile (E358), as is fanning out a field the schema also declares multiple: true — the attribute collects the occurrences into one array, the fan-out spends them one per record, and a field cannot be both. On an XML source, two entries may not name nested element groups either (Item and Item.part): that reader assigns each element to one occurrence group by document position, and a nested pair leaves the inner group’s membership ambiguous.

A JSON source accepts a nested pair to produce a two-level expansion, but only when the outer entry declares mode: split. mode: extract lifts the occurrence’s own keys to the top level, which removes the dotted path the inner entry addresses — the inner entry would then match nothing and fan nothing out, so the pairing is rejected (E358):

    split_to_rows:
      - { field: orders, mode: split }
      - { field: orders.items, mode: split }

Several values in one cell: split_values

split_values parses a delimited cell into the several values a multiple: column holds. It takes the same bare-name-or-mapping shorthand:

    split_values:
      - tags                      # shorthand: default delimiter `;`
      - field: codes              # full form
        delimiter: "|"
    schema:
      - { name: tags,  type: string, multiple: true }
      - { name: codes, type: string, multiple: true }

The delimiter defaults to ;. A split_values field the schema does not declare multiple: true is rejected at compile (E358): splitting produces several values, and only a multi-value column can hold them. So is an entry naming a column’s exposed name when that column reads a differently-named input field — the split runs against the document’s own field names, so name the source_name.

The entry is read by the JSON and XML readers, and — on a single-schema source — by the CSV and fixed-width readers. On a multi-record CSV or fixed-width source, or on any segment format, declaring it is rejected (E358) rather than silently ignored: those readers are never handed it, so the cell would arrive unsplit with nothing to say so.

Migrating from array_paths

array_paths: was the earlier form of these declarations. It is no longer read, and a source still carrying it is rejected at compile (E360) rather than running with the fan-out silently dropped. An explode path becomes a split_to_rows: entry, a delimited cell becomes a split_values: entry, and a path kept as an array becomes multiple: true on the schema column.

mode: extract (the default) reproduces the old projection — the element’s own fields lifted to the top level of each output record. It does not reproduce the old cardinality: explode dropped a record whose array was empty, while keep_empty defaults to true here and keeps it with the element’s fields unset. Add keep_empty: false to the entry to migrate row-for-row.

Format notes

split_to_rows is honored by the JSON and XML readers, over a file path and over a rest response body alike. split_values is honored by those two and also by the single-schema CSV and fixed-width readers. Declaring a knob on a reader that is never handed it — split_to_rows on any delimited-cell or segment format, split_values on a multi-record or segment source — is rejected at compile (E358) rather than accepted and inert, and a multiple: true column no split_values entry can cover is rejected by E361.

Composition body files are gated by the same four checks as the pipeline that calls them.

JSON — the field names the key holding the array (line_items, or order.line_items for an array nested under an object). A field present but holding a single object or scalar rather than an array counts as one occurrence, and is projected exactly as a one-element array would be — many producers unwrap a lone element, and XML cannot express the difference at all, so the two readers agree on this. A field that is absent, holds an empty array, or is explicitly null has no occurrence and is governed by keep_empty: an explicit null is how many producers write “no value”, and it counts as none rather than as one. For the same reason a multiple: true column holding an explicit null stays null rather than becoming [null], so size() over it reads the same as it does for a field the document omits.

XML — the field is the repeated child element’s dotted path relative to the record element (Item, or Items.Item when nested). Repetition and absence are indistinguishable in XML, so a record with no occurrence of the element is governed by keep_empty exactly as an empty array is. See XML Format for the full rules.