Skip to main content

Crate krafka

Crate krafka 

Source
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 zstd compression feature links against zstd-sys and 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 with Arc
  • Consumer - can be shared across tasks with Arc
  • AdminClient - can be shared across tasks with Arc

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

FeatureDefaultDescription
compressionyesEnables pure-Rust compression codecs (gzip + snappy + lz4).
compression-allnoEnables all compression codecs, including zstd.
gzipvia compressionGzip record batch compression via flate2.
snappyvia compressionSnappy compression via snap.
lz4via compressionLZ4 compression via lz4_flex.
zstdnoZstd compression via zstd (requires C toolchain).
aws-msknoAWS MSK IAM authentication with SDK credential chain.
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.
socks5noSOCKS5 proxy support via tokio-socks.
telemetrynoOpenTelemetry exporter for producer/consumer metrics.
share-groupsyesKIP-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-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.
ringyesrustls crypto backend using ring (pure Rust).
rustls-aws-lc-rsnorustls crypto backend using aws-lc-rs. Preferred on AWS Graviton and for FIPS deployments.
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.
test-brokernoIn-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,ring

Re-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_consumershare-groups
Share consumer implementation (KIP-932).
telemetrytelemetry
KIP-714 client telemetry: subscription polling and metric push (feature-gated).
testingtest-broker
In-process fake Kafka broker for deterministic client tests.
tracing_ext
Tracing extensions for observability.
util
Utility functions for Krafka.

Structs§

LazyRecordBatch
A lazily-decoded record batch for improved performance.
LazyRecordIterator
Iterator that decodes records on demand from raw bytes.
Record
A Kafka record within a batch.
RecordBatch
A Kafka record batch (v2 format).
RecordBatchBuilder
Builder for creating record batches.
RecordHeader
A Kafka record header.
SaslAuthenticator
SASL authenticator for handling authentication handshakes.
SecureConnectionConfig
Extended connection config with authentication.
SecureConnectionConfigBuilder
Builder for SecureConnectionConfig.

Enums§

ChallengeResponse
Response from processing a SASL challenge.
Compression
Compression codec.
MetadataRecoveryStrategy
Strategy for recovering when metadata refresh fails for too long.

Type Aliases§

ApiVersion
Kafka protocol API version.
BrokerId
Kafka broker ID.
CorrelationId
Kafka correlation ID for request/response matching.
Offset
Kafka offset.
PartitionId
Kafka partition ID.
Timestamp
Kafka timestamp (milliseconds since epoch).