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:latestThe integration tests use testcontainers with the Apache Kafka module, so any runnable example also spins up a container automatically.
Install
cargo add shove --features kafkaFor 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-sslFor 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-gssapiThis 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-iamConnect
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-iamMSK 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.
The DLQ topic of an external topology is infra's as well, under its default {queue}-dlq name or the one dlq_named gives it.
declare probes it exactly like the main topic and returns ShoveError::Topology when it is missing; shove never creates it and never expands its partitions.
On a topology that is not external, declare creates its DLQ with eight partitions and expands an existing one to that.
shove still manages its consumer group.
Provision the topic and its DLQ before the consumer starts.
The registry path is the only one that runs declare for you.
A consumer that starts without it, through KafkaConsumer::run, a ConsumerSupervisor, run_batch, run_fifo, a broadcast subscription or run_dlq, runs the same metadata probe once at startup, before its reconnect loop and before it consumes a record.
It probes the main topic when the topology is external, and the DLQ topic wherever the path publishes to it, owned or external: the standard, batch and FIFO consumers.
A broadcast subscription publishes no dead letters, so it probes the external main topic only.
run_dlq reads from the DLQ and publishes nowhere. On an external topology it probes the main topic and the DLQ topic, both infra's; on a shove-owned one it probes nothing and waits for the DLQ topic, as a consumer waits for an owned main topic.
A missing topic is ShoveError::Topology naming it, returned before any record is consumed.
A broker that cannot be reached, or does not answer in time, says nothing about the topic and is not a Topology error: the probe retries under the consumer's usual backoff and max_reconnect_attempts, and a stop ends it.
A shove-owned main topic is never probed, because another process may still declare it.
With with_producer_auto_create_topics(true) a missing shove-owned DLQ is not probed either, because the first dead-letter publish creates it.
An external topology with a DLQ is refused outright on a client with with_producer_auto_create_topics(true), at register and at every consumer start that publishes dead letters (a publish-only declare is not refused, since it never writes a dead letter): a dead-letter publish after that topic disappeared would recreate the infra-owned topic with broker defaults. Consume it through a client without producer auto-creation.
The probe is a metadata request, so it builds no producer and needs no permission beyond Describe on the topic.
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.
shove's producer is idempotent.
librdkafka requests its producer id 500 ms after the producer is created, before any publish, with an InitProducerId request.
A broker authorizes that request as IdempotentWrite on the cluster, or since Kafka 2.8 (KIP-679) as Write on a topic, and a consumer holds neither.
The client builds that producer at the first publish, not at connect.
The first publish may come from publisher() or from a consumer's retry, defer or DLQ republish.
A service that never publishes, through publisher() or a republish, therefore never builds a producer and needs no permission beyond the consumer's.
connect still builds the consumer-type client that ping uses, so a connection setting librdkafka refuses fails at connect as before.
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 topology infra creates the DLQ, so shove needs only Describe and Write on 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 DLQ.
That is the guarantee, and it covers declaring and consuming.
Publishing is outside it: a Publisher on this topology writes to the topic.
The producer pins allow.auto.create.topics=false, so a publish to a topic that does not exist fails after the produce timeout and creates nothing.
Only with_producer_auto_create_topics(true) lets a publish create a topic, with broker defaults, on a broker running auto.create.topics.enable=true.
Provision the topic before any publisher starts.
Ack and Reject behave as everywhere else, because a Reject publishes to the 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.
| Outcome | RetryStrategy::Republish, the shove-owned default | RetryStrategy::InPlace, implied by an external topology |
|---|---|---|
Ack | commits the offset | the same |
Reject | publishes to the DLQ when one is declared, then commits | the same, to the DLQ topic infra provisioned |
Retry | republishes into the topic after the tier's delay, with the retry count in a header | waits 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 |
Defer | republishes into the topic after the first tier's delay | waits 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 record the consumer had fetched but not yet handed to a handler when shutdown fired is dropped unhandled too.
It is redelivered on restart, behind the record that was waiting.
A revoke of the partition ends the wait at once with nothing completed.
The partition's next owner is handed the record again from the committed offset.
A handler result carries the assignment it was made under, and one from an earlier assignment of the partition is ignored.
With one slot, no later record of a partition reaches the handler while a record of that partition waits in place.
That holds while the consumer runs, on a stop, and across a rebalance.
A waiting handler holds its slot until the delay ends, shutdown starts or its partition is revoked.
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 pinnedmessage.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 atmin(batch.size, message.max.bytes), andmessage.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, built at the first publish.
publisher() and the consumer's retry, defer, and DLQ republishes share it.
A long linger also delays those republishes, and the offset commits gated on them for sequenced topics.
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. Default10.with_max_retries(u32)— retry budget before dead-lettering. Default10.with_handler_timeout(Duration)— per-message wall-clock deadline. Default30s(see Handlers & Context).with_concurrent_processing(bool)— dispatch each fetched message to its own tokio task (rejected for sequenced topics). Defaultfalse.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), orNone(refuse silent replay/skip on a fresh group). The same setting is available onConsumerOptions::<Kafka>for the direct and supervisor paths, which have no group config to carry it. UnderNone, a group member with no usable committed offset ends with aShoveError::Topologynaming librdkafka'sAutoOffsetResetanswer. 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 inSupervisorOutcomeat 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 underNonewith 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, withreset_consumer_group_offsets(below), or pick another policy. The DLQ drain (run_dlq_with_options) refuses the setting because it hard-codesearliest. 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 is500ms. 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. A value written to the publickafka_commit_intervalfield past the setter is refused there withShoveError::Topology, never with a panic. 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, soregister_fiforefuses a config that sets an interval. The DLQ drain commits each dead letter as it settles and refuses the setting onConsumerOptions::<Kafka>for the same reason. See Offset commit semantics. The same setting is available onConsumerOptions::<Kafka>for the direct and supervisor paths. It is shorthand forwith_commit_policy(CommitPolicy::Interval(..)), below.with_commit_policy(CommitPolicy)- how each consumer commits the offsets its handlers completed.CommitPolicy::Interval(Duration)is the default, at500ms, andwith_commit_interval(d)is shorthand for it.CommitPolicy::PerRecordcommits every completion before the next record is handed out. The commit runs synchronously on a thread of its own, and the consumer waits for the coordinator's answer while it keeps polling. Every record costs a commit round trip and a fetch round trip, so throughput is bounded by that latency. It is an opt-in for a handler that is not idempotent. It needs one prefetch permit:with_prefetch_count(1), or concurrent processing off, which clamps the count to one. With more than one permit a completion above an unfinished lower offset confirms nothing.registerandKafkaConsumer::runtherefore refuse such a configuration before anything is spawned. The same setting is available onConsumerOptions::<Kafka>for the direct and supervisor paths. A FIFO consumer, the DLQ drain and a broadcast subscription refuse a policy as they refuse an interval. See Offset commit semantics.
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:
| Consumer | No fan-out group | for_consumer_group("settlement-audit") |
|---|---|---|
| Standard | order-settlement-consumer | order-settlement-settlement-audit-consumer |
| FIFO (sequenced) | order-settlement-fifo | order-settlement-settlement-audit-fifo |
| DLQ drain | order-settlement-dlq-consumer | order-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_policy, its shorthand with_commit_interval, and with_auto_offset_reset at subscribe().
It commits nothing and assigns every partition at an explicit offset, so none of those settings would change anything.
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 isoffsets_for_times, exactly asreset_consumer_group_offsetsresolves 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.
Under CommitPolicy::Interval, the default, 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.
Under CommitPolicy::PerRecord, set with with_commit_policy on the same two configs, every completion is committed at once.
The next record is handed out only once the coordinator has accepted that commit.
The commit runs synchronously on a thread of its own, so it never blocks the runtime thread.
The consumer pauses its assignment the moment it takes a record, before anything decides its fate, and keeps polling.
A schema registry stall on that record runs inside that pause, and the commit of the record resumes the assignment, not the stall.
Rebalance callbacks are therefore served, and the member stays inside max.poll.interval.ms.
A record dropped before the handler, oversize or undecodable, commits its position the same way.
It holds the next record back until that commit is accepted.
It resumes once the commit is accepted and refetches from the position, so every record costs a fetch round trip beside its commit.
The wait for the answer ends on shutdown, and otherwise when librdkafka reports the result of the whole commit operation.
librdkafka waits for that result without a timeout of its own (bundled librdkafka 2.12.1, rdkafka_offset.c:403-406).
It retries a request unanswered for socket.timeout.ms up to two times (rdkafka_buf.c:161-162, rdkafka_proto.h:47, rdkafka_request.c:1784-1787).
A commit with no coordinator waits for one up to session.timeout.ms (rdkafka_cgrp.c:3744-3746).
The wait can therefore run for several socket timeouts before the commit counts as rejected.
A rejected commit is re-offered as under Interval, spaced by the default 500 ms.
The loop therefore polls between re-offers, and a rebalance behind a REBALANCE_IN_PROGRESS can complete.
The assignment stays paused until a re-offer is accepted, so no record is handed out behind a rejected position.
A partition revoked while paused and handed back later arrives paused, because librdkafka keeps a partition's pause flag across a revoke.
The consumer re-applies its intent over the whole assignment on every assign event a pass drains.
A record that comes through recv() needs no such step.
The pause when it is taken and the resume when its commit lands cover the whole assignment.
A returned partition is therefore paused or resumed with the rest, either way.
Two commits on one connection cannot overtake each other.
librdkafka sends a member's OffsetCommit requests to its coordinator over one connection and in order (bundled librdkafka 2.12.1, rdkafka_cgrp.c:4209, rdkafka_broker.c:2763-2771).
The Kafka protocol processes one connection's requests in the order they were sent.
A request unanswered for socket.timeout.ms fails its connection and is retried on the next one (rdkafka_broker.c:1010-1066, rdkafka_request.c:1757-1786, rdkafka_buf.c:392-416).
The broker may still hold that abandoned request, so it can apply it after the retry and after later commits, as under Interval.
That moves the committed position back to the older one, and a crash before the next commit then replays every record since it.
shove does not re-issue the position after such a reconnect: the next record's commit is what repairs it.
The fenced consumer detector keeps its 60 s floor under PerRecord, because there is no interval to scale it by.
A stop that lands while a commit is in flight waits for its answer inside the shutdown deadline.
The shutdown section below says what happens when that answer is late.
The policy needs one prefetch permit, and both register and KafkaConsumer::run refuse a configuration with more.
See the option above.
The pattern is a per-record acknowledgement, what Spring Kafka calls AckMode.RECORD: the position is committed after each record, before the next is taken.
Kafka transactions close the replay window differently: sendOffsetsToTransaction commits the position together with the records a handler produces, and applies to a consume-then-produce handler only.
An idempotent receiver removes the need for either, by making a replayed record a no-op, which is what Interval assumes of a handler.
PerRecord trades throughput for the one-record replay bound where a handler cannot be made idempotent and produces nothing Kafka could tie to the commit.
What the policy does, condition by condition, with the source line as of this change and the test that asserts it:
| Condition | Effect | Source line | Conformance assertion |
|---|---|---|---|
A completion under PerRecord | The position is committed synchronously on a thread, and the next record waits for the coordinator's answer. | src/backends/kafka/consumer.rs:3467, commit_confirmed | per_record_issues_one_offset_commit_per_completion |
| A record is dropped before the handler | Its position is committed like a completion, and the next record waits the same way, through a refused thread or a rejection. | src/backends/kafka/consumer.rs:5543-5556, the pause when a record is taken | a_record_dropped_before_the_handler_holds_the_next_record_until_its_commit_is_accepted, a_record_dropped_before_the_handler_holds_the_next_record_while_its_commit_is_rejected |
A schema registry stall under PerRecord | The stall runs inside the pause the take holds, and its end leaves the resume to the commit, so no record is fetched and put back behind it. | src/backends/kafka/consumer.rs:5616, RegistryStall::inside_pause | a_per_record_stall_leaves_the_resume_to_the_commit |
| The coordinator rejects the commit | The position is re-offered every 500 ms with the assignment paused, and the fence judges the streak at the 60 s floor. | src/backends/kafka/consumer.rs:198, AsyncCommitGate::mark_rejected, and src/backends/kafka/consumer.rs:2895, fence_threshold_over | a_rejected_commit_is_re_offered_until_the_fence_fires, a_rejected_commit_recovers_on_a_re_offer_before_the_next_record_is_handed_out |
| The answer is late | The next record waits for it, and the member keeps polling meanwhile. | src/backends/kafka/consumer.rs:3530, the recv() arm of commit_confirmed | a_late_answer_holds_the_next_record_while_the_member_keeps_polling |
A Retry or Defer republish fails | The member ends with a connection error before the next record is handed out, and the reconnect redelivers the record. | src/backends/kafka/consumer.rs:2024, run_delayed_republish | a_republish_that_fails_under_per_record_ends_the_member_before_the_next_record |
| A partition is handed back paused | A pass that drains the assign pauses or resumes it with the rest, and a record the receive arm takes covers it with its own pause and resume cycle. | src/backends/kafka/consumer.rs:3387, reconcile_pause_after_assign, and src/backends/kafka/consumer.rs:5543-5556 | a_partition_handed_back_paused_is_resumed_when_a_pass_drains_the_assign, a_partition_handed_back_paused_is_resumed_when_the_receive_arm_drains_the_assign |
| A stop finds a commit in flight | It is waited for up to PENDING_COMMIT_BUDGET; an accepted answer is followed by the final commit, a rejected one is re-offered by it, and no answer ends the stop with ShoveError::Commit of the Deadline kind. | src/backends/kafka/consumer.rs:5243-5317, the shutdown arm of KafkaConsumer::run | a_stop_during_a_commit_waits_for_its_answer_and_leaves_the_group_before_returning, a_discard_behind_a_commit_rejected_at_the_stop_is_counted_once_by_the_final_commit, a_per_record_stop_against_a_frozen_coordinator_reports_the_commit_it_could_not_finish |
| A prefetch above one with concurrent processing on | Refused at KafkaConsumer::run and at KafkaConsumerGroupRegistry::register. | src/backends/kafka/consumer.rs:2998, check_per_record_prefetch, and src/backends/kafka/consumer_group.rs:1221 | run_rejects_per_record_commits_with_two_permits, register_rejects_per_record_commits_with_more_than_one_permit |
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.
An in-place delay does not extend it: shutdown cuts the delay short, and the record stays uncommitted.
It then drains the republish queue and commits its final position synchronously.
That commit carries every partition's current position, whether or not an earlier asynchronous commit already offered it.
It therefore offers every safe position, including one an asynchronous commit still in flight may not have landed.
Only the broker's answer confirms them.
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 at most 20 s (SHUTDOWN_COMMIT_DEADLINE in kafka::constants).
That wait covers the commit and the close that follows it on the same thread, and the thread reports the two separately.
A close that outlives the deadline is logged at warn level and is never reported as a failed commit.
Under CommitPolicy::PerRecord a commit in flight when shutdown fires is waited for first, for at most three quarters of that deadline.
If it is accepted, the final commit carries the position again, as every final commit does.
If it is rejected, the final commit re-offers its position.
If it has no answer in time, the stop ends with ShoveError::Commit of the Deadline kind carrying that share, and issues no second commit.
The member then leaves its group once that commit returns on its thread, and a group run counts the member under errors.
The deadline is chosen against the 30 s termination grace Kubernetes gives a Pod.
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.
A final commit the broker does not confirm ends the member with ShoveError::Commit instead of a clean exit.
The error reads offset commit at shutdown, because under CommitPolicy::PerRecord the commit it names can be the one in flight when the stop landed.
The error names the topic, the position the commit carried per partition, and a CommitFailure.
Rejected carries the error the commit returned: the broker's answer, or a librdkafka local error.
A broker's answer says this commit did not land.
A local timeout after the request was sent does not prove that, because the broker may have accepted the commit.
Deadline means nothing answered within the time it carries, normally the 20 s deadline.
NoThread means no thread could be spawned for the commit, so this commit was never made.
Deadline says its result is unknown, because the detached thread may still land it.
Redelivery then starts from the last position the broker accepted for each partition, which an earlier asynchronous commit may have advanced.
The error carries positions and the error text only, never a record's payload, key or headers.
A consumer group counts the error once under SupervisorOutcome::errors, so exit_code() is 1 where it was 0 in 0.15.0.
Nothing else about the run changes: it still ends on its stop signal.
A drain that times out while the commit is still waiting aborts the member instead.
The run then reports timed_out, so exit_code() is 3 and no error is counted.
A member the autoscaler retires on scale-down makes the same final commit.
The rebalance that retirement starts can make the coordinator reject it with a rebalance in progress or a stale generation.
For now that rejection counts under errors like any other.
A member with nothing to commit has no commit to fail and still ends clean.
An acknowledged offset on a partition a rebalance revoked before the stop is not part of the final commit and raises no error.
Its tracker went with the partition, and the new owner picks the record up.
What that means for replay:
- After a crash, the records completed since the last accepted commit are redelivered.
Under
Intervalthat is about one interval's worth when commits are accepted and every earlier offset on the partition had completed. UnderPerRecordit is every record since the commit the coordinator applied last. On a stable coordinator connection that is one record. That record is in the handler, in aRetryorDeferrepublish that had not landed, or in a commit still unanswered or rejected. Nothing is handed out behind an unaccepted position, so no other record joins it. It can be more when a commit was in flight across a coordinator reconnect. The abandoned copy can then be applied after later commits, as the commit section above describes. ARetryorDeferwhose republish fails ends the member with a connection error, and the reconnect redelivers that record. UnderIntervala rejected commit still being re-offered holds the position and widens that window. So does an earlier record whose handler is still running. - 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.
A record the consumer had fetched but not yet handed to a handler is not handled during the drain, and is redelivered too.
When the final commit is rejected or hits its deadline, redelivery starts from the last position the broker accepted, and the member reports that with
ShoveError::Commit. 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-registryBuild 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
SchemaRegistryAuthother thanNone; user:pass@userinfo in the base URL — the HTTP client lifts it out and replays it as anAuthorization: Basicheader, so it is a credential even withSchemaRegistryAuth::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(®istry))
.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(®istry))
.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 reasonschema_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(®istry)))
.await?;
// `publisher.publish::<OrdersTopic>(&order).await?` now emits SR-framed bytes.
Details:
- Subject — defaults to the Confluent TopicNameStrategy
"{topic}-value"; override withKafkaPublisherConfig::with_subject("..."). - Schema id — resolved once per subject via
GET /subjects/{subject}/versions/latestand cached. shove carries no schema text; the subject must already be registered. - Codecs — framing applies only to
JsonCodecandProtobufCodectopics; 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 exceedsmax_message_sizebefore deserialization.deserialize— JSON/codec decode failed; routed to DLQ.timeout— handler exceededhandler_timeout; retried.max_retries_exceeded— retry budget exhausted; routed to DLQ.rejected— handler returnedOutcome::Reject; routed to DLQ.schema_validation- the registry answered and ruled the record out. The death reason isschema_validation_failedwhen the subject is not in the accepted set underEnforcemode. It isschema_resolve_failedwhen the registry does not know the schema id. Routed to DLQ. Requireskafka-schema-registry.schema_frame- the Confluent frame cannot be used, so no lookup is made. The death reason isschema_frame_invalidfor a malformed frame. It isschema_unsupported_codecwhen the topic uses a codec other than JSON or Protobuf. It isschema_message_index_rejectedwhen the protobuf message index is not the onerequire_schema_message_indexpins. Routed to DLQ. Requireskafka-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. Requireskafka-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. Setwith_default_replication_factor(3)on the registry (orwith_replication_factoron the topology declarer) before the first declaration. The default exists so the no-config dev path works; the constant lives inkafka::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.msto 10 seconds andmax.poll.interval.msto 5 minutes (constants inkafka::constants); neither is settable viaKafkaConfig. If handlers are slow, raise the handler timeout withwith_handler_timeout(per group) orwith_default_handler_timeout(registry-wide) instead. - On Windows, the
cmake-buildfeature is required forrdkafka(the underlying C library). This is handled automatically via the[target.'cfg(windows)'.dependencies]entry inCargo.toml. - SASL and TLS require the
kafka-sslfeature. 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, whichkafka-sslselects. Shove exposes OAUTHBEARER only asKafkaSasl::MskIamunderkafka-msk-iam. There is no generic token provider. GSSAPI/Kerberos needs Cyrus SASL, which the separatekafka-gssapifeature links for directrdkafkause.KafkaSaslexposes no GSSAPI variant, so that feature gates no shove code and shove itself cannot authenticate with Kerberos. For AWS MSK IAM, use thekafka-msk-iamfeature instead. - For publish failures on missing topics, see Declare topology.
librdkafkasystem dependencies on Debian/Ubuntu CI runners:libsasl2-devis required only withkafka-gssapi.kafka-sslalone links OpenSSL and no SASL library, and the shoveci.ymlasserts that withcargo tree.
Examples
- Basic — publish/consume round trip with hold queues and DLQ
- Sequenced — partition-key ordering
- Audited Consumer —
MessageHandlerExt::auditedwrapping - 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::pinginto a k8s health endpoint. - Broadcast - per-instance fan-out via a groupless
assign(), leaving__consumer_offsetsuntouched. It starts at the tail by default, or at the head or a timestamp withwith_broadcast_start.