๐ฆ Krafka
A pure Rust, async-native Apache Kafka client designed for high performance, safety, and ease of use.
โจ Features
- ๐ฆ Pure Rust by default: No librdkafka or C dependencies; the optional
zstdcompression feature requires a C toolchain viazstd-sys - โก Async-native: Built on Tokio for true async I/O
- ๐ Zero unsafe: Safe Rust by default
- ๐ High performance: Zero-copy buffers, inline hot paths, efficient batching, concurrent batch flushing
- ๐ฆ Full protocol support: Kafka protocol with all compression codecs
- ๐ค Real version negotiation:
ApiVersionsis negotiated, not pinned โ KIP-511 client identity reaches the broker and KIP-584 cluster feature levels are cached per connection - ๐ Incremental fetch sessions: KIP-227 fetch sessions for bandwidth-efficient multi-partition consumers
- ๐งฉ KIP-848 consumer groups: server-side assignment with validated reconciliation โ revoke-before-assign, epoch fencing, and no partition ever owned by two members at once
- ๐ TLS/SSL encryption: Using rustls for secure connections
- ๐ SASL authentication: PLAIN, SCRAM-SHA-256/512, OAUTHBEARER mechanisms
- ๐ฏ Transactions: Exactly-once semantics with KIP-447 zombie fencing, and a state machine that refuses the KAFKA-17754 abort-after-commit-timeout hazard
- ๐งญ Truncation detection (KIP-320): leader epochs are sent on Fetch and ListOffsets and persisted through OffsetCommit, so the check survives restarts and rebalances
- โ๏ธ Cloud-native: First-class AWS MSK support including IAM auth
- ๐ก๏ธ Security hardened: Secret zeroization, constant-time auth (
subtle), decompression bomb protection, decode loop bounds (MAX_DECODE_ARRAY_LEN), RFC 3986 path encoding on every outbound HTTP target - ๐ Built-in retry: Exponential backoff with metadata refresh on leader changes
- ๐ Metrics: Lock-free counters/gauges/latency with bounded per-topic cardinality
- ๐งช Fuzz + property tested: 6 cargo-fuzz targets and proptest round-trips across the protocol layer
Minimum Broker Version: Krafka requires Apache Kafka 3.9+. Protocol versions older than the Kafka 3.9 baseline have been removed.
๐ Quick Start
Add Krafka to your Cargo.toml:
[]
= "0.14.0"
= { = "1", = ["full"] }
# For AWS MSK IAM authentication with full SDK support:
# krafka = { version = "0.14.0", features = ["aws-msk"] }
Producer
use Producer;
use Result;
async
Consumer
use ;
use Result;
use Duration;
async
Admin Client
use ;
use Result;
use Duration;
async
Transactional Producer
For exactly-once semantics across multiple partitions:
use TransactionalProducer;
use Result;
async
Authentication
Connect to secured Kafka clusters with SASL, SCRAM, OAUTHBEARER, or AWS MSK IAM โ available on all client types:
use Producer;
use Consumer;
use AdminClient;
// Producer with SASL/SCRAM-SHA-256
let producer = builder
.bootstrap_servers
.sasl_scram_sha256
.build
.await?;
// Consumer with SASL/PLAIN
let consumer = builder
.bootstrap_servers
.group_id
.sasl_plain
.build
.await?;
// Producer with SASL/OAUTHBEARER
let producer = builder
.bootstrap_servers
.sasl_oauthbearer
.build
.await?;
// Admin with AWS MSK IAM
use AuthConfig;
let auth = aws_msk_iam;
let admin = builder
.bootstrap_servers
.auth
.build
.await?;
๐ฆ Modules
| Module | Description |
|---|---|
producer |
High-throughput message production with batching and compression |
consumer |
Consumer groups with rebalancing, offset management, and static membership |
admin |
Cluster administration (topics, groups, records, configuration, ACLs) |
interceptor |
Producer and consumer interceptor hooks for observability |
protocol |
Kafka wire protocol implementation |
auth |
Authentication (SASL/PLAIN, SASL/SCRAM, SASL/OAUTHBEARER, AWS MSK IAM) |
๐๏ธ Compression
Krafka supports all Kafka compression codecs, individually feature-gated:
use Producer;
use Compression;
let producer = builder
.bootstrap_servers
.compression // Fast compression
.build
.await?;
| Codec | Cargo Feature | Crate | Characteristics |
|---|---|---|---|
Compression::Gzip |
gzip |
flate2 | Best ratio, slower |
Compression::Snappy |
snappy |
snap | Good balance |
Compression::Lz4 |
lz4 |
lz4_flex | Fastest |
Compression::Zstd |
zstd |
zstd | Best modern choice (requires C toolchain) |
The default compression feature enables the pure-Rust codecs: gzip, snappy,
and LZ4. Zstd remains available through the explicit zstd or
compression-all feature because it requires a C toolchain via zstd-sys.
To select only what you need:
# Option 1: enable only the codecs you need
# `default-features = false` also drops the default `ring` TLS backend, so a
# crypto backend must be named explicitly.
= { = "0.14.0", = false, = ["lz4", "snappy", "ring"] }
# Option 2: enable all compression codecs, including zstd
# krafka = { version = "0.14.0", features = ["compression-all"] }
TLS crypto backend
Krafka uses rustls and needs exactly one crypto backend. ring is the
default; rustls-aws-lc-rs selects aws-lc-rs instead, which is the better
choice on AWS Graviton and in FIPS-oriented deployments:
= { = "0.14.0", = false, = ["rustls-aws-lc-rs", "compression"] }
The two backends are additive, not mutually exclusive โ a transitive
dependency may well enable the other one, and Cargo would have no way to
resolve a conflict if they were exclusive. When both are compiled in,
aws-lc-rs deterministically wins. Krafka always selects the provider
explicitly rather than letting rustls infer it from crate features, so the
combination cannot produce a runtime panic. Installing a process-wide provider
with CryptoProvider::install_default() overrides the choice for the whole
application, Krafka included.
๐ ๏ธ Development
Tasks are driven by just. The justfile is the single
source of truth for what the checks are โ CI calls the same recipes, so a check
cannot pass locally and fail in CI because the two drifted apart.
Individual recipes mirror one CI job each: fmt-check, clippy, check,
api-parity, test, test-ring, test-cross-platform, minimal-features,
doc, deny, integration, msrv.
just api-parity is worth calling out: it checks mechanically that every option
on an internal *ConfigBuilder is reachable from the public *Builder that
Client::builder() returns. An option present on one and not the other is
implemented, documentable and uncallable โ which is how KIP-848 selection and
the producer's dead-letter queue both shipped unreachable.
โก Performance Tuning
High Throughput Producer
use ;
use Compression;
use Duration;
let producer = builder
.bootstrap_servers
.acks
.compression
.batch_size // 1MB batches
.linger // Allow batching
.build
.await?;
Low Latency Consumer
use Consumer;
use Duration;
let consumer = builder
.bootstrap_servers
.group_id
.fetch_min_bytes
.fetch_max_wait
.build
.await?;
๐ฏ Delivery Semantics
Pick the mode that matches what your data is worth.
| Mode | Guarantee | Cost |
|---|---|---|
acks=0 |
At-most-once. No durability โ the record may never reach the log. | Lowest latency |
acks=all + idempotence (default) |
At-least-once, ordered and gap-free per partition. Batches for a partition are serialised in seal order; sequence numbers are monotonic and never reused. | One round trip to the ISR |
Transactions + send_offsets_to_transaction |
Exactly-once across a read-process-write cycle. | Transaction coordinator round trips |
Under acks=0, RecordMetadata::confirmation reports Unacknowledged โ do not
mistake the returned metadata for a durability guarantee.
Exactly-once
send_offsets_to_transaction takes a [ConsumerGroupMetadata], not a bare group
ID. That metadata is what lets the group coordinator fence a zombie: an
instance that was partitioned away, lost its partitions to a rebalance, and came
back still holding a transaction. Without it the coordinator accepts the zombie's
commit and it overwrites the position of the member that now owns the partition.
// Re-read for every transaction. The generation changes on every rebalance,
// so a cached value stops fencing at exactly the moment it matters.
let group_metadata = consumer
.group_metadata
.await
.ok_or?;
producer.begin_transaction?;
producer.send.await?;
producer.send_offsets_to_transaction.await?;
producer.commit_transaction.await?;
See examples/exactly_once.rs for the full
read-process-write loop.
๐งฉ Consumer Groups
Both rebalance protocols are supported: the classic JoinGroup/SyncGroup protocol
and KIP-848 (group.protocol = consumer).
Four assignors ship: Range, RoundRobin, Sticky (eager), and
CooperativeSticky. The default is the preference list
[Range, CooperativeSticky], matching the Java client โ every member advertises
both, so a group moves from eager to cooperative rebalancing in a single rolling
bounce rather than a full stop-the-world restart.
Rebalances do not wait on your poll loop. The background group task keeps
heartbeating through a rebalance and sends JoinGroup/SyncGroup itself, so a
consumer that is idle or busy between poll() calls does not hold the rest of
the group up. The new assignment is applied โ and your rebalance listener
called โ on the next poll(), so callbacks and record delivery stay on one
thread and an offset commit cannot race a revocation.
max.poll.interval.ms is still enforced: an application that genuinely stops
polling leaves the group so its partitions are reassigned promptly, rather than
holding them while a background task vouches for it. Static members
(group.instance.id) instead keep their assignment until the session expires,
so a restart can reclaim it.
๐ Scope
Krafka speaks the client side of the Kafka protocol: 60+ API keys covering
produce, fetch, group coordination, transactions, and administration.
Broker-internal APIs (LeaderAndIsr, UpdateMetadata, Vote, FetchSnapshot,
BrokerHeartbeat, โฆ) are deliberately absent โ a client does not speak them.
The schema-registry module implements the Confluent and AWS Glue wire formats (magic byte, schema ID framing, caching). Bring your own Avro/Protobuf/JSON-Schema codec; krafka does not impose one.
Tokio is the async runtime.
Authentication
SASL/PLAIN, SASL/SCRAM-SHA-256/512 (with RFC 5929 channel binding), SASL/OAUTHBEARER (with proactive token refresh), AWS MSK IAM, and mTLS.
GSSAPI/Kerberos is outside the scope of this client.
๐งช Testing Against a Fake Broker
Enable the test-broker feature to get an in-process Kafka broker your tests can
drive directly. Real Producer/Consumer/AdminClient instances connect to it
over a real TCP socket, so you exercise the actual client โ no Docker, no
containers, and failure modes you cannot reproduce against a healthy cluster.
use ;
use ApiKey;
use ErrorCode;
let broker = start.await?;
// Make the next CreateTopics land on a non-controller and assert the client
// refreshes metadata and retries instead of surfacing the error.
broker.on;
let admin = builder
.bootstrap_servers
.build
.await?;
Control covers Error, Delay, Disconnect, Silence and pass-through, and
the cluster can be manipulated mid-test: set_leader, bump_leader_epoch,
set_group_coordinator, set_txn_coordinator, set_controller,
set_broker_online. That makes leader moves, coordinator failover, late
responses and controller churn ordinary unit tests.
๐ Documentation
Full documentation is available at hupe1980.github.io/krafka
- Getting Started
- Producer Guide
- Consumer Guide
- Admin Client
- Configuration Reference
- Performance Tuning
- Architecture Overview
- Metrics & Observability
- Error Handling
- Interceptors
- Authentication
๐ฎ Examples
Run the examples with:
# Producer example
# Consumer example
# Advanced consumer example (pause/resume, seek, manual commits)
# Admin client example
# Transactional producer example
# Exactly-once read-process-write (KIP-447 zombie fencing)
# Authentication examples (SASL, SCRAM, MSK IAM)
๐ค Contributing
Contributions are welcome!
๐ License
Licensed under either the MIT License or the Apache License 2.0, at your option.