Links and workers

Declarative materialization with exactly-once apply, run by a stateless worker pool.

Links are Operon's replacement for connector pipelines: declared, continuously maintained materializations from streams into tables, collections and graphs. Workers are the stateless pool that runs links and every other background task.

The design uses SQL DDL, with a stateless transform over each batch:

CREATE LINK tickets_search
  FROM STREAM support_tickets
  INTO COLLECTION tickets
  TRANSFORM (
    SELECT key AS _pk,
           json_get_str(value, '$.subject') AS subject,
           json_get_str(value, '$.body')    AS body,
           timestamp AS ts
  )
  WITH (batch_interval = '2s', on_error = 'dead_letter');

Transforms cover projections, filters, JSON extraction, casts and scalar functions, including an optional embed() that calls an external embedding endpoint. Stateful joins and windowed aggregations are out of scope; run RisingWave alongside Operon for those.

What exists today

M0 has the link framework with one target kind, counter, which sums integer values per key. It proves exactly-once apply end to end. SQL DDL and the collection, table and graph targets come with their milestones. See the quickstart to try it.

Exactly-once apply

A link task leases its source partitions in the metastore, with an epoch.

It reads from the last applied offset recorded in the target's latest commit.

It builds the target's files, then commits the target together with the new applied offsets, in one conditional pointer swap.

A zombie task with a stale epoch fails its commit. A restarted task resumes from the committed applied offset, so work repeated after a crash is discarded, never applied twice.

What workers run

TaskOutput
SegmenterStream segments; WAL objects deleted
RetentionSegments and files past their policy deleted
Link applyTarget commits
Split merge, Lance compaction, Iceberg compactionLarger, fewer files
Hot artifact buildHNSW graphs and projections
Graph sidecar buildCSR and CSC sidecars
Garbage collectionUnreferenced objects deleted after a grace period
Metastore snapshotA snapshot in the bucket

Read more in §09 Links and workers.

On this page