Connectors¶
Reference documentation for clink's source and sink connectors. Each connector has its own page covering the dependency and pinned version, the CMake build knob, the exact factory names, every configuration option, SQL usage where available, an example, delivery semantics, and limitations.
Every connector is an optional module. It is gated by a CLINK_WITH_<NAME>
CMake option (default AUTO: built when its client library is found, skipped
otherwise; set ON to require it or OFF to exclude it). Most connectors link
a system client library obtained via apt (Debian) or brew (macOS); a few ride
the from-source toolchain (Apache Arrow/Parquet 24.0.0, iceberg-cpp v0.3.0,
aws-sdk-cpp 1.11.795, Pulsar client 4.2.0, DataStax cpp-driver 2.17.1,
clickhouse-cpp 2.5.1, Avro C++ 1.12.1),
which is compiled at exact versions into CLINK_DEPS_PREFIX on both the host
and the Debian image. Versions are recorded per connector and in
scripts/versions.env.
The SQL connector= column lists the string to use in a SQL
CREATE TABLE ... WITH (connector='...') statement. A dash means the connector
is reachable through the programmatic API only.
Messaging and streaming¶
| Connector | I/O | Client dependency | Version | SQL connector= |
|---|---|---|---|---|
| Apache Kafka | source + sink | librdkafka | system pkg | kafka |
| Apache Pulsar | source + sink | Pulsar C++ client | 4.2.0 |
pulsar |
| RabbitMQ (AMQP 0-9-1) | source + sink | rabbitmq-c | system pkg | rabbitmq |
| NATS JetStream | source + sink | nats.c | system pkg | nats |
| MQTT | source + sink | libmosquitto | system pkg | - |
| WebSocket | source | none (OpenSSL for wss://) |
- | websocket |
Object storage and table formats (Parquet)¶
| Connector | I/O | Client dependency | Version | SQL connector= |
|---|---|---|---|---|
| Amazon S3 (Parquet) | source + sink | Arrow S3FileSystem + aws-sdk-cpp | Arrow 24.0.0, aws-sdk 1.11.795 |
s3_parquet |
| Amazon S3 (raw objects) | sink (+ programmatic line source) | aws-sdk-cpp | 1.11.795 |
s3 |
| Google Cloud Storage | source + sink | Arrow GcsFileSystem (ARROW_GCS) |
Arrow 24.0.0 |
gcs_parquet |
| Azure Blob Storage | source + sink | Arrow AzureFileSystem (ARROW_AZURE) |
Arrow 24.0.0 |
azure_parquet |
| WebHDFS / HttpFS | source + sink | clink::http_connector (vendored httplib) | Arrow 24.0.0 |
webhdfs_parquet |
| Apache Iceberg | source + sink | iceberg-cpp + Arrow | iceberg-cpp v0.3.0, Arrow 24.0.0 |
iceberg |
| Delta Lake | sink | SQL frontend + Arrow (aws-sdk-cpp for s3:// roots) |
Arrow 24.0.0 |
delta |
| Local files and Parquet | source + sink | core (Arrow for Parquet) | built in | file, filesystem, parquet |
Databases and key-value stores¶
| Connector | I/O | Client dependency | Version | SQL connector= |
|---|---|---|---|---|
| PostgreSQL | source + sink | libpq | system pkg | postgres |
| MySQL / MariaDB | source + sink | mariadb-connector-c | system pkg | mysql |
| ClickHouse | source + sink | clickhouse-cpp | 2.5.1 (from source) |
clickhouse |
| Cassandra / ScyllaDB | sink | DataStax cpp-driver | 2.17.1 |
cassandra |
| MongoDB | source + sink | mongo-cxx-driver | system pkg | - |
| Redis | source + sink | hiredis | system pkg | redis |
Cloud services and HTTP¶
| Connector | I/O | Client dependency | Version | SQL connector= |
|---|---|---|---|---|
| AWS (Kinesis / Firehose / DynamoDB) | source + sink | aws-sdk-cpp | 1.11.795 |
kinesis, firehose, dynamodb |
| HTTP (Elasticsearch, OpenSearch, Splunk, InfluxDB, Prometheus, poll, Pub/Sub) | source + sink | cpp-httplib (vendored) | vendored | http, elasticsearch, opensearch, splunk, influxdb, prometheus, http_poll, pubsub |
Serialization¶
| Format | Role | Dependency | Version | SQL connector= |
|---|---|---|---|---|
| Apache Avro | encoding (codecs) | Avro C++ | 1.12.1 (from-source pin) |
- |
| Schema Registry formats | value formats on Kafka (format='avro', 'protobuf', 'json-schema') |
Avro C++, libprotobuf + libprotoc, cpp-httplib (vendored) | Avro 1.12.1; protobuf system pkg |
kafka |
Built-in sinks (no dependency)¶
Compiled into the SQL frontend itself; always available when
CLINK_BUILD_SQL=ON. One shared page: built-in connectors.
| Connector | I/O | SQL connector= |
|---|---|---|
| Blackhole (discard) | sink | blackhole |
| Changelog netting | sink | changelog |
| Print (stdout) | sink | print |
| Collect (Arrow to host, embedded only) | sink | collect |
| Queryable state (another job's live state) | source | queryable_state |
Generator (GeneratorSource<T>, tests and benchmarks) |
source | - |
The capability manifest¶
Every compiled connector declares a machine-readable capability record next
to its factory registration: identity, formats, boundedness, recovery model
(replayable offset or broker redelivery), delivery guarantee as implemented
(not as the external system could theoretically provide), transactionality,
idempotency-key requirements, auth/TLS surface, and limitations. clink
--capabilities prints the manifest for the binary at hand, the delivery
analyser computes end-to-end guarantees from it at submission, and each
record's internal coherence is checked by self_check(). Coverage is
enforced mechanically for the impls/ modules: a build-generated list of
enabled connector modules feeds a gate test
(tests/test_connector_manifest_gate.cpp), so a new connector module cannot
register factories without either declaring its record or being explicitly
classified as a non-connector module. The gate does not reach the sinks the
SQL frontend registers itself: file, parquet, generator and blackhole
carry records, but the Delta Lake sink and the print,
changelog, collect and queryable_state built-ins do not, so they are
absent from clink --capabilities and unknown to the delivery analyser.
Notes on delivery semantics¶
Guarantees vary by connector and are stated on each page. In summary:
- Exactly-once sinks require the
on_barrier/on_committwo-phase-commit contract. The Kafka transactional sink (kafka_2pc_sink_string) implements it. The object-store and WebHDFS Parquet connectors offer both: the default single-object sink is at-least-once, and a 2PC variant (<connector>_2pc_*_sink, ordelivery_guarantee='exactly_once'in SQL) stages one file per checkpoint under<prefix>/stagingand atomically promotes it to<prefix>/committedonly when the checkpoint completes globally. - Sources that record their position as operator state replay from the last checkpoint on recovery; the exact mechanism (Kafka offsets, Postgres LSN, object index, row index) and any caveats are documented per connector.
- Messaging sources (RabbitMQ, NATS, Pulsar) acknowledge at the checkpoint barrier rather than post-commit; unacknowledged messages are redelivered.