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
ccfinds forring, nothing else. Every codec decodes in pure Rust; the opt-inzstd,rustls-aws-lc-rsandaws-mskfeatures 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
unsafecode;panic,unwrapandexpectare 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 ownedMetricssnapshot, including the connection counters of its pool;Kafka::metricssums a handle’s clients.Metrics::prometheus_textrenders one for a scrape endpoint. - Spans. The clients emit spans through
tracing, on OpenTelemetry messaging semantic conventionsOTEL_SEMCONV_VERSION(v1.44.0, whose messaging conventions are still Development): asendspan per record fromsend/enqueueto its outcome (kindproducer), apollspan perpoll/recvand acommitspan per commit (kindclient), and a krafka-specificrebalancespan per assignment change (kindinternal). Record keys and values, and hashes of keys, are never recorded; a key’s size is. krafka depends on no OpenTelemetry crate: bridgetracingto 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 bymetrics_pushon 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
| Feature | Default | Description |
|---|---|---|
ring | yes | rustls crypto backend using ring. |
rustls-aws-lc-rs | no | rustls crypto backend using aws-lc-rs; offers post-quantum X25519MLKEM768 key exchange first. Compiles C and needs CMake. |
zstd | no | Zstd encoding via zstd (compiles C). Zstd decoding is always available. |
aws-msk | no | AWS MSK IAM authentication with the SDK credential chain (compiles C and needs CMake, via aws-lc-sys). |
oauth-oidc | no | Built-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-roots | no | Load platform-native root certificates via rustls-native-certs. |
tls-encrypted-keys | no | Passphrase-encrypted PKCS#8 client keys (ssl.key.password) via the RustCrypto pkcs8 crate. |
unstable-protocol | no | Enables 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-broker | no | In-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-rsRe-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).
- testing
test-broker - In-process fake Kafka broker for deterministic client tests.
Structs§
- Broker
Info - Information about a broker.
- Close
Options - How a client’s
close_withshuts down. - Kafka
- A connection to a Kafka cluster that every client is built from.
- Kafka
Builder - Builder for
Kafka: every connection setting. - Partition
Info - Information about a topic partition.
- Proxy
Config - A SOCKS5 proxy every broker connection of a
Kafkahandle is tunnelled through. - Topic
Info - Information about a topic.
Enums§
- Compression
- Compression codec.
- Group
Membership Operation - What a closing
Consumerdoes about its group membership (KIP-1092), set withCloseOptions::group_membership_operation. - Metadata
Recovery Strategy - 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§
- Broker
Id - Kafka broker ID.
- Headers
- Record headers, the same type on produce and consume: keys in order,
duplicates kept, and
Nonefor a null value (distinct from an empty one). - Offset
- Kafka offset.
- Partition
Id - Kafka partition ID.
- Timestamp
- Kafka timestamp (milliseconds since epoch).