๐ฆ krafka
A pure-Rust, async-native Apache Kafka client. No librdkafka, no C toolchain,
no unsafe, no panics โ enforced by the compiler, not by convention. Protocol
parity with Apache Kafka 4.3, checked in CI against Kafka's own schemas.
โจ Why krafka
Pure Rust, and it stays that way. No librdkafka, no C toolchain, no
cross-compilation surprises. The optional zstd feature is the single
exception, and it is opt-in.
The safety posture is enforced by the compiler. unsafe_code = "deny"
crate-wide, plus panic, unwrap and expect denied across the whole crate.
A malformed broker response cannot panic the process, and every allocation
from untrusted input is bounded twice โ by the declared count and by the
bytes actually available.
Protocol currency is a build failure, not a bug report. Every API version
is tracked against Apache Kafka 4.3 โ Fetch v18 (KIP-1166), Produce v13,
Metadata v13, DescribeLogDirs v5 (KIP-1066), DescribeQuorum v2
(KIP-836/853) โ and CI diffs krafka's version table against Kafka's own
message schemas. ApiVersions is negotiated rather than pinned, so a 3.9
broker gets 3.9-era versions with no configuration.
Correctness where it is hardest. KIP-320 truncation detection on all three
legs โ Fetch and ListOffsets and persisted through OffsetCommit, so it
survives restarts and rebalances. KIP-447 zombie fencing with a transaction
state machine that refuses the KAFKA-17754 abort-after-commit-timeout hazard.
OUT_OF_ORDER_SEQUENCE_NUMBER verified head-of-line before any rewind, so a
silent gap raises a fatal error instead of reporting success.
What is in the box
| Clients | Producer (batching, compression, idempotence, transactions) ยท Consumer (classic and KIP-848 server-side assignment with validated revoke-before-assign reconciliation) ยท ShareConsumer (KIP-932 at Kafka 4.2 parity, incl. KIP-1222 Renew and KIP-1206 ShareAcquireMode) ยท full AdminClient |
| Security | rustls TLS/mTLS with hot certificate reload (KIP-1288) ยท SASL PLAIN, SCRAM-SHA-256/512, OAUTHBEARER ยท built-in OIDC provider for client_credentials (KIP-768) and RFC 7523 client assertions (KIP-1258), with no cryptography dependency added ยท AWS MSK IAM ยท every mechanism composes with TLS through one with_tls, asserted reachable over both SASL_PLAINTEXT and SASL_SSL at compile time |
| Consistency | Every client shares one configuration surface and one operational surface (close, rebootstrap, update_seed_brokers, refresh_tls, metrics) โ asserted at compile time, builders included. One builder per client, with build_config() to validate without a broker and build() to validate and connect, both through the same validator. The transactional producer mirrors the plain one setter for setter, minus the two settings transactions fix (acks, idempotent) |
| Tuning | Per-codec compression levels (Gzip 0โ9, Zstd through 22) validated against the selected codec at build time, so a level set on a codec that has none is rejected rather than ignored โ on both send paths and on the transactional producer |
| Transport | One TransportConfig on every builder โ pass the same instance to every client that shares a network path: TCP keepalive, response ceiling, in-flight cap, idle eviction, file-descriptor cap ยท KIP-227 incremental fetch sessions ยท SOCKS5 |
| Observability | Lock-free counters, gauges and latency histograms with bounded per-topic cardinality ยท Prometheus export ยท producer interceptors and dead-letter queues on both producers and both send paths ยท OAUTHBEARER token-fetch counters and expiry gauge ยท OpenTelemetry semantic conventions |
| Hardening | Secret zeroization ยท constant-time comparison (subtle) ยท decompression-bomb limits ยท decode-loop bounds ยท RFC 3986 path encoding on every outbound HTTP target ยท CI forbids any credential-bearing type from deriving Debug |
| Testing | 2 440+ tests ยท 6 cargo-fuzz targets ยท proptest round-trips across the protocol layer ยท an in-process fake broker with fault injection that serves the full transaction protocol โ KIP-360 fencing, commit/abort markers, read_committed isolation, TV1 and KIP-890 TV2 โ so even exactly-once tests need no Docker |
Broker versions: krafka requires Apache Kafka 3.9+; protocol versions below that baseline have been removed. Features needing a newer broker (KIP-848 consumer groups, KIP-932 share groups, KIP-1066 cordoned log dirs) say so where they are documented and fail with a clear
UnknownApiVersionrather than silently degrading.
๐ Quick Start
Add krafka to your Cargo.toml:
[]
= "0.17.0"
= { = "1", = ["full"] }
# For AWS MSK IAM authentication with full SDK support:
# krafka = { version = "0.17.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_SSL + SCRAM-SHA-512 (the usual managed-Kafka listener)
use ;
let producer = builder
.bootstrap_servers
.auth
.build
.await?;
// Any mechanism composes with TLS through `with_tls`
let auth = sasl_scram_sha256
.with_tls;
// 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
let auth = aws_msk_iam;
let admin = builder
.bootstrap_servers
.auth
.build
.await?;
๐ฆ Modules
| Module | Description |
|---|---|
producer |
Batching, compression, idempotence, transactions, partitioners |
consumer |
Consumer groups (classic + KIP-848), offsets, rebalancing, compacted-topic tables |
share_consumer |
KIP-932 share groups โ queue semantics on a Kafka topic (unstable-protocol) |
admin |
Cluster administration: topics, partitions, groups, configs, ACLs, quotas, tokens |
client |
KrafkaClient โ one connection pool and metadata cache shared by several clients |
auth |
SASL PLAIN / SCRAM / OAUTHBEARER / AWS MSK IAM, TLS and mTLS |
serdes |
Serializer / Deserializer hooks applied on the way to and from the wire |
interceptor |
Producer and consumer hooks for tracing and enrichment |
dlq |
Dead-letter queues for records that exhaust their retries |
metrics |
Lock-free counters, gauges and latency histograms; Prometheus export |
telemetry |
KIP-714 broker-driven client telemetry and OTLP export (telemetry) |
tracing_ext |
OpenTelemetry semantic-convention fields for tracing spans |
testing |
In-process fake broker with fault injection (test-broker) |
error |
KrafkaError, ErrorCode, ProtocolErrorKind and retriability classification |
util |
Backoff policy, varint codecs, CRC32C, bootstrap-server parsing |
prelude |
One glob import for the common types (use krafka::prelude::*) |
Three more modules are public but #[doc(hidden)] โ protocol, network and
metadata. They are reachable for advanced use (custom authenticators, raw
record batches, benchmarks) but are not part of the stable API surface.
๐๏ธ 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.17.0", = false, = ["lz4", "snappy", "ring"] }
# Option 2: enable all compression codecs, including zstd
# krafka = { version = "0.17.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.17.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,
protocol-parity, secret-debug, test, test-ring,
test-cross-platform, minimal-features, doc, deny, integration, msrv.
Checks that exist because a review found what they now catch
tests/builder_surface.rs โ a compile-time assertion that every client
builder accepts a TransportConfig, offers a synchronous build_config()
alongside the async build(), exposes refresh_tls(), and keeps metrics and
version negotiation callable without an async context. Every line fails to
compile if the method it names disappears.
It also asserts two matrices that a per-client check could not see. Every
SaslMechanism must be constructible under both SASL_PLAINTEXT and
SASL_SSL from the public API alone โ SASL_SSL + SCRAM, the default secured
listener on most managed Kafka offerings, was unreachable from outside the crate
because the _ssl constructors were a hand-maintained list and SCRAM was missing
from it. And both producer builders must expose the same configuration surface โ
the transactional one was missing seventeen setters, including the
build_config() this file already promised for every client.
This replaced a Python parity script. krafka used to have two builders per
client โ 72 hand-maintained forwarding methods whose config half nothing outside
the crate's own tests ever called โ and the script checked that they stayed in
sync. It found two real defects, then missed a third: both producer builders had
a compression method, but only the unused one validated it, so
.compression(Zstd) without the zstd feature built a producer that failed on
its first send. A parity check compares surfaces; the divergence had moved
underneath it. Deleting the duplication removed both the defect class and the
need for the script, and Rust checks reachability better than a regex can.
just protocol-parity โ the API version table is diffed against Apache
Kafka's own message schemas: names and keys agree, MIN is still a version Kafka
accepts, MAX neither overstates (claiming a version marked
latestVersionUnstable) nor understates (declining a stable one), and the
flexible-version boundary matches. This is how Fetch v17/v18 sat implemented,
documented and unreachable for two Kafka releases.
It reads a vendored snapshot, so it needs no network and cannot flake. Track a newer Kafka release deliberately:
just secret-debug โ no credential-bearing type may derive Debug.
Debug is the quiet way secrets reach a log aggregator: a tracing field, an
error context or a panic message that formats the enclosing struct is enough,
and nobody has to log the secret deliberately. Two instances shipped before this
check existed โ the OIDC client secret, and SaslAuthenticateRequest.auth_bytes,
which for SASL/PLAIN is \0username\0password in cleartext.
โก 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?;
Transport tuning
Socket- and pool-level settings live on one TransportConfig, accepted by every
builder (Producer, Consumer, AdminClient, TransactionalProducer,
ShareConsumer, KrafkaClient). The defaults reproduce krafka's historical
behaviour exactly, so this is opt-in.
use TransportConfig;
use Consumer;
use Duration;
let transport = builder
// Beat the idle timeout of whatever NAT gateway or load balancer sits
// between you and the brokers โ the usual cause of "the consumer stops
// receiving after exactly N minutes".
.tcp_keepalive
// Kafka returns at least one full record batch per partition even when it
// exceeds fetch.max.bytes. Raise this above the topic's max.message.bytes
// or that partition stalls permanently.
.max_response_size
// Bound worst-case memory: the per-connection ceiling is
// max_response_size ร max_in_flight_requests.
.max_in_flight_requests
// Bound file descriptors on a cluster whose broker count can jump.
.max_connections
// Re-read certificates from disk hourly (KIP-1288).
.tls_reload_interval
.build?;
let consumer = builder
.bootstrap_servers
.group_id
.transport
.build
.await?;
TLS certificate rotation (KIP-1288)
Two paths, because rotation happens two ways:
// Event-driven: an inotify watch or a sidecar signal fired.
producer.refresh_tls.await?;
// Unattended: set `tls_reload_interval` above and krafka reloads on a timer.
Existing TLS sessions keep the certificates they handshaked with and are replaced as connections cycle. A reload that fails โ a half-written PEM caught mid-rotation โ is logged and the previous material stays active, so a non-atomic rotation converges on the next attempt instead of breaking every new connection in between.
๐ฏ 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, but they are no longer equals.
Prefer GroupProtocol::Consumer (KIP-848). It has been production ready
since Apache Kafka 4.0: the coordinator computes assignments server-side, a
rebalance reconciles incrementally instead of stopping every member, and a slow
member affects only its own partitions. Apache Kafka 4.3 began deprecating
the classic protocol (KIP-1274 phase 1 โ warn in 4.3, default flips in 5.0,
removed in 6.0), and krafka logs the same warning once per process when a group
starts on it.
use ;
let consumer = builder
.bootstrap_servers
.group_id
.group_protocol // KIP-848
.build
.await?;
Classic remains the default for now, so that upgrading krafka is never itself
a protocol migration: krafka supports Kafka 3.9 brokers, and KIP-848 needs 4.0
(or 3.7โ3.9 with group.coordinator.new.enable=true). The two protocols cannot
mix within one group on pre-4.0 brokers โ move every member together, or
upgrade the cluster first.
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, share groups, and administration, tracked against Apache Kafka 4.3.
Broker-internal APIs (LeaderAndIsr, UpdateMetadata, Vote, FetchSnapshot,
BrokerHeartbeat, the share-group state persister, โฆ) are deliberately absent
โ a client does not speak them. They are still named in the ApiKey enum, so
an ApiVersions response from a modern broker decodes to something readable
rather than Unknown(87).
Not implemented: KIP-1071 Streams group protocol (keys 88โ89) and KIP-1258 OAuth client assertion.
Schema registries are out of scope, as they are for every comparable client
โ Java's kafka-clients has none, librdkafka has none, franz-go keeps pkg/sr
out of kgo. A registry is a different service with a different protocol, auth
model and release cadence. krafka provides the hook (serdes::Serializer /
Deserializer, the equivalent of Java's key.serializer); pair it with
schemreg for Confluent, AWS Glue or
Apicurio, and Avro / Protobuf / JSON codecs. See the
Cookbook.
Tokio is the async runtime.
Authentication
SASL/PLAIN, SASL/SCRAM-SHA-256/512 (with RFC 5929 channel binding), SASL/OAUTHBEARER
(with proactive token refresh, a built-in OIDC client_credentials provider and
KIP-1258 client assertions behind the oauth-oidc feature), 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, DelayThen, Disconnect, Silence,
CorruptRecords and pass-through. The cluster is mutable mid-test
(set_leader, bump_leader_epoch, set_group_coordinator, set_controller,
set_broker_online), and set_api_versions makes the broker advertise an
older API range so the client's degradation branches are reachable.
It serves the produce/fetch path, both group protocols โ classic and KIP-848 with real revoke-before-assign reconciliation โ KIP-932 share groups with the share-partition state machine, KIP-584 feature administration, and the full transaction protocol.
Transactions are modelled end to end, which is what makes exactly-once testable
without a cluster: InitProducerId returns a stable producer ID per
transactional ID with an epoch that rises on every re-initialisation (KIP-360
fencing), EndTxn writes real commit and abort control batches, offsets staged
by TxnOffsetCommit apply only on commit, and a read_committed fetch stops at
the last stable offset and reports aborted transactions so the consumer's own
filtering runs for real. set_transaction_version(2) finalizes the
transaction.version feature and the client negotiates KIP-890 TV2 from it โ
the same route a real cluster takes โ so both protocols are reachable.
broker.set_transaction_version;
// ...run a transaction through a TransactionalProducer...
assert_eq!;
See the Testing guide for what it does and, just as importantly, what it deliberately does not model.
โฌ๏ธ Upgrading
Release-by-release detail lives in CHANGELOG.md.
Upgrading to 0.17
Breaking โ schema registry moved out
krafka::schema_registry is gone, with the schema-registry and
aws-glue-schema-registry features. The registry client now lives in
schemreg, which additionally supports
Apicurio and ships Avro / Protobuf / JSON codecs krafka never had.
Every comparable client draws the line here โ Java's kafka-clients has no
registry support, librdkafka has none, franz-go keeps pkg/sr out of kgo. A
registry is a different service; coupling it to the protocol client meant a
registry API change could force a Kafka client release.
What krafka keeps is the hook, generalised:
SchemaEncoder/SchemaDecoderโkrafka::serdes::Serializer/Deserializer(encodeโserialize,decodeโdeserialize).key_encoder/value_encoderโkey_serializer/value_serializer;key_decoder/value_decoderโkey_deserializer/value_deserializer.KrafkaError::SchemaRegistryโKrafkaError::Http.
Since the traits are plain Bytes -> Bytes, they now cover encryption and
compression as well as schema framing. The ~20-line schemreg adapter is in the
Cookbook.
Breaking โ consumer offset accessors
Consumer::cached_end_offsetis isolation-aware. Underread_committedit returns the last stable offset rather than the high watermark, because the broker will not deliver a record at or above the LSO. Use the newcached_high_watermark()if you specifically want the log-end offset.read_uncommitted(the default) is unchanged.Consumer::position()reports the delivered offset, not the read-ahead. It is the value a commit writes, soposition()andcommit()cannot disagree. The read-ahead value is the newfetch_position().
Fixed
- A
seek()could move the committed offset backwards. Every reposition path left already-fetched records in the receive buffer, and a commit is clamped down to the lowest still-buffered offset โ correct on its own, and what stops an undelivered record from being acknowledged. Afterseek_to_end()on a partition with buffered offset 100 and a new position of 5 000, the next commit wrote 100. Viaauto.offset.resetthe clamped offset could be one the log no longer holds, producing a reset โOFFSET_OUT_OF_RANGEloop that never converged. read_committedreported permanent phantom lag.last_stable_offsetwas decoded from every fetch response and never read, solag(),is_caught_up()and thelagmetrics compared against the high watermark. An open transaction kept a fully drained consumer reporting a backlog it could never close, andis_caught_up()could never returntrue.pause()was bypassed byrecv()/batch_recv().poll()withheld paused partitions; the buffer drain did not, so the same client gave two answers depending on which read API was used. Withheld records are held, not discarded โ the fetch position has already advanced past them.- A commit marker could end a
read_committedabort filter early. The aborted-transaction filter deactivated on any control batch without reading the marker's type field, so aborted records could reach the application. - A transactional commit could orphan a record into the next transaction.
commit_transaction()drained the accumulator before closing the transaction to new records, so a concurrentsend()could slip in behind the flush and stay buffered until afterEndTxnโ landing in the following transaction, and vanishing if that one aborted. - A commit could write
EndTxnwhilesend_offsets_to_transactionwas still in flight, committing the consumer's offsets outside the transaction. The output records stayed atomic with each other but not with the position that produced them. assign()leaked state for partitions it dropped. Narrowing a manual assignment left the old partitions' positions, watermarks and buffered records behind โ and the stale buffer entry dragged back the commit for the partitions still being consumed.- A share-consumer flush could strand acknowledgements.
poll()holds the pending acks out of the map for the duration of itsShareFetch, so a concurrentcommit_sync()orclose()flushed an empty map and reported success. The documentedwakeup()โclose()shutdown hits exactly that window. Both flush paths now wait for in-flight polls first.
Changed
- Lag counts records read ahead into the buffer โ fetched is not delivered.
position(),lag(),current_lag(),is_caught_up()andcommit()are now all derived from one boundary.
Faster
- Fetch responses are read ahead into a prefetch buffer. A 50 MB response was fully decoded and then truncated to 500 records, with the surplus dropped and re-decoded next poll โ roughly 100ร the necessary decode work per poll on a 50-partition assignment. Each fetch now decodes one delivery's worth plus the buffer's free capacity and parks the surplus, so the next poll is served from memory with no Fetch on the wire: half the round trips, and network latency out of every other poll.
- Partition fetch order is a real round robin. Fairness previously depended
on unspecified
HashMapiteration order; partitions now rotate by one position per poll, matching the Java client'sPartitionStates.moveToEnd.
New
Consumer::cached_high_watermarkandConsumer::cached_last_stable_offsetโ the gap between them is the volume of in-flight transactional data.Consumer::fetch_positionโ where the next fetch starts, as opposed to where delivery is.KafkaDeadLetterQueueโ the DLQ implementation everyone was writing by hand. Attaches provenance headers, drops the source partition index, and counts what it could not save.krafka::preludeโ one glob import for the common types.krafka::interceptor::CommitOffsetsโ names the mapon_committakes, so implementors no longer needahashin their own manifest.
Documentation
just docs-testcompiles the guide snippets. It was referenced by the doc tooling for two releases without existing; 192 of 321 Rust blocks are now compile-checked in CI. It found broken examples in the README and Getting Started on its first run โ including the admin quick-start, which chained a method onto aResult.
Upgrading to 0.16
Breaking
AwsMskIamCredentials::with_session_tokenis now a builder method, not a four-argument constructor. Build withAwsMskIamCredentials::new(id, secret, region).with_session_token(token). The old form fails to compile rather than changing meaning.
Fixed โ three settings that silently did nothing
compression_levelwas dropped on the batching path. It applied only atlinger = 0, so the throughput-tuned configuration โ and everyTransactionalProducer, which always batches โ encoded at the codec's default. Now applied on both paths and on both producers.dead_letter_queuewas direct-send only. Configuring a DLQ alongside any batching silently disabled it. Now invoked on both paths, and on the transactional producer.close()tore down a shared connection pool. AProducer,ConsumerorTransactionalProducerbuilt with.with_client(..)calledclose_all()unconditionally, killing every sibling client's connections. Every client now reportsowns_pool()and leaves a borrowed pool to itsKrafkaClient.SecureConnectionConfigBuilder::tls()likewise lost its TLS configuration if called before a SASL setter; order no longer matters.
New
AuthConfig::with_tls(TlsConfig)โ every SASL mechanism composes with TLS through one method.SASL_SSL+ SCRAM, the default secured listener on most managed Kafka offerings, was previously unreachable from outside the crate.sasl_scram_sha256_ssl/sasl_scram_sha512_ssladded for symmetry;AuthConfig::from_envgainedKAFKA_SSL_*material and theOAUTHBEARERandAWS_MSK_IAMmechanisms.AwsMskIamCredentials::with_regionandfrom_env_with_regionโ change the region without losing the session token, and load keys from the environment with the region from your own configuration.TransactionalProducerBuilderreaches parity withProducerBuilderโbuild_config(),compression_level,topic_compression,delivery_timeout,dead_letter_queue,interceptor/add_interceptor,state_store,with_client, the metadata cache TTLs andsasl_oauthbearer_provider.acksandidempotentstay excluded because transactions fix both.TransactionalProducer::flush()โ so code generic over "a producer" need not special-case which one it holds.ShareConsumerBuilder::with_clientโ the one client that could not share aKrafkaClient's pool now can.- OAUTHBEARER token-lifecycle metrics โ
oauth_token_fetches,oauth_token_fetch_failures,oauth_token_fetch_latencyandoauth_token_expiry_epoch_msonConnectionMetrics, plus aWARNon every failed fetch. A misconfiguredtoken_endpointis no longer indistinguishable from an unreachable broker. - The fake broker serves the full transaction protocol โ KIP-360 fencing,
commit/abort control batches,
read_committedisolation, TV1 and KIP-890 TV2. Exactly-once is now testable without Docker.
๐ Documentation
Full documentation: hupe1980.github.io/krafka ยท API reference: docs.rs/krafka ยท Release history: CHANGELOG.md
| Start here | Clients | Integration | Operations | Reference |
|---|---|---|---|---|
| Getting Started | Producer | Authentication | Metrics | Protocol Support |
| Cookbook | Consumer | Performance | Architecture | |
| Configuration | Share Consumer | Interceptors | Testing | |
| Admin Client | Error Handling |
The site is built with Zola from site/. Run
just site-serve for a local preview with live reload.
๐ฎ 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.