Upgrading to 0.27

Every breaking change in krafka 0.27 with the old API, the new API and what to do: the Kafka handle, send/enqueue, recv, metrics.

This page maps the 0.26 API to 0.27, area by area. The complete list is the release's Breaking and Removed sections in CHANGELOG.md. Start with construction, then each client you use, then metrics and tests.

Building clients

Connection settings are on KafkaBuilder; each client is built from the handle and shares its connection pool.

0.260.27
Producer::builder().bootstrap_servers(b)Kafka::builder(b).connect().await? then kafka.producer()
Consumer::builder().bootstrap_servers(b).group_id(g)kafka.consumer(g)
a consumer without a groupkafka.consumer_without_group()
ShareConsumer::builder().group_id(g)kafka.share_consumer(g)
AdminClient::builder()…build()kafka.admin()
KrafkaClient, with_client(&client), owns_poolthe Kafka handle; a second pool is a second Kafka
client_id, auth/sasl_*, proxy, transport/TransportConfig, request_timeout, connect_timeout, metadata settings on a role builderthe same setting on KafkaBuilder (security replaces auth)
refresh_tls, update_seed_brokers, rebootstrap on a clientthe same method on Kafka
Producer::close() returning (); close_with_timeout(d)close() returns Result<()>; close_with(CloseOptions::new().timeout(d))
Arc<dyn Interceptor> and friendspass by value (impl Trait); an Arc<T> still works
use krafka::Kafka;
use krafka::auth::{AuthConfig, TlsConfig};
use std::time::Duration;

let kafka = Kafka::builder("broker:9093")
    .client_id("payments")
    .security(AuthConfig::sasl_scram_sha512("user", "secret").with_tls(TlsConfig::new()))
    .request_timeout(Duration::from_secs(30))
    .connect()
    .await?;

let producer = kafka.producer().build().await?;
let consumer = kafka.consumer("payments-group").build().await?;
let admin = kafka.admin();

producer.close().await?;
consumer.close().await?;
admin.close().await?;

protocol, network and metadata are private modules. BrokerInfo, PartitionInfo, TopicInfo, MetadataRecoveryStrategy, ProxyConfig, TimestampType and Compression are re-exported at the crate root.

Cargo features

0.260.27
gzip, snappy, lz4, share consumer, SOCKS5, telemetry featuresalways compiled in; remove them from features = [...]
zstd needed the zstd feature to readzstd decodes in every build; zstd is needed only to encode
FIPS wording on rustls-aws-lc-rsno FIPS claim; the feature offers post-quantum X25519MLKEM768 first

The features are ring (default), rustls-aws-lc-rs, zstd, aws-msk, oauth-oidc, native-tls-roots, tls-encrypted-keys, unstable-protocol and test-broker.

Errors

0.260.27
KrafkaError::InvalidStateIllegalState
KrafkaError::Httpan OIDC failure is Auth with a source
fencing reported as a broker codeFenced
abort-required reported as a messageTransactionAbortable; check requires_abort()
no committed offset with auto_offset_reset = NoneNoOffset
max_connections reached: Configretriable Network
TLS handshake reset or EOF: Authretriable Network

New kinds: Closed, Wakeup, UnknownTopic, DeliveryTimeout { possibly_written } and OutOfOrderSequence. Branch on is_retriable(), is_fatal() and requires_abort() rather than on kinds where you can. See Error Handling.

Producer

0.260.27
send("t", key, value), send_with_headers, send_record(rec)send(Record::new("t", value).key(k).header(..))
ProducerRecordRecord (key, header, null_header, partition, timestamp)
RecordHeaderskrafka::Headers (Vec<(String, Option<Bytes>)>)
retries(n)removed: retries run until delivery_timeout
linger default 0 on Producer5 ms on every producer
buffer_memory(0) = unboundedrejected; pass a size
key_serializer/value_serializer, async SerializerTypedProducer<K, V> with the synchronous serdes::Serializer<T>
dead_letter_queue, DeadLetterQueue, KafkaDeadLetterQueueremoved from the producer; build a dead-letter record from a consumed one with dlq::record_for
DefaultPartitioner, StickyPartitioner, HashPartitioner, UniformStickyPartitionerthe built-in partitioner (no type to name); implement Partitioner for custom routing
Partitioner::on_new_batchremoved
state_store, ProducerStateStoreremoved
metrics_handle(), connection_metrics()metrics()
on_acknowledgement(.., DeliveryConfirmation::Failed ..)on_acknowledgement(topic, partition, result: Result<&RecordMetadata, &KrafkaError>, headers, ctx)

delivery_timeout must be at least linger + request_timeout, or build() fails. A partitioner that returns a partition outside [0, partition_count) fails that send with Config.

Transactions

0.260.27
TransactionalProducer::builder()…transactional_id(id).build() then init_transactions()kafka.producer().….build_transactional(id).await? (registers the id)
init_transactions_keeping_prepared()build_transactional, then prepared_transaction()
begin_transaction, commit_transaction, abort_transactionbegin, commit, abort
prepare_transaction, complete_transactionprepare, complete
send_record(rec)send(rec)
send_offsets_to_transactionsend_offsets(&offsets, &group_metadata)
TransactionalDeliveryHandleDeliveryHandle from enqueue

A commit after any failed send of the transaction refuses with TransactionAbortable, whether or not you awaited that send's handle; call abort(). abort() fails the transaction's buffered records instead of sending them.

use krafka::{Kafka, Record};

let kafka = Kafka::builder("localhost:9092").connect().await?;
let producer = kafka.producer().build_transactional("orders-tx").await?;

producer.begin()?;
producer.send(Record::new("orders", "o-1").key("c-7")).await?;
match producer.commit().await {
    Ok(()) => {}
    Err(e) if e.requires_abort() => producer.abort().await?,
    Err(e) => return Err(e),
}

Consumer

0.260.27
recv() -> Result<ConsumerRecord>, RecvErrorrecv() -> Result<Option<ConsumerRecord>>; Ok(None) means closed
batch_recv(n, timeout), BatchRecvOutcomepoll(timeout)
subscribe(&["a", "b"])subscribe(["a", "b"]), any iterator of strings
commit_sync()commit()
commit_async(), OffsetCommitHandlespawn commit() on a task, or commit less often
commit_with_metadata(..)commit_offsets(offsets) with OffsetAndMetadata::with_metadata; any iterator of pairs, &HashMap included
seek_many(&HashMap<(String, i32), i64>), initial_offsets(HashMap<(String, i32), i64>)an iterator of (TopicPartition, offset) pairs
ahash maps and sets in return typesstd::collections::HashMap / HashSet
current_lag(), is_caught_up(), fetch_end_offset(), cached_*lag() → HashMap<TopicPartition, PartitionLag>
revocation_timeout, max_cooperative_rebalance_roundsremoved; on_partitions_revoked is awaited to completion
PartitionAssignmentStrategy::StickyCooperativeSticky, Range or RoundRobin
PartitionAssignor trait, ConsumerGroup, GroupCoordinatorremoved; assignment strategies are the enum
AutoOffsetReset::to_offset()removed; AutoOffsetReset::ByDuration(d) is new
ConsumerRecord::topic: StringArc<str>
async Deserializersynchronous deserialize(&self, topic, headers, payload, is_key)

Behaviour to check in your code:

  • Position. A partition's position is the next record poll/recv will hand out. It moves only when records are returned; committing never moves it.
  • Seeking. seek and its variants fail with IllegalState for a partition that is not assigned. seek_to_beginning goes to the log start offset.
  • Subscribing. subscribe() returns before the group is joined; poll() applies the assignment and calls your listener.
  • Poll-interval expiry. The next poll() reports the partitions to on_partitions_lost once, rejoins and returns normally.
use krafka::Kafka;
use krafka::consumer::{OffsetAndMetadata, TopicPartition};

let kafka = Kafka::builder("localhost:9092").connect().await?;
let consumer = kafka.consumer("billing").build().await?;
consumer.subscribe(["invoices"]).await?;

while let Some(record) = consumer.recv().await? {
    let offsets = [(
        TopicPartition::new(record.topic.as_ref(), record.partition),
        OffsetAndMetadata::with_metadata(record.offset + 1, "billing-v2"),
    )];
    consumer.commit_offsets(offsets).await?;
}

Share consumer

0.260.27
acknowledge(&rec, AcknowledgeType::Accept).awaitack(&rec) (synchronous)
AcknowledgeType::Release / Rejectrelease(&rec) / reject(&rec)
no renewalrenew(&rec) (Kafka 4.2+ with KIP-1222)
acknowledge_by_offsetremoved; acknowledge the record
commit_sync(), commit_sync_with_timeout, commit_async, ShareCommitHandlecommit() → CommitResults, one Result per partition
close_with_timeout(d)close_with(CloseOptions::new().timeout(d))
max_buffered_records, max_records, session_timeout, heartbeat_intervalremoved; max_poll_records sets MaxRecords on every fetch

Implicit mode accepts the previous delivery when the next poll()/recv() starts, and on close(). See Share Consumer.

Admin client

Every operation is one method taking krafka types and an *Options struct (Default, with a timeout). Multi-item operations return HashMap<Item, Result<T, KrafkaError>>.

0.260.27
describe_consumer_group_offsetslist_consumer_group_offsets (require_stable replaces OffsetVisibility)
alter_topic_configincremental_alter_configs with ConfigResource and ConfigOp
describe_configs(topic names), describe_configs_per_resource, topic_configdescribe_configs over ConfigResources
describe_topic_partitions, describe_topic, partition_countdescribe_topics
alter_partition_reassignments_opts, abort_transaction_with_epochan option of alter_partition_reassignments / abort_transaction
list_client_metrics_resourceslist_config_resources
GroupListingListConsumerGroupsOptions
retries on the admin builderremoved; calls are bounded by default_api_timeout (60 s) or the options' timeout
write_txn_markers, get_controller_connection, poolremoved

list_consumer_groups and list_transactions return a result per broker. default_api_timeout and retry_backoff are set on the client kafka.admin() returns. See Admin Client.

Security

0.260.27
AuthConfig::sasl_plain_ssl(u, p, tls) and the other _ssl constructorsAuthConfig::sasl_plain(u, p).with_tls(tls)
AuthConfig::sasl_plain(..)? (fallible)infallible; bad credentials fail at connect()
with_scram_channel_binding, ChannelBindingremoved; SCRAM sends n,,
OAuthBearerTokenProvider, AwsMskIamCredentialProviderauth::CredentialProvider<C>; an async closure still works
AssertionSource variantsAssertionSource::file, fixed, provider
a key passphrase without a client certificate was ignoredTlsConfig build fails

See Authentication.

Metrics and tracing

0.260.27
KrafkaMetrics, per-client snapshot types, metrics_handle(), connection_metrics()metrics() → one owned krafka::metrics::Metrics with producer, consumer, connections; kafka.metrics() sums the handle
PrometheusExporter, JsonExporter, MetricsExporterMetrics::prometheus_text(); Metrics is plain data for any other format
LatencySnapshot (min, percentiles)Latency { count, sum, max }
Prometheus nameskrafka_producer_*, krafka_consumer_*, krafka_connections_*, each with a client_id label; update dashboards
tracing_ext modulespans are emitted by the clients; see Metrics
KIP-714 telemetry offon by default for producers and consumers; metrics_push(false) turns it off

The fake broker

krafka::testing is outside semver.

0.260.27
TransactionStateBrokerTransaction (status: TxnStatus, is_open())
ApiKey from the protocol modulekrafka::testing::ApiKey
RecordedRequestgains an at field
lenient sequence, transaction and share-session checksenforced as Kafka does; a request Kafka rejects is rejected

See Testing.