Expand description
§Krafka
A pure Rust, async-native Apache Kafka client.
Krafka provides high-performance, safe, and idiomatic Rust APIs for producing and consuming messages from Apache Kafka clusters.
§Features
- Pure Rust by default: No librdkafka or C bindings; the optional
zstdcompression feature links againstzstd-sysand requires a C toolchain - Async-native: Built on Tokio for non-blocking I/O
- High-performance: Zero-copy buffers, minimal allocations
- Safe: No unsafe code by default
- Cloud-native: First-class AWS MSK support including IAM auth
§Thread Safety
All main types in Krafka implement Send + Sync:
Producer- can be shared across tasks withArcConsumer- can be shared across tasks withArcAdminClient- can be shared across tasks withArc
This allows safe concurrent access from multiple Tokio tasks:
use std::sync::Arc;
use krafka::producer::Producer;
let producer = Arc::new(Producer::builder()
.bootstrap_servers("localhost:9092")
.build()
.await?);
// Spawn multiple tasks sharing the producer
for i in 0..10 {
let producer = producer.clone();
tokio::spawn(async move {
let _ = producer.send("topic", None, Some(b"message")).await;
});
}§Quick Start
§Producer
use krafka::producer::Producer;
let producer = Producer::builder()
.bootstrap_servers("localhost:9092")
.build()
.await?;
producer.send("my-topic", Some(b"key"), Some(b"value")).await?;
// A `None` value is a tombstone: on a compacted topic it deletes the key.
producer.send("my-topic", Some(b"key"), None).await?;§Consumer
use krafka::consumer::Consumer;
let consumer = Consumer::builder()
.bootstrap_servers("localhost:9092")
.group_id("my-group")
.build()
.await?;
consumer.subscribe(&["my-topic"]).await?;
loop {
match consumer.recv().await {
Ok(msg) => println!("{:?}", msg),
Err(krafka::RecvError::Closed) => break,
Err(krafka::RecvError::Error(e)) => return Err(e),
Err(_) => break,
}
}§Cargo Features
| Feature | Default | Description |
|---|---|---|
compression | yes | Enables pure-Rust compression codecs (gzip + snappy + lz4). |
compression-all | no | Enables all compression codecs, including zstd. |
gzip | via compression | Gzip record batch compression via flate2. |
snappy | via compression | Snappy compression via snap. |
lz4 | via compression | LZ4 compression via lz4_flex. |
zstd | no | Zstd compression via zstd (requires C toolchain). |
aws-msk | no | AWS MSK IAM authentication with SDK credential chain. |
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. |
socks5 | no | SOCKS5 proxy support via tokio-socks. |
telemetry | no | OpenTelemetry exporter for producer/consumer metrics. |
share-groups | yes | KIP-932 share consumer: queue semantics on a Kafka topic. GA in Apache Kafka 4.2; needs a 4.2+ broker. Same semver promise as the rest of the 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. |
ring | yes | rustls crypto backend using ring (pure Rust). |
rustls-aws-lc-rs | no | rustls crypto backend using aws-lc-rs. Preferred on AWS Graviton and for FIPS deployments. |
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. |
test-broker | no | In-process fake Kafka broker for testing your own code against a real client. Not for production builds. |
§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:
[dependencies]
# `ring` (or `rustls-aws-lc-rs`) is required — without it the build fails.
cargo add krafka --no-default-features --features lz4,ringRe-exports§
pub use error::KrafkaError;pub use error::ProtocolErrorKind;pub use error::RecvError;pub use error::Result;
Modules§
- admin
- Admin client for Apache Kafka.
- auth
- Authentication for Kafka connections.
- client
- Shared transport for connection pooling and metadata across multiple clients.
- consumer
- Kafka consumer implementation.
- dlq
- Dead-letter queue (DLQ) support for routing failed records to an error topic.
- error
- Error types for Krafka.
- interceptor
- Interceptor hooks for producers and consumers.
- metrics
- Metrics and observability for Krafka clients.
- prelude
- The types you need to write a producer, a consumer or an admin client.
- producer
- Kafka producer implementation.
- serdes
- Pluggable serialization applied on the way to and from the wire.
- share_
consumer share-groups - Share consumer implementation (KIP-932).
- telemetry
telemetry - KIP-714 client telemetry: subscription polling and metric push (feature-gated).
- testing
test-broker - In-process fake Kafka broker for deterministic client tests.
- tracing_
ext - Tracing extensions for observability.
- util
- Utility functions for Krafka.
Structs§
- Lazy
Record Batch - A lazily-decoded record batch for improved performance.
- Lazy
Record Iterator - Iterator that decodes records on demand from raw bytes.
- Record
- A Kafka record within a batch.
- Record
Batch - A Kafka record batch (v2 format).
- Record
Batch Builder - Builder for creating record batches.
- Record
Header - A Kafka record header.
- Sasl
Authenticator - SASL authenticator for handling authentication handshakes.
- Secure
Connection Config - Extended connection config with authentication.
- Secure
Connection Config Builder - Builder for SecureConnectionConfig.
Enums§
- Challenge
Response - Response from processing a SASL challenge.
- Compression
- Compression codec.
- Metadata
Recovery Strategy - Strategy for recovering when metadata refresh fails for too long.
Type Aliases§
- ApiVersion
- Kafka protocol API version.
- Broker
Id - Kafka broker ID.
- Correlation
Id - Kafka correlation ID for request/response matching.
- Offset
- Kafka offset.
- Partition
Id - Kafka partition ID.
- Timestamp
- Kafka timestamp (milliseconds since epoch).