Skip to content

Latest commit

 

History

1,157 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

clink

ci docs licence changelog

clink is an embedded-first, Arrow-native stream processing engine in modern C++ (C++23): stateful stream processing with engine-grade semantics - SQL, event time, keyed state, exactly-once checkpoints - that you run like a tool rather than operate like a platform.

The whole engine lives in one library, and the same pipeline runs two ways. In-process: clink run pipeline.sql executes it in a single process with no daemons and prints its first result about 155 ms after process start; libclink embeds the engine in any service behind a pure-C ABI with results as Arrow C streams; pyclink returns them as pyarrow tables in a notebook. At scale: the same SQL file, unchanged, submits to a distributed Coordinator/Worker cluster with parallelism, failover, and rescale.

Three capabilities follow from that design:

  • State is an open dataset. Snapshots are documented Arrow IPC: checkpoints and savepoints open directly in pyarrow, DuckDB or Polars, export to Parquet and Iceberg, and a running job's live state serves point lookups and Arrow scans over plain HTTP, no sink round-trip. See state and backends.
  • Incidents replay deterministically. A flight recorder captures what each operator consumed per checkpoint epoch; clink replay re-executes an operator over exactly those records, offline and byte-identically, and can freeze the incident into a permanent regression test. See replay determinism.
  • Nothing to manage under it. A single static binary with no managed runtime: cold start in milliseconds, one artefact to ship, and the same behaviour embedded, in CI, and on a cluster.

Measured, not asserted: across the 17-query nexmark suite on a five-node cluster, clink processes an event for 1.9x to 5.3x less CPU (median 2.45x) than a JVM stream processor producing identical, correctness-gated output. Method, caveats and raw per-run data: Benchmarks, priced out in instances, dollars and modelled CO2e at Cost and environmental footprint.

clink is heavily inspired by Apache Flink. Flink's model of typed operator DAGs, event-time processing, in-band watermarks and checkpoint barriers, keyed state, and exactly-once semantics is the conceptual foundation this engine builds on. clink reworks that model in modern C++ around Arrow-native columnar execution and JVM-free deployment.

Try it in one command

Nothing to clone, install or configure:

docker run --rm ghcr.io/orhaugh/clink-runtime:latest run /opt/clink/examples/sql/hello.sql

The pipeline and the data it reads both ship inside the image (amd64 and arm64). It groups a small stream of device readings by region and prints the result:

{"avg_reading":21.4,"peak":21.4,"readings":1,"region":"emea"}
{"avg_reading":22.65,"peak":23.9,"readings":2,"region":"emea"}
{"avg_reading":19.2,"peak":19.2,"readings":1,"region":"amer"}
...

Those rows are the point. A batch engine prints one line per region at the end; clink prints a line per update, because the aggregate is a stream and you are watching it change. Run your own file the same way by mounting it:

docker run --rm -v "$PWD:/work" ghcr.io/orhaugh/clink-runtime:latest run /work/mine.sql

Want to build something real? Follow the Kafka to clink to ClickHouse tutorial: one docker compose up gives you a broker, a database and a small clink cluster running an event-time windowed pipeline; you then kill the Worker mid-stream, watch it recover from its checkpoint, verify every result independently, and open the job's state as an Arrow table. The example lives in examples/kafka-to-clickhouse.

Build from source

clink builds against a pinned Apache Arrow toolchain in ~/.clink-deps. The bootstrap step downloads a prebuilt, checksum-pinned archive when one exists for your platform (macOS arm64, Linux x86_64/arm64; about a minute) and compiles it from source otherwise:

git clone https://github.com/orhaugh/clink && cd clink
scripts/build-arrow.sh && scripts/build-iceberg-cpp.sh   # one-time, cached
cmake -S . -B build && cmake --build build --parallel 10
ctest --test-dir build --parallel 8                      # optional

Then run a first pipeline - one process, no daemons:

printf '{"usr":"alice","amount":12}\n{"usr":"bob","amount":7}\n{"usr":"alice","amount":5}\n' > /tmp/orders.ndjson

./build/clink run -e "CREATE TABLE orders (usr VARCHAR, amount BIGINT) \
      WITH (connector='file', format='json', path='/tmp/orders.ndjson'); \
    SELECT usr, SUM(amount) AS total FROM orders GROUP BY usr"

Supported platforms: macOS (Apple Silicon is the primary development platform) and Linux (Debian-family, exercised in CI). Windows is not supported. ./build_and_test.sh is the reproducible CI-matching path, including the sanitizer matrix; CONTRIBUTING.md has the details.

One engine, two execution models

Embedded. clink run pipeline.sql runs the whole engine in one process: SQL frontend, operators, state backends, checkpointing, connectors. First result in about 155 ms from process start, a figure gated by a Release-build test so it cannot silently regress. The same engine embeds behind a pure-C ABI (libclink), from Python (pyclink), and over Arrow Flight SQL for any ADBC/JDBC client.

Distributed. The same SQL file or compiled job plugin submits to a Coordinator/Worker cluster: parallel subtasks, hash-partitioned keyed shuffles, exactly-once checkpoints, failover from the last completed checkpoint, hot per-operator rescale, and rolling upgrades through savepoints. A Helm chart and Kubernetes operator ship in-tree.

clink_node --role=coordinator --port=6123 --http-port=8081 &
clink_node --role=worker --coordinator-host=127.0.0.1 --coordinator-port=6123 &
clink run pipeline.sql --coordinator-host=127.0.0.1 --coordinator-port=8081

Workers join the coordinator on its RPC port; a SQL submission goes to its HTTP port, and a compiled job plugin submits over RPC (clink run --job=pipeline.so --coordinator-host=127.0.0.1 --coordinator-port=6123).

Typed C++ pipelines (the fluent Pipeline / DataStream<T> API, keyed process functions, windows, joins, CEP) are documented with runnable programs under docs/consumer-examples/.

Production qualification

Benchmarks say how fast an engine is; they say nothing about whether its guarantees hold when processes die at the worst possible instant. clink runs a standing qualification programme: long campaigns on disposable multi-host rigs with faults injected into the narrowest windows of the engine's own protocols, judged by an independent oracle, published only when green and with the evidence retained. A few of the published results:

Qualified Result
Kafka exactly-once 755/755 windows byte-exact after two hours of kills inside commit windows, coordinator SIGKILLs, broker outages and partitions
Large keyed state 29 GiB of keyed state on a disaggregated backend, every key correct under the same fault battery
Wide job graphs 147 operators as 292 subtasks, exactly once under faults, 28-second recovery from a worker kill at that width
Rolling upgrade An engine upgrade with exactly-once continuity: 2 s savepoint, 2 s restore, every event across the boundary counted once
Semantic comparison 19 of 19 queries content-equal with an independent reference engine under pre-declared judgement classes

The full table, the rig, the method and every campaign's honesty-bounded claims: Qualification.

The protocol those campaigns exercise is also written down as a TLA+ specification and model-checked on every push, with every defect the campaigns ever found kept as a mutant the checker must refute. The model found three further interleavings the campaigns had not, fixed before any rig paid for them. The engine also records a protocol trace on request, and every trace the test suite leaves is model-checked against the specification on each push, so the code's recorded behaviour and the model are held in agreement rather than assumed to be. It proves the model and the recorded runs, not the code in general: Exactly-once specification.

Status and maturity

clink is young and pre-1.0. Its guarantees are qualified within the published bounds above, not battle-tested by years of third-party production deployments. The public surface is tiered ahead of 1.0 (Compatibility): the Stable tier of the C++ API, the C ABI and the SQL dialect is declared in tracked manifests and frozen conformance suites, and every change to it before 1.0, and every Evolving-tier change after, is called out in the CHANGELOG. Durable state is treated conservatively: snapshots carry schema versions with a migrate-at-restore path, so an upgrade does not silently invalidate checkpoints or savepoints. What the engine can do, feature by feature with caveats stated, lives in the capability catalogue.

Road to 1.0

The path to 1.0 is primarily about stability and operational evidence rather than adding another broad layer of features.

Settled, each under a design record and held by gates in CI:

  • Stable public APIs. The C++, C and SQL-facing interfaces that carry compatibility guarantees across the 1.x line are declared by tier and held by conformance gates: a tracked header manifest, an append-only C symbol manifest, compile-only conformance units and a frozen SQL corpus (design record 011, Compatibility). 1.0 freezes the manifest.
  • Stable extension model. A compiled job or plugin loads on any engine build whose declared extension surface matches. The contract is a tracked header manifest plus the build options, pinned Arrow version and toolchain identity that shape its layout; a refusal names the differing headers, an incompatible submit is refused before any bytes ship, and out-of-tree modules build with the packaged clink_add_job_module() (design record 010).
  • Native type ergonomics. One CLINK_FIELDS declaration per type derives the byte codec, the Arrow schema and columnar batcher, distributed type registration, and a shape fingerprint that refuses a restore whose field list changed without a declared version bump (design record 009).

Still open before 1.0:

  • Long-duration qualification. Complete multi-day production-style campaigns and continue expanding the measured state, scale and recovery envelope.
  • Operational maturity. Incorporate real-world deployment feedback, strengthen upgrade and diagnostic tooling, and close issues found by external workloads.

1.0 will mean that these compatibility and operational contracts are ready to be relied upon. It will not mean that feature development stops.

See GitHub issues for work currently in progress.

Documentation

Everything below is published at orhaugh.github.io/clink:

  • Capability catalogue: the complete shipped feature surface, with caveats.
  • Diagnosing a pipeline with an agent: the diagnostic surface as MCP tools, read-only, walked through a real incident by any MCP client.
  • SQL reference: the supported SQL surface, from DDL through windows, joins, MATCH_RECOGNIZE, UDFs and SQL-native ML.
  • Connectors: twenty-plus sources and sinks (Kafka, Postgres incl. CDC, ClickHouse, S3/GCS/Azure Parquet, Iceberg, Delta Lake, MQTT, NATS, Pulsar, RabbitMQ, Redis, MongoDB, Cassandra, HTTP, Avro, WebSocket and more), each with dependencies, options and delivery semantics; Kafka speaks the Confluent Schema Registry wire format (Avro, Protobuf, JSON Schema).
  • Internals: every subsystem documented the way its code is structured, citing sources.
  • Design decisions: why the engine is built this way, trade-offs included.
  • Benchmarks and cost footprint.
  • Qualification.
  • Tutorial: Kafka to clink to ClickHouse on your machine, with a deliberate Worker kill and independent verification (examples/kafka-to-clickhouse).
  • Runnable examples: every core feature as a standalone buildable program, from hello-pipeline to the testing framework and state-as-data workflows.

Installing

clink ships as a regular CMake package:

cmake -S . -B build -DCLINK_BUILD_TESTS=OFF -DCLINK_BUILD_EXAMPLES=OFF \
                    -DCMAKE_INSTALL_PREFIX=/usr/local
cmake --build build --parallel 10
sudo cmake --install build

Downstream projects consume it with find_package(clink REQUIRED) and link clink::clink (or per-impl targets such as clink::kafka); docs/consumer-examples/ is a complete find_package-based project to copy from. The install carries the clink CLI and the clink_node cluster daemon. Optional connectors and backends are CLINK_WITH_<NAME> CMake options (default AUTO: used when the dependency is found); the SQL frontend is CLINK_BUILD_SQL=ON by default. clink --capabilities prints what any given binary was built with.

Getting help

Licence and attribution

clink is licensed under the Apache License 2.0. See LICENSE and NOTICE. You are free to use clink for any purpose, including in commercial and closed-source products. If you build on it, a link back to this repository is appreciated; for write-ups, please cite it via CITATION.cff.

About

Embedded-first, Arrow-native stream processing engine in modern C++ (C++23). Run SQL streaming pipelines in-process in milliseconds, or submit the same file to a distributed cluster: exactly-once, event time, state as open Arrow data, deterministic replay.

Topics

Resources

Contributing

Security policy

Stars

4 stars

Watchers

0 watching

Forks

Releases

Packages

Used by

Contributors

Languages