Skip to content

Shared Sinks

API Documentation

laktory.models.DataSinkSharedOptions

A shared sink is a table written by several writers: nodes of the same pipeline and/or other pipelines. For example, several feeds pooled into one table, nodes sending their quarantined rows to the same table, or one pipeline per client appending to a cross-tenant table.

Each writer owns its rows: a full refresh of a writer only deletes and reprocesses its own rows. Writers are therefore independent: they run in their own task, in any order, in parallel, alone or together.

Declaring a Shared Sink¤

Declare shared on every sink writing to the table (shared: true for the default options).

Several nodes of a pipeline:

nodes:
- name: feed_a
  sinks:
  - table_name: prices
    mode: APPEND
    shared: true
- name: feed_b
  sinks:
  - table_name: prices
    mode: APPEND
    shared: true

Nodes of a pipeline writing to the same table without shared fail validation: declaring it makes the writer column added to the table explicit.

Several pipelines:

# pl-client-acme.yaml, pl-client-globex.yaml, ...
sinks:
- table_name: all_orders
  mode: APPEND
  shared: true

Shared sinks must be DELTA table or file sinks in APPEND mode: other modes (OVERWRITE, MERGE) would modify the rows of the other writers. Write to separate tables instead.

Identifying the Rows of a Writer¤

Writer column (default): each row carries its writer in a _laktory_writer column (first column of the table).

_laktory_writer feed
pl-prices.feed_a a
pl-prices.feed_b b
Option Default Description
writer_id {pipeline_name}.{node_name} writer identifier
column _laktory_writer writer column name

Predicate: to avoid adding a column, declare the rows owned by the writer with where. A full refresh deletes the rows matching it.

# pl-client-23.yaml
sinks:
- table_name: all_orders
  mode: APPEND
  shared:
    where: client_id = 23    # full refresh: DELETE FROM all_orders WHERE client_id = 23
  • The predicate is written in the SQL of the backend: Spark SQL, or Polars / deltalake SQL with the Polars backend (e.g. IN, BETWEEN, LIKE).
  • A predicate also works for a single writer, when other processes (backfills, manual loads) write to the same table.
  • All the writers of a table must use the same kind: writer column or where.
  • Predicates must not overlap, and a writer must only write rows matching its predicate: other rows are not deleted by its full refresh.
  • The number of deleted rows is logged: check it to catch a wrong predicate.

Reading a Shared Sink¤

Source Reads
node_name: feed_a the output of feed_a only: its rows, without the writer column
table_name: prices the whole table: the rows of all the writers

node_name returns the same data whether the node output is read from memory (same run) or from the table (e.g. a separate job task).

  • With declarative orchestrators, rows carry no writer: reading a writer with node_name (or {nodes.x} in a transformer) fails validation, since it would return the rows of all the writers. Read the table with table_name, or write the node to its own table.

Streaming Readers¤

A full refresh of a writer deletes its rows, and a streaming read of a Delta table fails on deleted rows (DELTA_SOURCE_IGNORE_DELETE), including a node_name read of another writer. Choose per streaming reader:

Reader Full refresh of another writer Full refresh of the writer it reads
default fails: run a full refresh of the reader fails: run a full refresh of the reader
skipChangeCommits no effect: the deletes are skipped and the other writer's rows filtered out the writer's rows are read again: run a full refresh of the reader, otherwise they are duplicated
- name: gld_feed_a
  sources:
  - node_name: feed_a
    as_stream: true
    reader_kwargs:
      skipChangeCommits: true

With skipChangeCommits, a node reading a single writer (node_name) keeps running when the other writers are refreshed. A reader of the whole table (table_name) gets duplicates after a full refresh of any writer: keep the default there.

Resetting a Shared Table¤

Goal Run
Reprocess one writer refresh=FULL on its task (or node): deletes its rows, reset_mode doesn't apply
Drop or empty the whole table refresh=RESET with reset_mode=DROP or TRUNCATE, on any of its writers' tasks, then a normal run

When the whole table is reset:

  • The checkpoints of all its writers in the pipeline are reset too, whether or not they're part of the run: they reprocess all their data on their next run.
  • Other pipelines writing to the table need a full refresh: resuming from their checkpoints, they would skip the data they already wrote, and the table would silently miss their rows.
  • refresh=FULL with the override is rejected: a full refresh never deletes the rows of the other writers. Reset the table with refresh=RESET, then run normally.

Keeping Ownership Consistent¤

Situation What to do
Renaming a writer (node or pipeline) pin writer_id first: rows of the old identifier are no longer deleted by a full refresh
Removing a writer delete its rows (DELETE FROM ... WHERE _laktory_writer = '...') or reset the whole table
Table written before being shared (no writer column) writes fail until the table is dropped once (refresh=RESET, reset_mode=DROP)
Switching between writer column and where drop the table once
Parallel writers adding different columns declare the full schema, or order the writers with depends_on

Concurrent Writers¤

Writers of a shared table running at the same time can conflict: a writer deleting its rows (full refresh) while another one appends fails its Delta commit (ConcurrentAppendException with Spark, CommitFailedError with Polars). A failed commit writes nothing: Laktory retries it, up to 5 times with an increasing delay, and logs each retry.

To avoid most conflicts:

  • Databricks: use row-level concurrency (tables with deletion vectors) or liquid clustering, see Isolation levels and write conflicts.
  • Open-source Delta (Spark): partition the table by the writer column, so that a delete only reads the files of its writer:
sinks:
- table_name: prices
  mode: APPEND
  shared: true
  writer_methods:
  - name: partitionBy
    args: [_laktory_writer]

Validation¤

  • Every sink of a table written by several nodes of a pipeline declares shared.
  • Shared sinks are DELTA sinks in APPEND mode.
  • Writers of a table use the same kind of ownership, the same column, and distinct writer identifiers or predicates.

Pipelines writing to the same table are validated together when they are deployed together, in a Stack or a Databricks Asset Bundle:

Situation Result
A declarative pipeline writes to the table error: the engine owns the table
Sinks declaring shared identify their rows inconsistently (kind, column, identifiers) error
Some sinks don't declare shared warning: their full refresh or overwrite deletes the rows of the other writers

Sharing a table across pipelines is otherwise the responsibility of the user: pipelines deployed separately (other stacks, bundles or repos) can't be validated together, and a table is identified by its name or path as written in each pipeline.

Declarative Orchestrators¤

With Lakeflow / Spark Declarative Pipelines, the engine owns the tables: shared is not supported and a table written by a declarative pipeline can't be shared with other pipelines.

Several nodes can still write to the same table, without declaring shared: it's declared once as a streaming table, and each node appends to it through its own append flow, {table_name}__{node_name}. On a full refresh, the engine clears the table once and resets every flow.

  • All sinks must be streaming, non-CDC (MERGE) table sinks.
  • Table properties (comment, table_properties, format) can be set on any of the sinks, but must not conflict.
  • With Lakeflow Declarative Pipelines, expectations apply to the whole table: all its writers must declare the same expectations.
  • A flow is identified by its name: adding a second writer to a table, or renaming a node, restarts the existing flow from scratch. Run a full refresh of the table once after the change.