Are you an LLM? Read llms.txt for a summary of the docs, or llms-full.txt for the full context.
Skip to content

Apache Kafka

If you already have a Kafka cluster or need streaming semantics — replayable log, partition-key ordering, consumer-lag autoscaling — this backend fits naturally. Best for high-throughput pipelines where log retention and replay matter.

What you need

A Kafka cluster in KRaft mode or with Zookeeper. For local dev:

docker run --rm -p 9092:9092 confluentinc/cp-kafka:latest

The integration tests use testcontainers with the Apache Kafka module, so any runnable example also spins up a container automatically.

Install

cargo add shove --features kafka

For TLS and SASL (PLAIN, SCRAM-SHA-256, SCRAM-SHA-512). librdkafka implements PLAIN in every build. SCRAM and OAUTHBEARER come from its OpenSSL build, which is why shove keeps the whole KafkaSasl type behind this feature. OAUTHBEARER is compiled in, but shove exposes it only as KafkaSasl::MskIam under kafka-msk-iam (below). There is no generic token provider. For TLS, PLAIN and SCRAM:

cargo add shove --features kafka-ssl

For GSSAPI/Kerberos on an rdkafka client you build yourself in the same binary, which needs Cyrus SASL (libsasl2) linked into the build:

cargo add shove --features kafka-gssapi

This feature only links the library. KafkaSasl exposes no GSSAPI variant and KafkaConfig passes no raw client properties, so shove itself cannot authenticate with Kerberos. The feature gates no shove code and exists for direct rdkafka use.

For AWS MSK with IAM authentication:

cargo add shove --features kafka-msk-iam

Connect

    let broker = Broker::<Kafka>::new(KafkaConfig::new(&bootstrap)).await?;

KafkaConfig::new takes a bootstrap server address (e.g. localhost:9092). TLS and SASL are configured via KafkaTls and KafkaSasl on the config — available when the kafka-ssl feature is enabled.

Connecting to AWS MSK

The kafka-msk-iam feature adds IAM-based authentication for Amazon MSK clusters. It pulls in aws-config, aws-credential-types, and aws-msk-iam-sasl-signer. Use it alongside kafka-ssl:

cargo add shove --features kafka-ssl,kafka-msk-iam

MSK IAM clusters listen on port 9098 (SASL/IAM). SCRAM/PLAIN clusters use port 9096.

Minimal setup

MSK brokers use publicly-signed ACM certificates. KafkaTls::default() is correct — the OS trust store handles validation. No custom CA path is needed.

use shove::kafka::{KafkaConfig, KafkaSasl, KafkaTls};
 
let config = KafkaConfig::new("b-1.cluster.amazonaws.com:9098,b-2.cluster.amazonaws.com:9098")
    .with_tls(KafkaTls::default())
    .with_sasl(KafkaSasl::msk_iam("eu-west-2"));

KafkaSasl::msk_iam(region) resolves credentials from the standard AWS provider chain: environment variables, shared credentials file, EC2 instance metadata (IMDS), EKS pod identity / IRSA, and SSO. No explicit credential configuration is needed in most deployment environments.

For a non-default named profile, use KafkaSasl::msk_iam_with_profile:

let config = KafkaConfig::new(brokers)
    .with_tls(KafkaTls::default())
    .with_sasl(KafkaSasl::msk_iam_with_profile("eu-west-2", "production"));

The OAUTHBEARER mechanism and SASL_SSL security protocol are set automatically. Do not set sasl.mechanism or security.protocol manually, and do not set sasl.oauthbearer.config — token rotation is handled automatically by the library.

See examples/kafka/msk_iam.rs for a runnable walkthrough.

Declare topology

    broker.topology().declare::<OrderTopic>().await?;

topology().declare::<T>() creates the main Kafka topic and its configured DLQ. It does not create hold topics. Declaration is idempotent, so it is safe to call on every startup.

Producer topic auto-creation defaults to off (allow.auto.create.topics=false), so declare the topology before the first publish or provision the topic outside shove. A publish to a topic nobody declared fails with a Connection error after the produce timeout, and the topic stays absent. Set KafkaConfig::with_producer_auto_create_topics(true) to allow creation on first publish when the broker enables auto.create.topics.enable=true.

Bind to an infra-owned topic

Sometimes infra owns the topic: a Terraform module, a platform team's provisioning job, or a cluster whose credentials carry no Create or Alter permission. Then use external(), the backend-neutral flag NATS shares:

use shove::TopologyBuilder;

shove::define_topic!(
    PriceChanges,
    PriceChange,
    // The queue name must match the externally provisioned topic name.
    TopologyBuilder::new("price-changes")
        .external()
        .dlq()
        .build()
);

In this mode declare::<T>() does not create the main topic. It verifies that the topic, named after the queue, already exists. If it does not, declare returns a ShoveError::Topology, so a missing or misnamed topic fails fast at startup. Nothing falls back to an auto-created topic with default partitions and config. The check is a metadata request from a consumer-type client with no group.id. It therefore needs no group permission and cannot trigger the broker's topic auto-creation. shove never expands the topic's partitions and never reconciles its config. A consumer group registered with more max_consumers than the topic has partitions runs with the extra members idle. declare logs a warning when that happens. shove still creates and owns its DLQ topic, and still manages its consumer group. Provision the topic before the consumer starts.

A consumer then needs these ACLs: Describe and Read on the topic, and Describe and Read on the consumer group. The group is {queue}-consumer unless for_consumer_group or with_group_id names another. A publisher additionally needs Write on the topic. A declared DLQ needs Create, Describe and Write on the DLQ topic, so that shove can create it and route to it.

On an external topic, shove's consumer never writes into the topic when it settles an outcome: a Retry or a Defer waits in place, and a Reject publishes to the shove-owned DLQ. That is the guarantee, and it covers declaring and consuming. Publishing is outside it: a Publisher on this topology writes to the topic. Until the producer pins allow.auto.create.topics=false, a separate change in the same release, a publish to a topic that does not exist can create it with broker defaults on a broker with auto.create.topics.enable=true, so provision the topic before any publisher starts. Ack and Reject behave as everywhere else, because a Reject publishes to the shove-owned DLQ and not into the topic. How Retry and Defer are carried out is the consumer's retry strategy, RetryStrategy::Republish or RetryStrategy::InPlace, set with with_retry_strategy on ConsumerOptions::<Kafka> or KafkaConsumerGroupConfig. Ownership implies it: an external topology runs with InPlace and refuses Republish, because the republish would write into a topic infra owns. A shove-owned topology runs with Republish by default and may opt into InPlace, so a consumer that must not duplicate a retried record for the other groups on a fan-out topic can wait in place too. A FIFO consumer refuses InPlace, which it does not implement.

OutcomeRetryStrategy::Republish, the shove-owned defaultRetryStrategy::InPlace, implied by an external topology
Ackcommits the offsetthe same
Rejectpublishes to the DLQ when one is declared, then commitsthe same, the DLQ is shove's own topic
Retryrepublishes into the topic after the tier's delay, with the retry count in a headerwaits the tier's delay inside the handler's task, holding its prefetch slot, then hands the same record back with the retry count kept in memory; an exhausted budget goes to the DLQ as usual, with the in-memory count written to the dead letter's Shove-Retry-Count header
Deferrepublishes into the topic after the first tier's delaywaits the first tier's delay in place, then hands the same record back

The redelivered MessageMetadata carries redelivered: true and the in-memory retry count. A Defer leaves that in-memory budget unchanged, as it does under the republish strategy; on an external NATS stream, by contrast, the count is the broker's redelivery count, so a Defer there consumes budget. While every prefetch slot is held and at least one holder is a waiting handler, the consumer pauses its assignment and keeps polling. A long wait therefore does not evict the member from its group. Slots held only by running handlers do not pause the assignment by themselves, because a pause purges the fetch queue and refetches on resume. A record that arrives while every slot is held by a running handler is decoded and waits for a slot while the consumer keeps polling. A further record that arrives during that wait is handed back to the broker, and the assignment is paused until a slot frees. A partition the group assigns to the member during a pause is paused too. A record it delivers first is handed back to the broker and arrives again, in order, once a slot frees. The pause lasts until a slot frees: the slot that ends the wait serves the record already in hand, so with one slot the assignment stays paused through that record's handler as well. A shutdown during a wait completes nothing: the record stays uncommitted and is redelivered on restart. A waiting handler holds its slot for the whole delay, the contract a broadcast subscription has always had. Size prefetch_count for the number of records you accept to have waiting at once.

external() is incompatible with sequenced() and with every topic-config method. Those are with_topic_config, with_retention, with_retention_forever, with_retention_bytes, with_cleanup_policy and with_max_message_bytes. They configure topic creation and reconciliation, which external mode skips, and build() panics if they are combined. dlq(), dlq_named(), hold_queue(), for_consumer_group() and broadcast() stay available.

Publish

    let publisher = broker.publisher().await?;
    for i in 0..3 {
        publisher
            .publish::<OrderTopic>(&OrderCreated {
                order_id: format!("ORD-{i}"),
                amount: 99.99 + i as f64,
            })
            .await?;
        println!("Published order ORD-{i}");
    }

publisher().await? returns a Publisher<Kafka>. Messages are produced via rdkafka. The message key is derived from the topic's partition-key logic (or from T::sequence_key() for sequenced topics). For topic provisioning, see Declare topology.

Producer tuning

shove pins the correctness-critical producer settings (acks=all, enable.idempotence=true) and keeps them non-configurable. Idempotence caps max.in.flight.requests.per.connection at 5, so sustained throughput is bounded by messages per request — i.e. by batching. librdkafka's defaults (linger.ms=5, no compression) favor low latency; high-rate pipelines can raise the throughput ceiling with three opt-in knobs on KafkaConfig:

use shove::kafka::{KafkaCompression, KafkaConfig};
 
let config = KafkaConfig::new(brokers)
    .with_producer_compression(KafkaCompression::Lz4)
    .with_producer_linger_ms(25)
    .with_producer_batch_size(500_000);
  • with_producer_compression(KafkaCompression) — compression.type (None, Gzip, Snappy, Lz4, Zstd). Compresses each batch on the client, cutting producer→broker bytes as well as broker storage and replication traffic.
  • with_producer_linger_ms(u32) — linger.ms, how long the producer waits to accumulate a batch. Higher values trade a little latency for materially larger (and better-compressed) batches. Must stay below the pinned message.timeout.ms (5000): linger time counts toward the message timeout, so values at or above it would expire every publish. Rejected at connect.
  • with_producer_batch_size(u32) — batch.size, maximum bytes accumulated per batch (1..=i32::MAX, checked at connect). librdkafka caps the effective batch at min(batch.size, message.max.bytes), and message.max.bytes (default 1 MB) is not exposed — so this knob can lower the cap but not raise it above 1 MB.

Unset throughput settings keep librdkafka's defaults. acks, enable.idempotence, and max.in.flight.requests.per.connection are deliberately not exposed.

The client holds a single producer, shared by publisher() and by the consumer's retry, defer, and DLQ republishes — a long linger also delays those republishes (and the offset commits gated on them for sequenced topics), so keep linger modest when consumers republish on the same client.

Consume

    let mut group = broker.consumer_group();
    group
        .register::<OrderTopic, _>(
            ConsumerGroupConfig::new(KafkaConsumerGroupConfig::new(1..=1)),
            || OrderHandler,
        )
        .await?;
 
    // Stop after 3 s for demo purposes, or on ctrl-c.
    let outcome = group
        .run_until_timeout(
            async {
                tokio::select! {
                    _ = tokio::time::sleep(Duration::from_secs(3)) => {}
                    _ = tokio::signal::ctrl_c() => {}
                }
            },
            Duration::from_secs(10),
        )
        .await;

consumer_group() registers a native Kafka consumer group. KafkaConsumerGroupConfig::new(min..=max) sets the autoscale bounds. Kafka handles partition assignment and rebalance automatically when the group membership changes.

Group configuration

use shove::kafka::{KafkaAutoOffsetReset, KafkaConsumerGroupConfig};

let cfg = KafkaConsumerGroupConfig::new(1..=8)
    .with_prefetch_count(20)
    .with_max_retries(5)
    .with_handler_timeout(Duration::from_secs(30))
    .with_concurrent_processing(true)
    .with_group_id("billing-orders-consumer")          // override the default `{queue}-consumer`
    .with_auto_offset_reset(KafkaAutoOffsetReset::Latest)
    .with_commit_interval(Duration::from_secs(5));     // commit at most every 5 s (default 500 ms)
  • with_prefetch_count(u16) — librdkafka in-flight cap per consumer task. Default 10.
  • with_max_retries(u32) — retry budget before dead-lettering. Default 10.
  • with_handler_timeout(Duration) — per-message wall-clock deadline. Default 30s (see Handlers & Context).
  • with_concurrent_processing(bool) — dispatch each fetched message to its own tokio task (rejected for sequenced topics). Default false.
  • with_group_id(impl Into<String>) — override the broker-side consumer group ID. Defaults to "{queue}-consumer". Set this when two independent services consume the same topic and must each receive every message (fan-out) — otherwise they share a group and compete for partitions. Prefer .for_consumer_group(...) on the topology (below), which sets the group and the DLQ/hold-queue names together.
  • with_auto_offset_reset(KafkaAutoOffsetReset) — Earliest (default, replay history), Latest (tail-only), or None (refuse silent replay/skip on a fresh group). The same setting is available on ConsumerOptions::<Kafka> for the direct and supervisor paths, which have no group config to carry it. Under None, a group member with no usable committed offset ends with a ShoveError::Topology naming librdkafka's AutoOffsetReset answer. A reconnect would only meet the same answer, so the member does not retry. The error ends that member, not the group: the group run keeps waiting for its shutdown signal. The error count the run reports in SupervisorOutcome at shutdown grows by one per member. Without autoscaling nothing replaces the member, so the group stays below its configured minimum for the life of the process. With autoscaling enabled, the autoscaler tick respawns the member through the respawn supervisor, and each replacement meets the same answer. The supervisor tops the group up in full on each of its first five rounds, at least 2, 4, 8 and 16 seconds apart. The fifth round opens the circuit and rests for 300 seconds. After that, each round spawns one probe member per 300-second cooldown. The circuit never becomes terminal, so a group under None with autoscaling on keeps probing until an operator intervenes. Lag-driven scale-up is a second path outside that circuit. A group with no committed offset reports its retained backlog as lag, so while records remain on the topic the autoscaler can add members. The 300-second cooldown therefore bounds the respawn probes, not every member the group creates. Commit a starting position first, with reset_consumer_group_offsets (below), or pick another policy. The DLQ drain (run_dlq_with_options) refuses the setting because it hard-codes earliest. That policy applies only while the drain's group has no usable committed offset. A fresh drain therefore never skips a dead letter, and a restarted drain resumes from its commit.
  • with_commit_interval(Duration) - how often each consumer commits the offsets its handlers completed. The default is 500ms. The interval must be positive and at most one hour; the setter panics otherwise, so an unrepresentable commit deadline is refused at configuration time and never met in the receive loop. The check runs again where the consumer starts, so a value written to the public kafka_commit_interval field past the setter is refused there. A longer interval sends fewer OffsetCommit requests and widens the replay window after a crash or rebalance by the same amount. Standard groups only: a FIFO consumer commits each message as it settles, so register_fifo refuses a config that sets an interval. The DLQ drain commits each dead letter as it settles and refuses the setting on ConsumerOptions::<Kafka> for the same reason. See Offset commit semantics. The same setting is available on ConsumerOptions::<Kafka> for the direct and supervisor paths.

Fan-out — a second reader on the same topic

A bare with_group_id splits the group but not the retry chain: both readers still derive {queue}-dlq and {queue}-hold-*, so each drains the other's dead and held messages. Declaring the second reader's topology with for_consumer_group splits both:

TopologyBuilder::new("order-settlement")
    .for_consumer_group("settlement-audit")
    .hold_queue(Duration::from_secs(5))
    .dlq()                                   // order-settlement-settlement-audit-dlq
    .build()

The resolved group IDs follow the topology:

ConsumerNo fan-out groupfor_consumer_group("settlement-audit")
Standardorder-settlement-consumerorder-settlement-settlement-audit-consumer
FIFO (sequenced)order-settlement-fifoorder-settlement-settlement-audit-fifo
DLQ drainorder-settlement-dlq-consumerorder-settlement-settlement-audit-dlq-consumer

Precedence is: an explicit with_group_id (on either KafkaConsumerGroupConfig or ConsumerOptions::<Kafka>) > the topology's fan-out group > the {queue}-consumer default. The explicit override staying on top means adding for_consumer_group to a topology cannot move an already-deployed consumer off the group it holds committed offsets under. The autoscaler resolves the same group ID as the broker-side consumer in every case, so lag is read from the group that is actually committing.

Re-anchoring a group (seek to tail / head / timestamp)

auto.offset.reset only decides where a group starts when it has no usable committed offset. Once the group has committed, the setting is inert — which is why "just seek to the tail" so often turns into minting a throwaway group ID (orders-v2, orders-20260812, …). That works, but it strands the old group's offsets and its lag metrics forever, and the generation suffix becomes a permanent piece of config nobody dares remove.

reset_consumer_group_offsets rewrites the group's committed offsets in place — the library-side equivalent of kafka-consumer-groups.sh --reset-offsets --execute:

use shove::kafka::{KafkaConsumerGroupConfig, KafkaOffsetReset};

let config = KafkaConsumerGroupConfig::new(1..=4);

// Operator-initiated: re-anchor at the tail before the consumers start.
if std::env::var("PRICES_SEEK_TO_TAIL").is_ok() {
    let report = broker
        .reset_consumer_group_offsets::<Prices>(&config, KafkaOffsetReset::Latest)
        .await?;
    tracing::warn!(?report, "re-anchored the prices group at the tail");
}

let mut group = broker.consumer_group();
group
    .register::<Prices, _>(ConsumerGroupConfig::new(config), || Handler)
    .await?;
  • KafkaOffsetReset::Latest — every partition's high watermark. The seek-to-tail case: a latest-value sink that must serve fresh data now rather than after crawling days of backlog.
  • KafkaOffsetReset::Earliest — every partition's low watermark: replay all retained history.
  • KafkaOffsetReset::Timestamp(ms) — the first record at or after that point, in milliseconds since the Unix epoch (the same unit as --to-datetime). Partitions with no record at or after it re-anchor at their high watermark.

The group ID is resolved from config and the topology exactly as register would resolve it — following the same precedence as the fan-out table above, plus the -fifo suffix for a sequenced topic — so the offsets rewritten are the ones the consumers will actually read.

The group must be inactive. Kafka only accepts an offset reset while the group has no live members; with consumers running the call returns ShoveError::Validation naming the active member count, and the broker enforces the same rule independently. Re-anchor at process start, before the group is registered. A group does not go inactive the instant its consumers stop — the coordinator drops each member as its LeaveGroup lands — so a reset issued immediately after run_until_timeout returns may need a brief retry.

The returned KafkaOffsetResetReport carries one entry per partition with the previous and new offsets (and delta(), positive for records skipped, negative for history replayed). It is the only record of where the group was before it moved, so log it. is_noop() reports that every partition already sat at its target.

Kafka is the only backend with this API: it is the only one shove supports where a group's read position is a broker-side committed offset an operator can rewrite. Redis Streams' XGROUP SETID is the nearest equivalent and is not yet exposed.

Starting a broadcast subscription elsewhere than the tail

A broadcast subscription assigns every partition itself and never commits, so it has no stored position to reset. Where it starts is a property of the subscription, set on its options:

use shove::BroadcastStart;

let mut subscriber = broker.broadcast_subscriber();
subscriber.subscribe::<CacheInvalidations, _>(
    Evict,
    ConsumerOptions::<Kafka>::new()
        // Replay everything the topic still retains, then keep tailing.
        .with_broadcast_start(BroadcastStart::Head)
        // The inert `group.id` the handle carries; see below.
        .with_group_id("cache-invalidations-broadcast"),
)?;

BroadcastStart is backend-neutral, and Kafka is the backend that honours all three variants on this version. A broadcast subscription refuses with_commit_interval and with_auto_offset_reset at subscribe(): it commits nothing and assigns every partition at an explicit offset, so neither setting would change anything, and a setting that changes nothing is refused rather than dropped. The competing-consumer entry points refuse with_broadcast_start for the same reason.

  • BroadcastStart::Tail, and no call at all - every partition at its end offset: deliver-new, the broadcast contract.
  • BroadcastStart::Head - every partition at its beginning, so a fresh instance replays the retained log before it tails. This is the right start for a subscriber that rebuilds an in-memory view from the topic instead of from a separate store.
  • BroadcastStart::Timestamp(ms) - the first record at or after that point, in milliseconds since the Unix epoch. The lookup is offsets_for_times, exactly as reset_consumer_group_offsets resolves it. A partition with no record at or after the point starts at its tail. A partition the lookup cannot answer for fails the subscription rather than silently starting at the tail.

A partition added while the subscription runs is picked up within a few seconds and assigned at the same start. A reconnect after a broker outage re-resolves the start as it stands then. A Head subscription therefore replays the retained log again, the price of having no stored position.

with_group_id names the inert group.id the groupless handle is configured with, and the default is {queue}-broadcast. Nothing joins under it and nothing commits to it either way. Set it when the cluster's ACLs grant group Describe on one prefix only. librdkafka looks up the configured group's coordinator even for an assign-only handle, and an unauthorised lookup draws a GroupAuthorizationFailed on every attempt. The subscription tolerates that error: it warns once and keeps fetching, because fetching never goes through the coordinator. An id under the granted prefix keeps the logs clean.

Replication factor

Topics created by declaration get replication factor 1 by default, which is fine for single-broker dev but unsafe in production. Set a default for topics declared through this registry:

let mut group = broker
    .consumer_group()
    .with_default_replication_factor(3);                // applied to topics created by declaration

Or set it per-declaration on the topology declarer:

broker
    .topology()
    .with_replication_factor(3)
    .declare::<Orders>()
    .await?;

create_topic is idempotent and will not lower an existing topic's replication. Pre-creating topics out of band (Terraform, MSK console) is also fine — the declarer is a no-op when the topic already exists.

Sequenced delivery

Messages for the same key stay in order. Kafka uses partition-key routing: T::sequence_key() becomes the Kafka message key, so messages with the same key always land on the same partition. A single consumer handles each partition at a time, guaranteeing that messages for the same key are never processed concurrently.

Ordering is partition-scoped: two messages with different keys may land on different partitions and be processed concurrently. The partition count is fixed at topic creation time and caps the maximum degree of parallelism across the consumer group.

See Sequenced Topics for the full ordering model, and Sequenced example for a runnable walkthrough.

Consumer groups + autoscaling

Kafka consumer groups are native — broker.consumer_group() creates a standard Kafka consumer group. Call register to associate a topic with handlers and bounds:

use shove::kafka::KafkaConsumerGroupConfig;
use shove::{Broker, ConsumerGroupConfig, Kafka};

let mut group = broker.consumer_group();
group
    .register::<OrderTopic, _>(
        ConsumerGroupConfig::new(KafkaConsumerGroupConfig::new(1..=4)),
        || MyHandler,
    )
    .await?;

The autoscaler measures consumer lag (the offset gap between the latest produced message and the latest committed offset) and adjusts the number of active consumers within the min..=max range.

Note: consumer-group rebalances occur during scale-up and scale-down events. Rebalances cause a brief delivery pause while Kafka reassigns partitions. Plan for this when setting autoscale bounds and drain timeouts.

See the Basic example for a full runnable walkthrough.

Offset commit semantics

Outcome::Ack marks the offset complete as soon as the handler returns. Outcome::Retry and Outcome::Defer mark it complete only once the broker has acknowledged the delayed republish to the hold topic. Under RetryStrategy::InPlace, which an external topology implies and a shove-owned one may opt into, there is no republish. The offset is marked complete only when the in-place retry or defer ends in a terminal outcome. That is an Ack, or a Reject or an exhausted retry budget that was dead-lettered or discarded. This closes the publish-then-commit race. A broker crash or process kill between handler return and republish redelivers the message on restart rather than silently dropping it.

Completions are tracked in memory per partition and committed asynchronously, at most once per commit interval. The interval is 500 ms by default, configurable with with_commit_interval on KafkaConsumerGroupConfig or ConsumerOptions::<Kafka>. The gate keeps the number of in-flight commits O(1) however fast the consumer runs. A longer interval trades coordinator requests for a wider replay window. A commit the coordinator rejects is re-offered on a later drain. A consumer whose commits keep being rejected with no rebalance resolving them is treated as fenced and reconnects. That threshold grows with the interval, because the streak can only clear on a drain.

The commit position is derived from the offsets the broker actually delivered, never from a run of consecutive integers. A hole left by log compaction or by a transactional producer's control records does not hold the position back. Compacted and transactional topics therefore commit past such a hole as soon as every delivered record below it has completed. That holds for a hole between delivered records. A trailing control record is different: with nothing delivered after it, the committed position stays one below the high watermark. An idle transactional topic therefore reads a lag of 1 until the next record arrives, and nothing is redelivered.

On shutdown the consumer first waits for its in-flight handlers to finish. That wait is bounded by the handler timeout and by any in-place delays. It then drains the republish queue and commits its final position synchronously. The commit runs on a dedicated thread that also owns the consumer's close. The thread is spawned before it is handed the consumer, so a failed spawn never closes the consumer on the runtime thread. If no thread can be spawned, nothing commits and the close moves to a second, close-only thread. If that thread cannot be spawned either, the handle is leaked and logged at error level. The leaked consumer keeps heartbeating, so while the process lives its partitions stay assigned until max.poll.interval.ms (5 min) passes without a poll. Once the process exits, the broker drops the member after the session timeout, as after a crash. The consumer waits for that commit for at most 20 s (SHUTDOWN_COMMIT_DEADLINE in kafka::constants). The deadline is chosen against the 30 s termination grace Kubernetes gives a Pod. Past the deadline the consumer logs that the batch may be redelivered and returns. The thread finishes on its own and never holds the runtime or the process, so a frozen coordinator cannot turn a deploy into a SIGKILL.

What that means for replay:

  • After a crash, the records completed since the last accepted commit are redelivered. That is about one interval's worth when commits are accepted and every earlier offset on the partition had completed. A rejected commit that is still being re-offered, or an earlier record whose handler is still running, holds the position and widens that window.
  • After a rebalance, a moved partition replays from its last accepted commit.
  • After a clean shutdown, a record whose handler returned a terminal outcome and whose offset reached the final commit is not redelivered. A record whose handler had not settled when shutdown fired stays uncommitted and is redelivered. So is the last batch when the final commit hits its deadline. An asynchronous commit librdkafka was still retrying can also land after the final one. A coordinator hiccup overlapping the shutdown can therefore replay a few records.

Outcome::Reject routes to the DLQ (if configured) and commits the offset; the original is not redelivered. Messages that exhaust max_retries follow the same DLQ-then-commit path.

Schema Registry

The kafka-schema-registry feature adds Confluent Schema Registry support for Kafka consumers (decode) and producers (encode). On the consume side, each incoming message is unwrapped from the Confluent wire frame (magic byte 0x00 + 4-byte big-endian schema id; for Protobuf, also the message-index array), the schema id is resolved against the registry (cached in memory, with single-flight deduplication of concurrent cold misses and a configurable negative-TTL for permanent failures), and the inner payload is decoded via the topic's existing Codec (JsonCodec or ProtobufCodec). On the produce side, an opt-in publisher wraps each encoded payload in the same wire frame — see Producer-side encoding.

The feature works with the Confluent Schema Registry and with Redpanda's built-in Schema Registry, which exposes the same Confluent-compatible REST API and wire format. The e2e test suite validates both paths.

Install

cargo add shove --features kafka-schema-registry

Build the registry client

use std::time::Duration;
use shove::schema_registry::{SchemaRegistry, SchemaRegistryAuth};

let registry = SchemaRegistry::builder("https://schema-registry:8081")
    .auth(SchemaRegistryAuth::Basic {
        user: "sr-user".into(),
        pass: "sr-pass".into(),
    })
    .timeout(Duration::from_secs(3))
    .max_retries(2)
    .negative_cache_ttl(Duration::from_secs(60))
    .build();

SchemaRegistry::builder(url) accepts any http:// or https:// base URL, with or without user:pass@ userinfo. Auth options are:

  • SchemaRegistryAuth::None — no authentication (default)
  • SchemaRegistryAuth::Bearer(token) — Authorization: Bearer <token>
  • SchemaRegistryAuth::Basic { user, pass } — HTTP Basic auth

Credentials require TLS. Configuring any credential against a base URL that is not https:// makes build() panic, because the secret would be sent in cleartext on every schema fetch. Three things count as a credential:

  • any SchemaRegistryAuth other than None;
  • user:pass@ userinfo in the base URL — the HTTP client lifts it out and replays it as an Authorization: Basic header, so it is a credential even with SchemaRegistryAuth::None;
  • any header set with .header(...), whose value is assumed to be a secret. Use .non_secret_header(...) for one that is not.

The base URL is parsed before it is judged, so alternate spellings of a plaintext scheme — http:registry:8081, http:/registry:8081 — are refused too rather than slipping past a prefix check. A base URL that does not parse is treated as plaintext.

A credential-bearing client also stops following redirects. A schema registry has no legitimate reason to redirect, and a 302 would otherwise replay a custom secret header to whatever host and scheme the Location names. An unauthenticated client is unaffected.

An unauthenticated http:// registry is unaffected. For a registry that is genuinely unreachable from an untrusted network — a local development stack — opt in explicitly:

use shove::schema_registry::{SchemaRegistry, SchemaRegistryAuth};

let registry = SchemaRegistry::builder("http://localhost:8081")
    .auth(SchemaRegistryAuth::Bearer("dev-token".into()))
    .allow_plaintext_credentials()
    .build();

The returned value is an Arc<SchemaRegistry>. Clone it to share the same schema cache across multiple consumers or a consumer group.

Per-consumer configuration

use std::sync::Arc;
use shove::schema_registry::{SchemaEnforcement, SchemaRegistry, SchemaRegistryAuth};
use shove::consumer::ConsumerOptions;
use shove::markers::Kafka;

let registry = SchemaRegistry::builder("http://schema-registry:8081").build();

let opts = ConsumerOptions::<Kafka>::new()
    .with_schema_registry(Arc::clone(&registry))
    .with_schema_enforcement(SchemaEnforcement::Enforce)
    .accept_schema_subjects(["orders-value"]);

Per-consumer-group configuration

Attaching the registry on a KafkaConsumerGroupConfig shares the Arc — and therefore the same in-memory schema cache — across every consumer spawned in the autoscaling group:

use std::sync::Arc;
use shove::kafka::{KafkaConsumerGroupConfig};
use shove::schema_registry::{SchemaEnforcement, SchemaRegistry};

let registry = SchemaRegistry::builder("http://schema-registry:8081").build();

let cfg = KafkaConsumerGroupConfig::new(1..=8)
    .with_schema_registry(Arc::clone(&registry))
    .with_schema_enforcement(SchemaEnforcement::Enforce)
    .accept_schema_subjects(["orders-value"]);

Enforcement modes

with_schema_enforcement controls what happens when a message's registered subject is not in the accepted set:

  • SchemaEnforcement::Enforce (default) — the message is routed to the DLQ with death reason schema_validation_failed. Choose this for production where a subject mismatch is a producer misconfiguration.
  • SchemaEnforcement::Permissive — the mismatch is logged and counted, and the message is decoded anyway. Use this during migration windows or when multiple producers share a topic with different subject conventions.

Accepted subjects

accept_schema_subjects([...]) pins the set of Confluent schema subjects that are accepted for a consumer. When not called, the default is derived from the queue name using the Confluent TopicNameStrategy: "{queue}-value". For a topic named orders, that resolves to orders-value.

Protobuf message index

A Confluent protobuf frame carries a message-index array naming which message of the schema file the bytes encode. [0] is the file's first top-level message, [1] the second, and [0, 2] the third message nested in the first. shove reads the array as Confluent writes it: zigzag varints, with the single-byte shorthand for [0]. By default the index is read and ignored, so a ProtobufCodec<M> decodes whatever message the producer framed as M. Two messages with compatible field numbers then decode into each other silently. require_schema_message_index([0]) pins the index a consumer accepts. A frame with any other index is routed to the DLQ with death reason schema_message_index_rejected, counted under reason="schema_frame". The check runs before the schema id is resolved, so a rejected frame costs no registry round trip. JSON frames carry no index and the setting is inert on them. The setter lives on ConsumerOptions::<Kafka>, BatchConsumerOptions::<Kafka> and KafkaConsumerGroupConfig, next to accept_schema_subjects. All three panic on an empty requirement, a negative index or more than 1024 indexes. The parser never yields such an index, so the requirement could never match. Refusing it at configuration time points at the cause at startup. A value written to the public schema_message_index field past the setter is checked when the options are handed to the consumer.

Producer-side encoding

The same feature lets a publisher emit Confluent-framed messages, symmetric with how the consumer is configured. Attach a SchemaRegistry to a KafkaPublisherConfig and obtain the publisher with Broker::publisher_with; the publisher then wraps each codec-encoded payload in the Confluent wire frame using the latest registered schema id for the subject. Like the consumer side, framing is a publisher-layer concern — it is not baked into the topic's Codec.

use std::sync::Arc;
use shove::kafka::KafkaPublisherConfig;
use shove::schema_registry::SchemaRegistry;

let registry = SchemaRegistry::builder("http://schema-registry:8081").build();

// `broker` is a `Broker<Kafka>`.
let publisher = broker
    .publisher_with(KafkaPublisherConfig::new().with_schema_registry(Arc::clone(&registry)))
    .await?;

// `publisher.publish::<OrdersTopic>(&order).await?` now emits SR-framed bytes.

Details:

  • Subject — defaults to the Confluent TopicNameStrategy "{topic}-value"; override with KafkaPublisherConfig::with_subject("...").
  • Schema id — resolved once per subject via GET /subjects/{subject}/versions/latest and cached. shove carries no schema text; the subject must already be registered.
  • Codecs — framing applies only to JsonCodec and ProtobufCodec topics; other codecs are published unframed. For Protobuf the message-index is [0] (the first message type).
  • Opt-in — a publisher built with the plain Broker::publisher() does no framing, so existing publishers are byte-for-byte unchanged.

When the registry cannot answer

A record whose schema id the registry cannot resolve right now is neither poison nor processed. The consumer keeps it, pauses its partition assignment, and asks again every second (REGISTRY_RETRY_DELAY in kafka::constants). Every partition the member holds is paused, healthy ones included, so one stalled record stops that member's delivery until the registry answers. It keeps asking until the registry answers or the consumer shuts down. The offset stays uncommitted throughout, so a shutdown or a crash mid-wait redelivers the record. Each wait counts shove_messages_failed_total{reason="schema_unavailable"} once and logs a warning with the schema id. Nothing is dead-lettered or discarded, so shove_messages_discarded_total does not move. "Cannot answer" means a transport failure, or a 5xx, a 408 or a 429 that outlasts the client's own retries. The batch consumer flushes the records ahead of the stalled one first. Its committed span therefore ends before that record until it decodes. If that flush comes back Retry or Defer, the span is sought back to its start. The stalled record is then handed back to the broker as well. It arrives again behind the rewound records, so every batch stays in offset order. The loop keeps polling while it waits, so the member stays in its group for the whole outage, however long. A partition the group assigns to the member during the wait is paused too. A record it delivers first is handed back to the broker and arrives again, in order, once the wait ends.

A registry that does answer, but wrongly for the deployment, is not an outage and is not waited on. A 401 or 403, a redirect the client cannot follow, another unexpected status, or an undecodable response ends the consumer. The error is a ShoveError::Topology naming the schema id and the fault. On the registry path the group's supervision restarts the member with backoff. The fault is therefore visible in the logs and the metrics rather than hidden behind an endless wait. A 404 is still the registry's definite answer that the id is unknown. Such a record goes to the DLQ with reason schema_resolve_failed, as before.

Codec requirement

Registry decoding is only supported for topics using a JsonCodec or ProtobufCodec. A message arriving on a topic whose codec is neither JSON nor Protobuf is routed to the DLQ with reason schema_unsupported_codec.

Redpanda compatibility

The Confluent wire format (magic byte + schema id) and REST API endpoints used by kafka-schema-registry are fully compatible with Redpanda's built-in Schema Registry. No additional configuration is needed — point SchemaRegistry::builder at the Redpanda Schema Registry URL.

DLQ consumer caveat

Metrics

The Kafka backend emits shove_messages_failed_total with the following reason labels (see Observability):

  • oversize — payload exceeds max_message_size before deserialization.
  • deserialize — JSON/codec decode failed; routed to DLQ.
  • timeout — handler exceeded handler_timeout; retried.
  • max_retries_exceeded — retry budget exhausted; routed to DLQ.
  • rejected — handler returned Outcome::Reject; routed to DLQ.
  • schema_validation - the registry answered and ruled the record out. The death reason is schema_validation_failed when the subject is not in the accepted set under Enforce mode. It is schema_resolve_failed when the registry does not know the schema id. Routed to DLQ. Requires kafka-schema-registry.
  • schema_frame - the Confluent frame cannot be used, so no lookup is made. The death reason is schema_frame_invalid for a malformed frame. It is schema_unsupported_codec when the topic uses a codec other than JSON or Protobuf. It is schema_message_index_rejected when the protobuf message index is not the one require_schema_message_index pins. Routed to DLQ. Requires kafka-schema-registry.
  • schema_unavailable - the registry could not answer for the record's schema id, so the record is kept and the lookup retried. It is never dead-lettered or discarded, and it counts once per wait. Requires kafka-schema-registry.

shove_backend_errors_total carries the backend="kafka" label for connection drops, produce failures, and consume-stream errors.

Gotchas

  • Replication factor defaults to 1 — safe for single-broker dev, unsafe in production. Set with_default_replication_factor(3) on the registry (or with_replication_factor on the topology declarer) before the first declaration. The default exists so the no-config dev path works; the constant lives in kafka::constants::DEFAULT_REPLICATION.
  • Partition count is set at topic creation time and is rarely changed after the fact. Choose carefully — partition count caps the maximum degree of parallelism a consumer group can achieve (one consumer per partition maximum).
  • Consumer-group rebalances during scale events cause brief delivery pauses. The larger the group and the more topics, the longer the rebalance. shove pins session.timeout.ms to 10 seconds and max.poll.interval.ms to 5 minutes (constants in kafka::constants); neither is settable via KafkaConfig. If handlers are slow, raise the handler timeout with with_handler_timeout (per group) or with_default_handler_timeout (registry-wide) instead.
  • On Windows, the cmake-build feature is required for rdkafka (the underlying C library). This is handled automatically via the [target.'cfg(windows)'.dependencies] entry in Cargo.toml.
  • SASL and TLS require the kafka-ssl feature. PLAIN, SCRAM-SHA-256 and SCRAM-SHA-512 are supported out of the box. librdkafka implements PLAIN itself in every build, and SCRAM and OAUTHBEARER in its OpenSSL build, which kafka-ssl selects. Shove exposes OAUTHBEARER only as KafkaSasl::MskIam under kafka-msk-iam. There is no generic token provider. GSSAPI/Kerberos needs Cyrus SASL, which the separate kafka-gssapi feature links for direct rdkafka use. KafkaSasl exposes no GSSAPI variant, so that feature gates no shove code and shove itself cannot authenticate with Kerberos. For AWS MSK IAM, use the kafka-msk-iam feature instead.
  • For publish failures on missing topics, see Declare topology.
  • librdkafka system dependencies on Debian/Ubuntu CI runners: libsasl2-dev is required only with kafka-gssapi. kafka-ssl alone links OpenSSL and no SASL library, and the shove ci.yml asserts that with cargo tree.

Examples

  • Basic — publish/consume round trip with hold queues and DLQ
  • Sequenced — partition-key ordering
  • Audited Consumer — MessageHandlerExt::audited wrapping
  • Stress — throughput benchmarking
  • MSK IAM — IAM authentication against Amazon MSK
  • Schema Registry — Confluent Schema Registry decode with subject enforcement

See also

  • Liveness Probes — wire Broker::ping into a k8s health endpoint.
  • Broadcast - per-instance fan-out via a groupless assign(), leaving __consumer_offsets untouched. It starts at the tail by default, or at the head or a timestamp with with_broadcast_start.