Skip to main content

Crate krafka

Crate krafka 

Source
Expand description

§Krafka

A pure Rust, async-native Apache Kafka client.

Producer, transactional producer, consumer, share consumer and admin clients for Apache Kafka.

§Features

  • No C library, no system dependency: the default build needs a Rust toolchain and the C compiler cc finds for ring, nothing else. Every codec decodes in pure Rust; the opt-in zstd, rustls-aws-lc-rs and aws-msk features compile more C
  • Async-native: built on Tokio
  • Shared buffers: decoded record keys, values and header values are slices of the fetched response (one allocation per uncompressed batch, for the record list), and each response frame is read into one allocation
  • Safe: no unsafe code; panic, unwrap and expect are denied
  • Cloud-ready: TLS, SASL PLAIN/SCRAM/OAUTHBEARER, a built-in OIDC provider and AWS MSK IAM

§Quick start

use krafka::{Kafka, Record};

let kafka = Kafka::builder("localhost:9092").connect().await?;
let producer = kafka.producer().build().await?;
producer.send(Record::new("orders", "hello").key("k")).await?;

let consumer = kafka.consumer("my-group").build().await?;
consumer.subscribe(["orders"]).await?;
while let Some(rec) = consumer.recv().await? { println!("{rec:?}"); }

Kafka holds every connection setting — bootstrap servers, client id, TLS and SASL, transport, proxy, metadata — and one connection pool. Every client is built from it and shares the pool: Kafka::producer, Kafka::consumer, Kafka::share_consumer, Kafka::admin. A separate pool is a second Kafka.

Every client is Send + Sync; share one across tasks with an Arc, and end it with close().await.

§Errors

Every fallible call returns KrafkaError, which answers three questions: is_retriable (the same call may succeed), is_fatal (rebuild the client) and requires_abort (abort the transaction; the producer stays usable).

§Observability

  • Metrics. Every client’s metrics() returns one owned Metrics snapshot, including the connection counters of its pool; Kafka::metrics sums a handle’s clients. Metrics::prometheus_text renders one for a scrape endpoint.
  • Spans. The clients emit spans through tracing, on OpenTelemetry messaging semantic conventions OTEL_SEMCONV_VERSION (v1.44.0, whose messaging conventions are still Development): a send span per record from send/enqueue to its outcome (kind producer), a poll span per poll/recv and a commit span per commit (kind client), and a krafka-specific rebalance span per assignment change (kind internal). Record keys and values, and hashes of keys, are never recorded; a key’s size is. krafka depends on no OpenTelemetry crate: bridge tracing to OpenTelemetry in the application.
  • KIP-714 telemetry. Every client pushes its metrics to brokers whose operator subscribed to them, as Java clients do with enable.metrics.push: on by default for producers and consumers, off for the admin client, switched by metrics_push on each role builder. A cluster without a client-telemetry plugin receives nothing.

§Stability

The public API follows semver. Outside that promise: testing (the fake broker, behind test-broker), the unstable-protocol feature, and the hidden __private module behind the tooling-only internal feature, which exists for krafka’s own benches and fuzz targets.

§Cargo Features

FeatureDefaultDescription
ringyesrustls crypto backend using ring.
rustls-aws-lc-rsnorustls crypto backend using aws-lc-rs; offers post-quantum X25519MLKEM768 key exchange first. Compiles C and needs CMake.
zstdnoZstd encoding via zstd (compiles C). Zstd decoding is always available.
aws-msknoAWS MSK IAM authentication with the SDK credential chain (compiles C and needs CMake, via aws-lc-sys).
oauth-oidcnoBuilt-in OIDC token provider for SASL/OAUTHBEARER: the client_credentials grant (KIP-768) and RFC 7523 client assertions (KIP-1258). Adds no cryptography dependency — assertions are supplied pre-signed.
native-tls-rootsnoLoad platform-native root certificates via rustls-native-certs.
tls-encrypted-keysnoPassphrase-encrypted PKCS#8 client keys (ssl.key.password) via the RustCrypto pkcs8 crate.
unstable-protocolnoEnables protocol versions Kafka marks latestVersionUnstable — a released broker does not advertise them without unstable.api.versions.enable=true. Covers ApiVersions v5 (KIP-1242) and InitProducerId v6 (KIP-939). APIs under this feature may change without semver notice.
test-brokernoIn-process fake Kafka broker for testing your own code against a real client. Not for production builds; outside semver.

Always compiled in: gzip, Snappy and LZ4 compression, the KIP-932 share consumer (needs a Kafka 4.2+ broker), SOCKS5 proxy support, and KIP-714 client telemetry.

§TLS crypto backend

Exactly one rustls crypto backend is used at runtime. ring is the default; rustls-aws-lc-rs selects aws-lc-rs instead.

These features are additive, as Cargo requires: if two crates in your dependency graph each select a different backend, the build still succeeds and rustls-aws-lc-rs wins. To pin a specific backend regardless of what your dependency graph enabled, install it as the process default before constructing any krafka client:

ⓘ
rustls::crypto::ring::default_provider().install_default().ok();

Enabling neither backend is a compile error — rustls cannot build a ClientConfig without a crypto provider.

To disable the default features and pick only what you need, remember that default-features = false also drops ring:

# `ring` (or `rustls-aws-lc-rs`) is required — without it the build fails.
cargo add krafka --no-default-features --features rustls-aws-lc-rs

Re-exports§

pub use consumer::ConsumerRecord;
pub use consumer::TimestampType;
pub use error::KrafkaError;
pub use error::ProtocolErrorKind;
pub use error::Result;
pub use producer::Record;

Modules§

admin
Admin client for Apache Kafka.
auth
Authentication for Kafka connections.
consumer
Kafka consumer.
dlq
Dead-letter helper for consumer-side poison pills.
error
Error types for Krafka.
interceptor
Interceptor hooks for producers and consumers.
metrics
Client metrics.
prelude
The types you need to write a producer, a consumer or an admin client.
producer
Kafka producer implementation.
serdes
Typed serialization for the producer, pluggable byte transforms for the consumer.
share_consumer
Share consumer (KIP-932).
testingtest-broker
In-process fake Kafka broker for deterministic client tests.

Structs§

BrokerInfo
Information about a broker.
CloseOptions
How a client’s close_with shuts down.
Kafka
A connection to a Kafka cluster that every client is built from.
KafkaBuilder
Builder for Kafka: every connection setting.
PartitionInfo
Information about a topic partition.
ProxyConfig
A SOCKS5 proxy every broker connection of a Kafka handle is tunnelled through.
TopicInfo
Information about a topic.

Enums§

Compression
Compression codec.
GroupMembershipOperation
What a closing Consumer does about its group membership (KIP-1092), set with CloseOptions::group_membership_operation.
MetadataRecoveryStrategy
Strategy for recovering when the client loses the cluster, i.e. Java’s metadata.recovery.strategy (KIP-899, KIP-1102).

Constants§

OTEL_SEMCONV_VERSION
The OpenTelemetry semantic-conventions version krafka’s span names and attributes follow. The messaging conventions are still marked Development; krafka moves to a newer version deliberately, as a breaking change.

Type Aliases§

BrokerId
Kafka broker ID.
Headers
Record headers, the same type on produce and consume: keys in order, duplicates kept, and None for a null value (distinct from an empty one).
Offset
Kafka offset.
PartitionId
Kafka partition ID.
Timestamp
Kafka timestamp (milliseconds since epoch).