KRaft & Modern Kafka

Understand KRaft mode: Kafka 4.x is post-ZooKeeper. New clusters should be KRaft-first. Also understand tiered storage splitting local hot and remote cold storage like S3.

Last generated

Lesson 4 of 10 available15 practice questions

SPACED REPETITION Β· 15 practice questions

Make this lesson stick.

Try 3 questions now. No account needed. Sample answers aren't saved.

By the end of this lesson you will be able to:

  • explain what KRaft replaced and why your producer and consumer code does not change;
  • stand up a single-node KRaft cluster for local and CI work, and avoid the advertised-listener trap;
  • say what DescribeClusterAsync() can and cannot tell you about a KRaft cluster;
  • size a controller quorum and estimate how long a failover takes, and budget client timeouts for it;
  • explain tiered storage β€” hot local disk, cold remote storage β€” and size local disk correctly with and without it.

Prerequisites: "Core Kafka Mental Model" (brokers, replication, ISR) and "Topics, Partitions & Keys". Most major sections include a short exercise β€” try it before opening the answer.

Watch the Lesson Video32 min
The whole lesson, narrated: what KRaft replaced, the one door into the cluster, a local KRaft cluster and the advertised-listener trap, clients that ride out a failover, tiered storage, and five common mistakes

Why Modern Kafka Changed the Rules for .NET Developers

If you've ever run docker-compose up on a Kafka project and watched two separate containers fight for memory before your producer could send a single message, you've felt the tax that ZooKeeper used to charge every Kafka deployment. For most of Kafka's history, a Kafka cluster was never really "a cluster" β€” it was two distributed systems duct-taped together, and .NET developers writing Confluent.Kafka code had to know at least a little about both. Why did the Kafka project decide an entire second coordination system was expendable? What actually breaks β€” and what stays exactly the same β€” in your C# code when that system disappears? And why does a change that sounds purely operational end up touching your docker-compose.yml, your CI pipeline, and your IAdminClient calls (Confluent.Kafka's interface for cluster-management operations such as creating topics)? This lesson answers those questions by walking through KRaft (Kafka Raft), the ZooKeeper-free architecture that is now the only way to run a new Kafka cluster, and by tracing exactly where that shift shows up in .NET application and tooling code.

The Two-System Problem ZooKeeper Created

In the original Kafka architecture, brokers didn't just talk to each other β€” they depended on an external coordination service, Apache ZooKeeper, for three critical jobs: electing a controller broker (responsible for managing partition leadership and cluster metadata), tracking which brokers were alive and registered, and storing the metadata describing every topic and partition. ZooKeeper is a general-purpose distributed coordination service, not something purpose-built for Kafka, which meant every Kafka deployment was really an exercise in running and tuning two distributed systems side by side.

That had real costs. ZooKeeper needed its own quorum of nodes (a group that decides by majority vote), its own JVM tuning, its own monitoring, and its own upgrade cadence, often out of sync with the broker upgrade cycle. Controller failover went through a ZooKeeper-mediated election involving watches (change notifications) and ephemeral znodes (ZooKeeper data entries that vanish when their owner's session ends) β€” a mechanism that worked, but added a layer of indirection between "a broker died" and "a new controller is active and metadata is consistent again." For a .NET developer this rarely showed up in producer or consumer code, but it showed up in the surrounding ecosystem: local dev tooling that shelled out to zkCli.sh, health checks that pinged ZooKeeper's client port, and runbooks with a section titled "if ZooKeeper is unhealthy, do not touch the brokers yet."

KRaft: One System Instead of Two

KRaft replaces ZooKeeper's job with a built-in Raft-based metadata quorum made up of Kafka nodes themselves. Instead of an external system tracking cluster metadata in znodes, the metadata lives in a replicated log inside Kafka, managed by the Raft consensus protocol (a majority of nodes elects one leader, which appends every change to a log the others copy). As of the Kafka 4.0 release, ZooKeeper support has been removed entirely β€” KRaft is no longer an opt-in alternative for the adventurous, it is the only supported mode for standing up a new cluster. If you're provisioning a fresh cluster today, whether that's a local container for integration tests or a production deployment, there is no ZooKeeper configuration to reach for.

🎯 Key Principle: KRaft doesn't change what Kafka's metadata is (which broker leads which partition, which topics exist, which configs apply) β€” it changes how that metadata is agreed upon and stored, replacing an external coordination service with a Raft log native to Kafka itself.

ZooKeeper-based cluster (removed in Kafka 4.0):
 Kafka Broker 1 ↔ ZooKeeper ensemble ↔ Kafka Broker 2
 Controller election, broker registration, and topic metadata
 all live in ZooKeeper znodes, external to the brokers.

KRaft-based cluster (the only supported mode going forward):
 Controllers 1, 2, 3  ── Raft metadata log (one active controller,
                          two hot standbys; a majority must be up)
        ↓ brokers replicate the metadata log
 Broker 1   Broker 2   Broker 3
 Controller election and metadata live in a Kafka-native
 replicated log; no external coordination service exists.
 (In small setups one process can play both roles.)

The diagram shows dedicated controller and broker nodes; a combined node that plays both roles is the third option. The process.roles and controller quorum settings that make this concrete for a docker-compose.yml are covered in "process.roles and the Controller Quorum" and put to work in "Spinning Up and Targeting a KRaft Cluster from a .NET Project," and how the Raft quorum elects a new leader β€” and how long that takes β€” is sketched in "Faster Failover Doesn't Mean No Failover" below.

What Actually Changes for a .NET Developer

Here's the part that matters most for your day job: the shift is, in one sense, remarkably boring. Open a Confluent.Kafka producer or consumer loop and nothing about the message-passing code changes. Producing a message, consuming a batch, committing an offset β€” all of that talks to Kafka brokers over the Kafka wire protocol, and it never talked to ZooKeeper directly even in the old architecture. .NET client libraries were always insulated from ZooKeeper by design; only the brokers and a narrow set of admin-adjacent tooling touched it.

What does change is everything around that core β€” the seams where your .NET project meets cluster infrastructure:

πŸ”§ Cluster bootstrap β€” how a cluster is initialized and assigned an identity now happens through a Kafka-native storage format step rather than auto-assignment via ZooKeeper. πŸ”§ The Admin API as the only window β€” questions that ZooKeeper-era scripts answered by reading znodes (which cluster is this, which brokers are in it) are now answered through the Kafka Admin API, e.g. DescribeClusterAsync(). (The call itself isn't a KRaft invention: metadata responses carried a cluster ID on ZooKeeper clusters too, and the .NET method arrived in Confluent.Kafka 2.3.) πŸ”§ Local dev tooling β€” Docker Compose files and Testcontainers definitions (Testcontainers is a library that starts Docker containers from test code) need different environment variables and a different startup sequence than a ZooKeeper-plus-broker pair. πŸ”§ Operational runbooks β€” failure modes, health checks, and recovery steps need to reflect a single coordinated system instead of two independently-monitored ones.

A minimal way to see this insulation is to look at how little a basic .NET consumer cares. The following is identical whether the brokers behind it run KRaft or, on an older deployment, still coordinate through ZooKeeper β€” which is exactly the point:

using Confluent.Kafka;

var config = new ConsumerConfig
{
    // Only the Kafka protocol endpoint matters here.
    // There has never been a ZooKeeper setting on this config object.
    BootstrapServers = "localhost:9092",
    GroupId = "orders-processor",
    AutoOffsetReset = AutoOffsetReset.Earliest
};

using var consumer = new ConsumerBuilder<string, string>(config).Build();
consumer.Subscribe("orders");

using var cts = new CancellationTokenSource();
Console.CancelKeyPress += (_, e) => { e.Cancel = true; cts.Cancel(); };

try
{
    while (!cts.IsCancellationRequested)
    {
        var result = consumer.Consume(cts.Token);
        Console.WriteLine($"Received: {result.Message.Value}");
    }
}
catch (OperationCanceledException) { /* graceful shutdown */ }
finally
{
    consumer.Close();
}

What is new is where you go for answers that used to come from ZooKeeper-side tooling. Confirming which cluster you are connected to, and which brokers it has, is something your .NET code asks the cluster directly:

using Confluent.Kafka;
using Confluent.Kafka.Admin;

var adminConfig = new AdminClientConfig { BootstrapServers = "localhost:9092" };
using var admin = new AdminClientBuilder(adminConfig).Build();

// DescribeClusterAsync reaches the cluster through the Kafka protocol itself β€”
// no external ZooKeeper client required. Note it is async-only; there is no
// synchronous DescribeCluster overload on IAdminClient.
var cluster = await admin.DescribeClusterAsync(
    new DescribeClusterOptions { RequestTimeout = TimeSpan.FromSeconds(5) });

Console.WriteLine($"Cluster ID: {cluster.ClusterId}");
Console.WriteLine($"Brokers:    {cluster.Nodes.Count}");

This is a preview, not the full treatment β€” a worked end-to-end producer/consumer/admin example against a live KRaft container is the focus of "Spinning Up and Targeting a KRaft Cluster from a .NET Project". The point here is narrower: operations that used to require stepping outside the Kafka client entirely, into ZooKeeper CLI tools, now go through the Kafka Admin API, mostly via the same IAdminClient your application already uses.

πŸ’‘ Mental Model: Think of the ZooKeeper-to-KRaft transition less as "Kafka got a new feature" and more as "Kafka internalized a dependency it used to outsource." The job didn't disappear, it moved inside the system you're already talking to β€” which is exactly why your producer and consumer code doesn't need to change, but your infrastructure code does.

What This Lesson Covers, and What It Deliberately Doesn't

This lesson stays focused on the .NET-facing consequences: the configuration keys that replaced ZooKeeper settings, how to point local and CI tooling at a KRaft cluster, resilience patterns for AdminClient calls during topology changes, and the mistakes teams carry over from ZooKeeper-era habits.

Two related topics get a lighter treatment. The internals of the Raft-based controller quorum β€” voters (the controllers that elect a leader and commit metadata) and observers (nodes such as brokers that only follow the log), the metadata log's replication, dynamic membership β€” matter more to operators than to application developers; this lesson covers only the parts with client-visible consequences (quorum size and failover timing), and the Kafka operations documentation covers the rest. Separately, tiered storage, which splits a topic's data between fast local disks and cheaper remote object storage, is a related-but-independent modernization; it changes how retention and disk sizing work but not how your producer or consumer code is written. It has its own section, "Tiered Storage: Hot Local Disk, Cold Remote Storage," later in this lesson.

⚠️ Common Mistake: Treating "ZooKeeper is gone" and "tiered storage exists" as the same upgrade, since both get bundled into the phrase "modern Kafka." They're independent: a cluster can run KRaft with no tiered storage configured at all, and the disk-filling consequences of assuming tiered storage is on are covered later in "Common Mistakes .NET Teams Make with Modern Kafka Clusters."

✍️ Exercise: what has to change?

Your team's Kafka cluster is being moved from ZooKeeper mode to a KRaft cluster on Kafka 4.x. For each item, say whether it must change, and why.

  1. The ConsumerConfig in your order service (BootstrapServers, GroupId, AutoOffsetReset).
  2. The ProduceAsync calls in your API.
  3. The docker-compose.yml your integration tests use (a zookeeper service plus a kafka service with KAFKA_ZOOKEEPER_CONNECT).
  4. A Kubernetes readiness probe on TCP port 2181.
  5. A CI script that runs kafka-topics.sh --zookeeper zk1:2181 --create ....
Check your answer
  1. No change (beyond the address if the cluster moves). Clients only ever spoke the Kafka protocol to bootstrap.servers.
  2. No change. Same reason.
  3. Must change. The ZooKeeper service goes, and the Kafka service needs KRaft settings instead: KAFKA_PROCESS_ROLES, a node ID, a controller listener (a named endpoint the node opens for controller traffic) and the quorum setting.
  4. Must change. Port 2181 is ZooKeeper's; nothing listens there any more, so the probe will keep the pod unready forever.
  5. Must change. The --zookeeper flag no longer exists; use --bootstrap-server (the Admin API), or do it from .NET with CreateTopicsAsync.

The pattern: code that talks to Kafka is untouched; infrastructure that talked around Kafka to ZooKeeper has to go.

From ZooKeeper Habits to KRaft Reality: What Client Code Needs to Know

If you've maintained a .NET service talking to Kafka for a few years, some of your instincts were formed by ZooKeeper's presence β€” even if your application code never called it. Config files, compose definitions, admin scripts, and mental models about "where cluster state lives" all quietly assumed a second distributed system behind the brokers. KRaft removes that assumption, and the differences show up in exactly the places a .NET developer touches.

Connection Strings: bootstrap.servers Is the Only Door

Under the ZooKeeper-era architecture, two connection strings coexisted in most Kafka projects. Producers and consumers used bootstrap.servers to reach brokers, but administrative tooling frequently used zookeeper.connect β€” pointing at the ZooKeeper ensemble (ZooKeeper's own cluster of servers) instead. Tools like zkCli.sh, and broker-side scripts invoked with --zookeeper, could read and write cluster metadata by talking to ZooKeeper's znode tree directly, bypassing the Kafka protocol altogether.

That second door is gone. In a KRaft cluster, zookeeper.connect is not deprecated-but-tolerated β€” it has no meaning at all, because there is no ensemble to connect to. Every operation goes through the Kafka protocol against bootstrap.servers. It doesn't just mean "use a different connection string": an entire class of tooling that read cluster state by inspecting znodes has no equivalent path anymore and must be rewritten against the Kafka Admin API.

// ZooKeeper-era .NET tooling often needed BOTH of these,
// because some operations only existed via ZK-aware scripts:
// var zkConnect = "zk1:2181,zk2:2181,zk3:2181"; // no longer applicable

// KRaft-mode configuration: only the Kafka protocol endpoint matters
var adminConfig = new AdminClientConfig
{
    BootstrapServers = "broker1:9092,broker2:9092,broker3:9092"
};

using var admin = new AdminClientBuilder(adminConfig).Build();

This looks almost too simple to be the point β€” and that's the point. AdminClientConfig carries no ZooKeeper-related properties at all, because Confluent.Kafka's admin client has only ever spoken the Kafka wire protocol. What changed is that the ZK-aware back door is gone: the Kafka Admin API β€” and the CLI tools built on it β€” is now the only way in. Not all of that API is exposed by Confluent.Kafka's IAdminClient, though: partition reassignment and controller-quorum management have no .NET method, so those stay with kafka-reassign-partitions.sh and kafka-metadata-quorum.sh.

Where Cluster Identity Comes From Now

In the ZooKeeper architecture, cluster identity was assigned automatically: the first broker to start registered itself, a cluster ID was generated and stored in ZooKeeper, and other brokers discovered it there. There was no explicit provisioning step β€” the cluster "became" a cluster the moment brokers connected to a shared ensemble.

KRaft inverts this. Before a node can start in KRaft mode, its storage directory must be formatted with kafka-storage.sh format, supplying a cluster ID generated once (typically with kafka-storage.sh random-uuid) and reused across every node in that cluster. This is a deliberate, one-time provisioning act rather than an emergent side effect of nodes finding each other. Skip it and the node refuses to start rather than silently forming an ad hoc cluster.

πŸ’‘ Pro Tip: You will not see this step in most Docker Compose files, and that's not an omission. The official apache/kafka image formats storage for you when it starts. If you don't set a CLUSTER_ID environment variable, its startup scripts fall back to a fixed default cluster ID baked into the image (the scripts suggest replacing it with one from kafka-storage.sh random-uuid). Two consequences: every local container you start without CLUSTER_ID shares the same ID, so an ID check can't tell them apart; and a container that loses its volume comes back with the same ID but empty data β€” a vanished consumer group (the consumers sharing a GroupId, whose committed offsets the cluster stores) after a restart means lost data, not a new cluster.

Anything that used to confirm "which cluster am I talking to" by reading a znode under /cluster/id has no ZooKeeper tree to inspect. The Kafka protocol exposes this directly:

// Confirm cluster identity and membership via the Kafka protocol,
// replacing tooling that used to read ZooKeeper's /cluster/id znode.
// DescribeClusterAsync is async-only β€” there is no synchronous overload.
var description = await admin.DescribeClusterAsync(
    new DescribeClusterOptions { RequestTimeout = TimeSpan.FromSeconds(10) });

Console.WriteLine($"Cluster ID: {description.ClusterId}");
foreach (var node in description.Nodes)
{
    Console.WriteLine($"Node {node.Id}: {node.Host}:{node.Port}");
}

DescribeClusterAsync() returns the cluster ID baked in at format time and the list of brokers β€” all over the same protocol your producers and consumers already use.

⚠️ The Controller property is a trap under KRaft. The result also has a Controller node, and ZooKeeper-era habit says it names the active controller. On a KRaft cluster it doesn't: brokers fill that field with a randomly chosen live broker (Kafka uses it only to route admin requests, which brokers then forward to the real controller). With dedicated controller nodes, the actual controllers never appear in Nodes at all. To see the quorum's active controller, use the CLI: kafka-metadata-quorum.sh --bootstrap-server broker1:9092 describe --status and read LeaderId. Confluent.Kafka has no .NET equivalent of that call.

⚠️ Common Mistake: Reaching for GetMetadata() and looking for a cluster ID on the result. Metadata exposes Brokers, Topics, and OriginatingBrokerId/OriginatingBrokerName β€” the last of which identifies the broker that answered your request, not the cluster. DescribeClusterAsync() is the only call that returns ClusterId.

process.roles and the Controller Quorum: Why This Matters for Your docker-compose.yml

ZooKeeper-era Kafka had a clean process-level separation: ZooKeeper processes handled coordination, Kafka broker processes handled data. KRaft collapses coordination into Kafka itself, but does so by introducing an explicit role each Kafka process plays, set via process.roles. A node can start with process.roles=broker, process.roles=controller, or β€” common in small and local deployments β€” process.roles=broker,controller, meaning one JVM does both jobs.

Alongside that, a node needs to know how to find the controller quorum. Two settings can do this:

βš™οΈ Setting🎯 ModelπŸ“Œ When to use
controller.quorum.votersStatic β€” every voter listed as nodeId@host:portFixed quorum membership; simplest for local dev. Works while the cluster's kraft.version (its KRaft feature level) is 0; must be removed once the cluster moves to level 1, which enables dynamic membership and deprecates this setting
controller.quorum.bootstrap.serversDynamic β€” an entry point, membership managed at runtimeKafka 3.9+ clusters using dynamic quorum membership, where controllers can be added or removed without a full-cluster config change; the form the current Kafka docs use for every node

Static voters are what most examples and local setups still use, and they still work in Kafka 4.x β€” the examples in this lesson use them for that reason. The Kafka docs, however, say that kraft.version=1 deprecates controller.quorum.voters and that it must be removed once a cluster runs that version, and converting an existing static quorum to a dynamic one is supported from Kafka 4.1. For a new production cluster, start with controller.quorum.bootstrap.servers; the kafka-metadata-quorum.sh add-controller / remove-controller workflow is described in the Kafka operations docs.

Either way, this is the setting that occupies the space zookeeper.connect used to fill:

# docker-compose.yml excerpt: a single combined broker+controller node,
# typical for local .NET integration test environments
services:
  kafka:
    image: apache/kafka:4.1.0
    environment:
      KAFKA_NODE_ID: 1
      KAFKA_PROCESS_ROLES: broker,controller
      KAFKA_CONTROLLER_QUORUM_VOTERS: "1@kafka:9093"
      KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
      # CONTROLLER is a listener NAME, not a protocol. With only PLAINTEXT in
      # use Kafka can infer this mapping, but anything else (SSL, SASL) makes it
      # mandatory β€” so write it down explicitly.
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT
      # One node can't hold the default 3 replicas of the consumer-offsets topic
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
    ports:
      - "9092:9092"

The environment variables that didn't exist in a ZooKeeper-based compose file β€” KAFKA_PROCESS_ROLES, KAFKA_NODE_ID, the CONTROLLER listener with KAFKA_CONTROLLER_LISTENER_NAMES, and the quorum setting β€” are exactly the ones replacing what ZooKeeper used to provide implicitly.

Metadata Propagation: Fewer Stale-Metadata Surprises

One behavioral difference .NET developers notice without changing a line of code is how quickly the cluster's view of partition leadership becomes consistent after a change. Under ZooKeeper, metadata changes were written to ZooKeeper; the controller broker watched ZooKeeper and pushed UpdateMetadata requests out to every other broker. That worked, but introduced multiple hops and windows where a client's cached metadata didn't yet reflect reality, producing NOT_LEADER_OR_FOLLOWER and similar errors on the next produce or fetch.

KRaft replaces that write-to-ZooKeeper-then-push path with a single Raft-replicated metadata log that all brokers consume directly. Instead of a controller reading from ZooKeeper and republishing, every broker tails the same ordered log of metadata records. One authoritative log rather than a coordination store plus a relay step means a shorter, more uniform propagation path, and disagreements between what a broker believes and what the controller just committed resolve faster.

The practical upshot: you may see fewer stale-metadata errors following a leadership change, and error windows around controller transitions tend to close more quickly. This does not mean such errors disappear β€” a client can still hold cached metadata momentarily behind the cluster's true state, and handling that is still the job of retry and timeout configuration, covered in "Writing Resilient .NET Clients for KRaft-Based Clusters." What changes is the frequency and duration of the condition, not its existence.

⚠️ Common Mistake: Assuming that because KRaft's metadata propagation is faster, a .NET service no longer needs sensible metadata-refresh or retry handling. Faster convergence reduces the odds of hitting a stale-metadata window during routine leadership changes; it does not eliminate transient errors or slow admin calls around an active controller election, which is exactly the edge case explored later.

✍️ Exercise: the alert that never stops

A platform team adds a CI health check to a production cluster with three brokers and three dedicated controller nodes:

var d = await admin.DescribeClusterAsync();
if (d.Controller?.Id != lastKnownControllerId)
    Alert($"Controller moved to {d.Controller?.Id}");

It fires on about two runs in three, although nothing is failing. It also never reports node IDs 100–102, which are the controller nodes. Explain both observations, and say what the check should use instead.

Check your answer

On a KRaft cluster the Controller field in the DescribeCluster (and Metadata) response is a randomly chosen live broker, not the active quorum controller β€” so it changes from call to call while nothing is wrong. Controller nodes never appear in it, or in Nodes, because clients only ever talk to brokers; brokers forward controller-bound requests.

To watch the real active controller, run kafka-metadata-quorum.sh --bootstrap-server ... describe --status (its LeaderId field) from the CI job, or use the controllers' own metrics. Keep DescribeClusterAsync() for what it is good at: the cluster ID and the broker list.

Spinning Up and Targeting a KRaft Cluster from a .NET Project

Once you understand that a KRaft cluster has no ZooKeeper ensemble behind it, the practical question is how to stand one up locally and how a .NET producer, consumer, or admin client finds it. The client-side wiring barely changes. The work that does change is entirely in configuring the broker container, because a KRaft node now has to know things ZooKeeper used to track on its behalf β€” its identity, its role, and who else votes on cluster metadata.

Configuring a Combined Broker+Controller Node for Local Development

For integration tests you don't need a multi-node quorum β€” a single node acting as both broker and controller via process.roles=broker,controller is the standard pattern. It gives you no controller fault tolerance at all: fault tolerance comes from having three or five voters, whatever roles they play. The Kafka docs advise against combined mode in critical deployments, so production clusters normally run those voters as dedicated controller nodes β€” for isolation, so a busy broker can't starve the controller and the two can be rolled and sized independently.

A combined node needs three pieces of identity ZooKeeper used to hand out: a node ID, a cluster ID (stamped into the storage directory at format time β€” the apache/kafka image does the formatting for you and uses its fixed default ID unless you set CLUSTER_ID), and the controller quorum setting listing which node IDs may vote, with their controller-listener addresses.

services:
  kafka:
    image: apache/kafka:4.1.0
    container_name: kraft-broker
    ports:
      - "9092:9092"
    environment:
      KAFKA_NODE_ID: 1
      KAFKA_PROCESS_ROLES: broker,controller
      KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
      KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT
      # Single-node cluster: no replication headroom, local dev/test only
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1

The format 1@kafka:9093 means "node ID 1, reachable at host kafka on port 9093 for controller traffic." Note that the client-facing listener (PLAINTEXT on 9092) and the controller listener (CONTROLLER on 9093) are deliberately separate β€” application traffic and Raft consensus traffic never share a listener.

Testcontainers, and the Advertised-Listener Trap

⚠️ Common Mistake β€” the one that costs an afternoon. KAFKA_ADVERTISED_LISTENERS is the address the broker tells clients to use after bootstrap. If you bind Kafka's port to a random host port (WithPortBinding(9092, true), the usual Testcontainers practice for avoiding collisions) but leave advertised listeners hardcoded to localhost:9092, your client connects to the random port, receives localhost:9092 back in the metadata response, and then fails to reach the broker. The symptom looks like a network or firewall problem and has nothing to do with either.

The cleanest fix is to use the dedicated Testcontainers.Kafka package, whose KafkaBuilder rewrites advertised listeners to match the mapped port for you:

using Testcontainers.Kafka;

// KafkaBuilder handles the advertised-listener rewrite and KRaft startup.
// Pin a real tag rather than "latest", so a new image can't change your tests
// overnight; the module detects the vendor (apache or confluentinc) from the name.
var kafkaContainer = new KafkaBuilder("apache/kafka:4.1.0")
    .Build();

await kafkaContainer.StartAsync();
var bootstrapServers = kafkaContainer.GetBootstrapAddress();

If you need the raw ContainerBuilder β€” say, to lay out the listeners yourself, since KafkaBuilder writes its own (any other KRaft variable can go on it with WithEnvironment) β€” bind a fixed host port so the advertised address stays truthful:

using DotNet.Testcontainers.Builders;

var kafkaContainer = new ContainerBuilder("apache/kafka:4.1.0")
    .WithPortBinding(9092, 9092)          // fixed, NOT random β€” must match advertised
    .WithEnvironment("KAFKA_NODE_ID", "1")
    .WithEnvironment("KAFKA_PROCESS_ROLES", "broker,controller")
    .WithEnvironment("KAFKA_LISTENERS", "PLAINTEXT://:9092,CONTROLLER://:9093")
    .WithEnvironment("KAFKA_ADVERTISED_LISTENERS", "PLAINTEXT://localhost:9092")
    .WithEnvironment("KAFKA_CONTROLLER_LISTENER_NAMES", "CONTROLLER")
    .WithEnvironment("KAFKA_CONTROLLER_QUORUM_VOTERS", "1@localhost:9093")
    .WithEnvironment("KAFKA_LISTENER_SECURITY_PROTOCOL_MAP",
        "PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT")
    .WithEnvironment("KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR", "1") // one node
    .WithWaitStrategy(Wait.ForUnixContainer().UntilInternalTcpPortIsAvailable(9092))
    .Build();

await kafkaContainer.StartAsync();
var bootstrapServers = "localhost:9092";

A fixed port means parallel test runs on one machine will collide, which is the trade you're accepting for a simpler listener configuration.

The wait strategy above checks only that port 9092 is open, and that is a fair readiness signal: Kafka 4.x opens that listener only after the broker has caught up with the metadata log and the controller has accepted it as live, so it can serve requests. What can still take a moment is work done on first use β€” the first consumer group makes the broker create the consumer-offsets topic. librdkafka retries through that on its own, so give the first Consume several seconds (the example below allows 10) rather than adding a sleep.

Producing and Consuming Against a KRaft Cluster

This is the reassuring part: ProducerConfig and ConsumerConfig need only BootstrapServers to find the cluster (plus GroupId for a consumer that subscribes). There is no zookeeper.connect equivalent and no separate discovery step β€” the client asks whatever broker it can reach for current cluster metadata, and that broker answers using information it learned from the Raft metadata log instead of from the ZooKeeper-era controller's pushed updates. From the client's point of view this is the same bootstrap-and-discover flow Kafka has always used.

using Confluent.Kafka;

var bootstrapServers = "localhost:9092"; // from Testcontainers or docker-compose

// Produce a single message
var producerConfig = new ProducerConfig { BootstrapServers = bootstrapServers };
using (var producer = new ProducerBuilder<string, string>(producerConfig).Build())
{
    var result = await producer.ProduceAsync("orders",
        new Message<string, string> { Key = "order-42", Value = "placed" });
    Console.WriteLine($"Produced to {result.TopicPartitionOffset}");
}

// Consume it back
var consumerConfig = new ConsumerConfig
{
    BootstrapServers = bootstrapServers,
    GroupId = "orders-test-group",
    AutoOffsetReset = AutoOffsetReset.Earliest
};
using (var consumer = new ConsumerBuilder<string, string>(consumerConfig).Build())
{
    consumer.Subscribe("orders");
    var consumeResult = consumer.Consume(TimeSpan.FromSeconds(10));
    Console.WriteLine($"Consumed: {consumeResult?.Message.Value}");
    consumer.Close();
}

Nothing here is KRaft-specific β€” precisely the point. What KRaft changes is what happens before this code runs and what tooling you use to inspect cluster state.

Verifying Cluster State with DescribeClusterAsync()

Once the container is running, confirm β€” rather than assume β€” that your client sees the cluster as expected. DescribeClusterAsync() returns the broker list and the cluster ID. This matters more under KRaft than before, because there's no external zkCli session to peek at; the AdminClient call is your main window into cluster metadata now.

using Confluent.Kafka;
using Confluent.Kafka.Admin;

var adminConfig = new AdminClientConfig { BootstrapServers = bootstrapServers };
using var admin = new AdminClientBuilder(adminConfig).Build();

var clusterInfo = await admin.DescribeClusterAsync(
    new DescribeClusterOptions { RequestTimeout = TimeSpan.FromSeconds(5) });

Console.WriteLine($"Cluster ID: {clusterInfo.ClusterId}");
foreach (var node in clusterInfo.Nodes)
{
    Console.WriteLine($"Node {node.Id}: {node.Host}:{node.Port}");
}

Against the single-node combined setup above, Nodes should contain exactly one entry. Scale up later and this same call is how a CI health check or diagnostic script confirms the cluster ID and that every broker it expects is alive β€” a broker the controller has marked as dead (fenced, covered below) is left out of Nodes. It is not how you find the active controller β€” as covered in "Where Cluster Identity Comes From Now," the Controller field is a random live broker on KRaft clusters; use kafka-metadata-quorum.sh ... describe --status for that.

Checking Client/Broker API Version Compatibility

Confluent.Kafka wraps librdkafka, which learns the broker's supported API version ranges on connect and picks the highest version both sides support β€” this is how the two agree on request/response formats. ⚠️ Common Mistake: pinning an old Confluent.Kafka NuGet package (and the librdkafka build bundled with it) and forgetting it while the brokers move on. KRaft itself asks nothing new of clients β€” brokers forward admin requests to the controller transparently β€” but Kafka 4.0 also removed a set of old protocol API versions (KIP-896; a KIP is a numbered Kafka Improvement Proposal). A client old enough to know only removed versions of some API can't use that API. librdkafka never sends a version the broker doesn't offer, so instead of the broker's UNSUPPORTED_VERSION the failure is local β€” for an admin call, ErrorCode.Local_UnsupportedFeature ("… not supported by broker"). An outdated client may also mishandle newer error codes, producing errors that look unrelated to versioning.

Inspect what actually got negotiated rather than guessing:

var debugConfig = new AdminClientConfig
{
    BootstrapServers = bootstrapServers,
    Debug = "broker,feature"
};
using var debugAdmin = new AdminClientBuilder(debugConfig)
    .SetLogHandler((_, message) => Console.WriteLine($"[{message.Level}] {message.Message}"))
    .Build();

await debugAdmin.DescribeClusterAsync(
    new DescribeClusterOptions { RequestTimeout = TimeSpan.FromSeconds(5) });

The feature debug context logs librdkafka's "Broker API support" lines: for each API, the version range the broker advertised in its ApiVersions response. (Only the broker advertises; librdkafka then picks the highest version that both it and the broker support.) If an API your code uses has no overlap with what your client supports, the NuGet pin is the thing to change β€” a stale package is a far more common root cause than a genuinely incompatible current broker.

One Connection Configuration, Two Environments

Because BootstrapServers is the only setting that says where the cluster is, keep it and any security settings entirely in configuration, never hardcoded, so one compiled binary can point at a local KRaft container in dev and a managed cluster in production.

{
  "Kafka": {
    "BootstrapServers": "localhost:9092",
    "SecurityProtocol": "Plaintext"
  }
}
public class KafkaOptions
{
    public string BootstrapServers { get; set; } = string.Empty;
    public string SecurityProtocol { get; set; } = "Plaintext";
}

// Program.cs / composition root
builder.Services.Configure<KafkaOptions>(
    builder.Configuration.GetSection("Kafka"));

builder.Services.AddSingleton<IProducer<string, string>>(sp =>
{
    var options = sp.GetRequiredService<IOptions<KafkaOptions>>().Value;
    var config = new ProducerConfig
    {
        BootstrapServers = options.BootstrapServers,
        SecurityProtocol = Enum.Parse<SecurityProtocol>(options.SecurityProtocol)
    };
    return new ProducerBuilder<string, string>(config).Build();
});

In production, an environment variable or secrets store overrides Kafka:BootstrapServers and typically flips SecurityProtocol to SaslSsl. SASL also needs credentials, so give KafkaOptions SaslMechanism, SaslUsername and SaslPassword properties and map them the same way; after that, the code building the producer never changes between environments.

✍️ Exercise: why can't the test connect?

A test uses the raw ContainerBuilder with .WithPortBinding(9092, true) (random host port) and keeps KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092. Docker maps container port 9092 to host port 55001, and the test's producer uses BootstrapServers = "localhost:55001".

  1. Walk through what the client does, step by step, and say where it fails.
  2. Give two fixes.
Check your answer
  1. The client connects to localhost:55001 and the bootstrap succeeds β€” it reaches the broker. It asks for metadata, and the broker answers "I am localhost:9092", because that is its advertised listener. From then on the client talks to localhost:9092, where nothing on the host is listening (the container port is mapped to 55001). Produce requests time out, and it looks like a network or firewall problem.
  2. Use Testcontainers.Kafka's KafkaBuilder, which sets the advertised listener to the mapped port for you; or bind a fixed port (.WithPortBinding(9092, 9092)) so the advertised address is true β€” accepting that parallel test runs on one machine will then collide.

Writing Resilient .NET Clients for KRaft-Based Clusters

A cluster that fails over its controller in a few seconds sounds like it should make client-side resilience less necessary. In practice the opposite happens: because the failure window is short, it's tempting to skip defensive tuning entirely, and then a producer with a short delivery budget fails its sends when a broker dies, or a deployment pipeline breaks the first time the controller quorum stays leaderless for longer than a normal election. The goal here is to make that window survivable by construction.

Faster Failover Doesn't Mean No Failover

Under ZooKeeper coordination, a controller failover after a crash could take around 20 s or more: ZooKeeper first had to expire the dead controller's session (18 s by default), and the new controller then re-read broker and partition state from ZooKeeper before acting, which took longer the more partitions the cluster had. KRaft's Raft-based quorum is quicker, because every standby controller already has the metadata log locally and doesn't need to reload state. The client-facing consequence: admin calls see shorter stalls around controller changes, not zero-length ones. Producers and consumers stall only while a partition leader has to move (the broker-crash case below), which is also faster than under ZooKeeper.

What "quicker" means with default settings:

  • Controller crash. Standby controllers notice only after controller.quorum.fetch.timeout.ms (2 s) without a successful fetch from the leader. The election that follows usually takes a few round trips; if the vote splits, candidates wait a randomized election timeout (1–2 s with the default controller.quorum.election.timeout.ms of 1 s) and try again. Budget roughly 2–4 s.
  • Controller shut down cleanly. The active controller resigns and a standby takes over almost at once β€” usually in under a second.
  • A broker crash β€” what producers actually feel. When a broker that leads partitions dies, the controller fences it (marks it dead, so it can no longer lead partitions) after broker.session.timeout.ms (9 s) without a heartbeat, then moves partition leadership. Until then, produces and fetches for those partitions fail with leader errors and are retried. A cleanly shut-down broker hands its leadership over first, so rolling restarts are much gentler.

A shorter window still needs a plan. librdkafka's timeouts are generous by default β€” tuned so that a transient problem doesn't fail a request, at the cost of taking a long time to surface a real one. Against a cluster that recovers within seconds, several of them are worth tuning β€” the timeouts tightened so your client reports real failures sooner, while still covering a normal failover:

πŸ”§ SettingπŸ“¦ librdkafka defaultβœ… SuggestedπŸ’¬ Why
message.timeout.ms
(MessageTimeoutMs)
300000 (5 min)30000Total delivery budget per message; 30 s still covers a crashed broker being fenced (~9 s) with margin, while 5 minutes delays real failures
socket.connection.setup.timeout.ms3000010000Move on to another broker sooner when one isn't coming back
topic.metadata.refresh.interval.ms300000 (5 min)Default is usually fine; shorten it (e.g. 60000) only to pick up new brokers or partitions soonerA lost partition leader already triggers an immediate metadata refresh (topic.metadata.refresh.fast.interval.ms); the periodic refresh mainly discovers new brokers, topics and partitions
retry.backoff.ms100250–500Avoid hammering a broker mid-election
message.send.max.retries2147483647leave as-isAlready effectively unbounded β€” let the timeout be the real budget

🎯 Key Principle: Treat message.timeout.ms as the single source of truth for "how long is this operation allowed to take," and leave the retry count effectively unbounded underneath it. A transient leadership change is then absorbed automatically instead of needing every retry knob tuned in lockstep.

(librdkafka accepts delivery.timeout.ms as an alias for message.timeout.ms, which is why Java-oriented material and .NET material sometimes appear to name different settings for the same budget.)

Idempotent Producers and a Resilient AdminClient Wrapper

Producer-side resilience starts with the idempotent producer: the broker deduplicates retried writes using a producer ID and sequence number, so the client library's own retries after a timeout don't create duplicates. (It does not deduplicate a record your code sends again β€” see "Core Kafka Mental Model".) Pairing idempotence with acks=all and a capped in-flight request count keeps ordering intact even when retries fire.

using Confluent.Kafka;

var producerConfig = new ProducerConfig
{
    BootstrapServers = "broker1:9092,broker2:9092,broker3:9092",
    EnableIdempotence = true,       // dedupes retried writes at the broker
    Acks = Acks.All,                // wait for all in-sync replicas
    MaxInFlight = 5,                // capped so ordering survives retries
    MessageTimeoutMs = 30000,       // overall delivery budget per message
    RetryBackoffMs = 300,           // spacing between retry attempts
};

using var producer = new ProducerBuilder<string, string>(producerConfig).Build();

MaxInFlight caps unacknowledged requests outstanding on one connection; with idempotence on, 5 or fewer is what lets the broker keep ordering intact across retried batches (it tracks the last five batches per producer and partition). Enabling idempotence in librdkafka sets this bound and acks=all for you, and Build() throws an InvalidOperationException if you set a conflicting value such as MaxInFlight = 50 β€” so setting them explicitly documents intent rather than adding behavior.

AdminClient operations need the same posture but fail differently. Instead of a message timing out, a call like CreateTopicsAsync can throw a timeout when the controller quorum is slow to answer. A CI provisioning step that lets that bubble up is treating a condition that clears on its own as a hard error.

using System.Linq;
using Confluent.Kafka;
using Confluent.Kafka.Admin;

static async Task CreateTopicWithRetryAsync(
    IAdminClient adminClient,
    TopicSpecification spec,
    int maxAttempts = 5)
{
    // Worth retrying around a controller or leadership change; anything else is a real failure.
    static bool IsTransient(ErrorCode code) =>
        code == ErrorCode.Local_TimedOut ||    // client's own timeout: the usual case
        code == ErrorCode.RequestTimedOut ||   // broker gave up first (RequestTimeout > 60 s)
        code == ErrorCode.NotController;       // ZooKeeper-mode brokers; KRaft brokers retry it

    var delay = TimeSpan.FromMilliseconds(500);

    for (var attempt = 1; ; attempt++)
    {
        try
        {
            await adminClient.CreateTopicsAsync(new[] { spec });
            return; // success
        }
        catch (CreateTopicsException ex) when (
            ex.Results.Any(r => r.Error.Code == ErrorCode.TopicAlreadyExists))
        {
            return; // idempotent for a provisioning step: already exists is fine
        }
        // CreateTopicsException's own Error is always Local_Partial ("see Results"),
        // so the per-topic results are what say whether a retry can help.
        catch (CreateTopicsException ex) when (
            attempt < maxAttempts && ex.Results.Any(r => IsTransient(r.Error.Code)))
        {
        }
        // A failure of the request as a whole arrives as a plain KafkaException.
        catch (KafkaException ex) when (attempt < maxAttempts && IsTransient(ex.Error.Code))
        {
        }

        await Task.Delay(delay); // only reached after a caught, retryable failure
        delay *= 2;
    }
    // On the last attempt no filter matches, so the real exception reaches the caller.
}

The clauses do different jobs: "already exists" counts as success for a CI step that may run twice; the transient checks retry only the error codes that show up around a controller or leadership change. Where to look matters: a CreateTopicsException's top-level Error is always Local_Partial, and the real per-topic codes are in Results; a request that failed as a whole (for example a client-side timeout) arrives as a plain KafkaException instead. On the final attempt no filter matches, so the caller sees the original exception with its real error code rather than a generic wrapper. ⚠️ Common Mistake: catching KafkaException broadly and retrying on any code β€” that also retries genuine configuration errors (an invalid replication factor, say) that will never succeed, turning a fast clear failure into a slow confusing one.

Tuning Metadata Refresh and Connection Setup for Fast Recovery

Retry logic only helps if the client's view of the cluster is current enough for the retry to succeed. Two settings control how quickly a client notices topology changes: the metadata refresh interval (topic.metadata.refresh.interval.ms), which sets how often the client re-fetches metadata, and socket.connection.setup.timeout.ms, bounding how long the client waits on a TCP connection before trying another broker.

Left at defaults, a .NET client can be slower than the cluster itself to recover β€” the cluster has already moved leadership while the client waits 30 seconds on a socket to a broker that isn't coming back.

var adminConfig = new AdminClientConfig
{
    BootstrapServers = "broker1:9092,broker2:9092,broker3:9092",
    SocketConnectionSetupTimeoutMs = 10000,      // fail fast on an unreachable broker
    TopicMetadataRefreshIntervalMs = 60000,      // discover new brokers/partitions sooner
    MetadataMaxAgeMs = 180000                    // hard ceiling on cached metadata age
};

using var admin = new AdminClientBuilder(adminConfig).Build();

All three are strongly typed on ClientConfig, so they're available on AdminClientConfig, ProducerConfig, and ConsumerConfig alike β€” no string-keyed Set() needed. In librdkafka, topic.metadata.refresh.interval.ms is the knob that actually drives periodic refreshes; metadata.max.age.ms is the outer bound and should stay comfortably larger than it.

πŸ’‘ Pro Tip: Apply the connection-setup timeout to producer and consumer configs too, not just AdminClient. You don't need an aggressive periodic refresh to follow leadership changes: when a produce or fetch hits a "not leader" error, librdkafka requests fresh metadata immediately and then backs off from topic.metadata.refresh.fast.interval.ms (defaults to retry.backoff.ms) up to retry.backoff.max.ms.

Edge Case: Provisioning Topics During an Active Controller Election

Picture a CI step that calls CreateTopicsAsync at the moment a controller election is in progress. Under KRaft the client never talks to a controller directly β€” it sends the request to a broker, and the broker forwards it to the active controller. During a normal election (2–4 s) the broker holds the request and re-sends it once a new controller is active, so the call just takes a few seconds longer. The retry wrapper is for when the quorum stays leaderless longer, for example after losing its majority: then the call fails with a timeout. With default settings that is Local_TimedOut in a plain KafkaException: the client's own 60 s request timeout starts before the broker's 60 s limit on holding the request, so it runs out first. RequestTimedOut in CreateTopicsException.Results appears only if the call's RequestTimeout is set above 60 s. Not a malformed request, not a permissions problem, just "ask me again in a moment."

Without the wrapper, one timed-out request fails an entire deployment pipeline over a condition that clears once the quorum has a leader again. With it, the request is retried after a short backoff, up to the attempt limit. The fix isn't clever engineering β€” it's refusing to treat a known, transient, documented error code as fatal.

Harder Variant: Provisioning Topics Mid Rolling Upgrade

A tougher version appears when a provisioning step runs against a cluster mid rolling upgrade β€” not a single election, but a sustained period where different brokers sit at different points in the metadata log. In that state CreateTopicsAsync can return success from the broker that answered, while a different broker a later pipeline step talks to hasn't caught up and briefly reports the topic as not found.

The retry-on-error-code pattern doesn't cover this, because the create call itself succeeds β€” the danger is a false negative in a later verification step. The fix is to poll until the topic shows up rather than trust a single success response. GetMetadata asks whichever broker the client picks, so this absorbs most short lags without proving that every broker has caught up:

using System.Linq;

static async Task VerifyTopicVisibleAsync(
    IAdminClient adminClient,
    string topicName,
    int maxAttempts = 6)
{
    var delay = TimeSpan.FromMilliseconds(500);

    for (var attempt = 1; attempt <= maxAttempts; attempt++)
    {
        // Note: on a cluster with auto.create.topics.enable left on, a metadata
        // request for a missing topic can create it. That's benign here because
        // we only call this for a topic we just created deliberately.
        var metadata = adminClient.GetMetadata(topicName, TimeSpan.FromSeconds(5));
        var topic = metadata.Topics.FirstOrDefault(t => t.Topic == topicName);

        if (topic is not null && topic.Error.Code == ErrorCode.NoError)
        {
            return; // topic is visible with no per-topic error
        }

        if (attempt < maxAttempts) // don't sleep after the final check
        {
            await Task.Delay(delay);
            delay *= 2;
        }
    }

    throw new InvalidOperationException(
        $"Topic '{topicName}' was created but is not yet visible in metadata.");
}

Calling CreateTopicWithRetryAsync then VerifyTopicVisibleAsync gives a pipeline two defenses: the first absorbs errors thrown during creation, the second absorbs most of the lag between creation succeeding and the cluster converging on that fact. ⚠️ Skipping verification and treating a successful create as the end of the story is a common way a CI pipeline passes while a smoke test against a lagging broker fails moments later β€” two failures that look unrelated unless you know they share a metadata-log catch-up window.

Taken together these patterns share one posture: assume the cluster will occasionally answer "not yet" rather than "yes" or "no," and treat "not yet" as a reason to wait and ask again rather than to fail. That costs a handful of retry lines and turns KRaft's fast-but-nonzero failover into something a .NET service rides out without failing.

✍️ Exercise: does the timeout budget cover a failover?

A topic has replication factor 3. At t = 0 the broker leading one of its partitions is hard-killed (no clean shutdown). Defaults are in force on the cluster (broker.session.timeout.ms = 9000). A producer sends a record to that partition at t = 0.

  1. Roughly when does the partition get a new leader?
  2. Will the record be delivered with MessageTimeoutMs = 5000? With MessageTimeoutMs = 30000?
  3. Harder variation: the dead node was a combined broker+controller and was also the active controller of a 3-voter quorum. What changes?
Check your answer
  1. Around t β‰ˆ 9 s: the controller fences the broker after broker.session.timeout.ms without a heartbeat (counted from its last heartbeat, so possibly a little sooner), then elects a new leader from the ISR, and the producer's next metadata refresh finds it.
  2. With 5000 ms the record's budget runs out at t = 5 s, before a new leader exists β€” the delivery report (what ProduceAsync returns or throws) fails with a local message-timeout error. With 30000 ms the budget lasts until t = 30 s; the retries after t β‰ˆ 9 s succeed, with about 20 s to spare.
  3. First the quorum needs a new active controller: about 2 s before the standbys give up on the old leader, plus the election β€” say 2–4 s. Only an active controller can fence the dead broker and move partition leadership, so the partition is leaderless for longer β€” plan for well over 10 s. The 30 s budget still covers it; 5 s never would. That is one more reason production clusters usually keep controllers on dedicated nodes.

Tiered Storage: Hot Local Disk, Cold Remote Storage

KRaft changed where metadata lives. Tiered storage (KIP-405) changes where data lives, and the two are independent: a KRaft cluster runs perfectly well without it.

Without tiered storage, every byte a topic retains sits on the brokers' local disks β€” retention.ms of 30 days means 30 days of segments (the files a partition's log is split into; only the newest, the active segment, is still being written) on each replica's disk. With tiered storage, a topic's log has two tiers:

  • Local (hot) tier β€” the brokers' disks, as before. The active segment and recent closed segments live here, and tail reads are served from here (usually straight from the page cache, the OS's in-memory copy of recently used file data).
  • Remote (cold) tier β€” an external store such as S3 or HDFS. Once a segment is closed, the broker copies it to remote storage; when it ages past the topic's local retention it is deleted from local disk but stays readable remotely until the topic's overall retention expires.

Consumers don't change: a consumer that asks for an old offset gets the data, and the broker fetches it from the remote tier behind the scenes (with higher latency than a tail read).

offset β†’  0 ........................ 8,000,000 ......... 9,000,000 (tail)
          |-- remote tier only ------|--- local + remote --|-active-|
          ↑ retention.ms / retention.bytes: overall limit
                                     ↑ local.retention.ms / local.retention.bytes: local limit

Turning it on takes three things, and all three are easy to forget:

  1. A storage plugin. Apache Kafka ships the RemoteStorageManager interface, not an S3 implementation; the broker needs a plugin, configured with remote.log.storage.manager.class.name and .class.path.
  2. The broker switch remote.log.storage.system.enable=true β€” off by default.
  3. The topic switch remote.storage.enable=true β€” also off by default, set per topic.

Then local.retention.ms / local.retention.bytes decide how much stays local, while retention.ms / retention.bytes keep their usual meaning as the total retention. If the local.* settings are unset, they fall back to the overall values, so nothing is removed from local disk early. Tiered storage does not support compacted topics (cleanup.policy=compact, which keep only the latest record per key).

For a .NET developer none of this touches producer or consumer code. It touches capacity planning β€” which is exactly where Mistake 4 below comes from β€” and the cost of reading old data.

✍️ Exercise: how much local disk?

A topic receives 20 MB/s on average and has retention.ms = 7 days. With tiered storage fully enabled it has local.retention.ms = 1 day. Estimate the local disk each replica of this topic needs (ignore compression and segment rounding):

  1. with tiered storage fully enabled (plugin, broker switch and topic switch);
  2. if the topic switch remote.storage.enable was never set.
Check your answer
  1. One day locally: 20 MB/s Γ— 86,400 s β‰ˆ 1.73 TB per replica. The other six days live only in remote storage.
  2. Seven days locally: 20 MB/s Γ— 604,800 s β‰ˆ 12.1 TB per replica β€” seven times as much. Without the topic switch the topic is not tiered at all, so retention.ms applies to local disk and local.retention.ms has no effect. Multiply by the replication factor for the cluster-wide total.

Two Other 4.x Changes Consumers Will Meet

"Modern Kafka" also brought two changes to how consumers are coordinated. Both are covered from the client side in "Consumer Groups & Offsets":

  • The next-generation consumer group protocol (KIP-848) is generally available since Kafka 4.0: the broker-side group coordinator computes partition assignments and rebalances incrementally, instead of a client acting as group leader. librdkafka supports it for production from 2.12 (in preview since 2.10) but still defaults to the classic protocol, so in Confluent.Kafka you opt in with GroupProtocol = GroupProtocol.Consumer on ConsumerConfig.
  • Share groups (KIP-932, "Queues for Kafka") became production-ready in Kafka 4.2: several consumers can share one partition with per-record acknowledgement, which a consumer group never allows. Confluent.Kafka (2.15 at the time of writing) does not expose a share consumer yet, although librdkafka 2.15 has one in preview, so check your client's release notes before designing around it.

Common Mistakes .NET Teams Make with Modern Kafka Clusters

Migrating an existing ZooKeeper cluster to KRaft is an operator project with its own procedure β€” it requires a 3.x "bridge" release (3.9 is the last), because Kafka 4.0 cannot migrate from ZooKeeper directly. The mistakes that bite .NET teams show up afterward, buried in manifests, package references, and CI pipelines nobody revisits once the cluster "just works."

Mistake 1: Zombie ZooKeeper Configuration ⚠️

The most common mistake isn't a KRaft misconfiguration β€” it's leftover ZooKeeper configuration nobody deleted. Teams point bootstrap.servers at the new cluster, confirm producers and consumers work, and call the migration done. But the compose file, Helm chart, or Kubernetes manifest that stood up the old cluster often still contains a zookeeper.connect variable, a separate ZooKeeper container, and β€” critically β€” a readiness or liveness probe pinging ZooKeeper's port before marking the broker pod healthy.

❌ Wrong thinking: "The broker starts fine, so the leftover ZK block is harmless dead weight." βœ… Correct thinking: A stale ZK health check is an active failure mode, not inert clutter β€” it can block pod readiness, fail CI health checks, or cause an orchestrator to restart a perfectly healthy broker because a service it no longer depends on isn't responding.

# ⚠️ Leftover from a pre-KRaft manifest β€” this probe targets a ZooKeeper
# port that no longer exists once process.roles is broker,controller
readinessProbe:
  tcpSocket:
    port: 2181   # ZooKeeper's client port
  initialDelaySeconds: 10
  periodSeconds: 5

If the ZooKeeper sidecar was removed from the broker pod but this probe wasn't (a tcpSocket probe checks the pod's own IP), the pod never reports ready, and any integration test or CI step waiting on readiness times out with an error that looks like a networking problem. The fix is a deliberate audit: grep every compose file, Helm values file, and manifest for zookeeper, zkClient, and port 2181, and remove them alongside --zookeeper flags in CI shell scripts.

Mistake 2: Treating the Controller Quorum Size as a Formality ⚠️

It's tempting to fill in the controller quorum with whatever node count feels convenient β€” one for a quick dev setup, or an even number because it "seemed reasonable" β€” without registering that this setting is the fault-tolerance boundary of the cluster's metadata plane.

A single controller voter gives zero redundancy: if that process dies, no new controller can be elected and metadata operations β€” topic creation, partition reassignment, ISR updates β€” stall until it's replaced. An even number doesn't buy anything over the next odd number down, because Raft-style quorums need a strict majority to make progress. Three voters tolerate one failure; four voters also tolerate only one, while requiring an extra machine and an extra vote in every round. Five tolerate two.

# βœ… Three-voter quorum: tolerates one controller failure
controller.quorum.voters=1@ctrl-1:9093,2@ctrl-2:9093,3@ctrl-3:9093

# ❌ Single voter: any controller loss halts metadata operations
controller.quorum.voters=1@ctrl-1:9093

πŸ’‘ Pro Tip: Treat quorum size as a majority vote, not a replica count β€” pick an odd number (three is the common baseline, five when you need to survive two simultaneous failures) and revisit it whenever you change the number of controller-eligible nodes.

✍️ Exercise: how many voters can you lose?
  1. Fill in the table: for 1 to 5 voters, the majority needed and the number of voters that can fail while the quorum still makes progress.
  2. Your quorum has 3 voters and one is down for OS patching. Is it safe to restart a second voter now?
Check your answer
  1. A majority of n is ⌊n/2βŒ‹ + 1, so the tolerated failures are n βˆ’ majority:
Voters Majority Can fail
1 1 0
2 2 0
3 2 1
4 3 1
5 3 2
  1. No. With one voter already down, restarting a second leaves one of three β€” no majority. The quorum can't elect a leader or commit metadata changes until a voter returns, so topic creation, leadership moves and ISR updates stall. Wait until the first voter has rejoined and caught up before touching the next one.

Mistake 3: Pinning a Years-Old Client Library Version ⚠️

Teams pin Confluent.Kafka (and transitively librdkafka) in a .csproj and leave it there, especially in services that "work fine." Then the brokers are upgraded to Kafka 4.x β€” often in the same project as the move to KRaft, which is why KRaft gets the blame β€” and a pin that is old enough can start failing in confusing ways: features reported as "not supported by broker", or errors that look like network issues.

The underlying issue is API version negotiation. Kafka's protocol evolves per API: the broker advertises the version range it supports for each API, and librdkafka picks the highest version both sides know. Kafka 4.0 removed a set of old API versions (KIP-896). A client old enough to know only removed versions of an API fails every call to it (for admin calls, with the local Local_UnsupportedFeature error described earlier), while APIs where it knows a newer version keep working. KIP-896 lists which client versions are affected; check it rather than guessing.

<!-- ❌ Old, unmaintained pin left in a .csproj for months -->
<PackageReference Include="Confluent.Kafka" Version="1.4.2" />

<!-- βœ… A current version (2.15.1 at the time of writing β€” check NuGet), pinned
     explicitly and verified against your broker. Prefer an exact version over a
     floating range so builds stay reproducible, and bump it deliberately. -->
<PackageReference Include="Confluent.Kafka" Version="2.15.1" />

⚠️ Common Mistake: Assuming that because produce and consume still work, the library is fully compatible. Version floors apply per API, so one kind of request can fail while the rest work β€” and the two librdkafka quirks KIP-896 called out were in Produce and JoinGroup (the request a consumer sends to join its consumer group), not in admin calls. The API-version check from "Checking Client/Broker API Version Compatibility" is the concrete way to confirm rather than guess.

πŸ’‘ Real-World Example: Kafka 4.0.0 removed JoinGroup v0–v1 under KIP-896. librdkafka before 2.11 (so Confluent.Kafka before 2.11) decides whether a broker supports Kerberos (SASL GSSAPI) by checking that it still offers JoinGroup v0, so services on those versions that log in with Kerberos could no longer authenticate against 4.0.0 brokers, with no code change on their side. librdkafka 2.11 fixed the check, and Kafka 4.0.1 and 4.1.0 put the old versions back (KAFKA-19444).

Mistake 4: Sizing Retention as If Tiered Storage Is Already On ⚠️

A subtler mistake: teams read about tiered storage β€” the split between fast local disk for recent segments and cheaper remote object storage for older ones β€” and assume it changes their retention math automatically. The .NET-facing error is specific: configuring retention as if remote offload is happening when the feature was never enabled.

Three settings are involved, and conflating them is the trap:

βš™οΈ SettingπŸ“ Scope🎯 Meaning
remote.log.storage.system.enableBrokerTurns the tiered storage subsystem on at all
remote.storage.enableTopicOpts this topic into remote offload
local.retention.ms / local.retention.bytesTopicHow much stays on local disk once offload is active

Set a generous 30-day retention.ms on a high-throughput topic reasoning that "old segments offload to cheap storage anyway," without tiered storage actually on for that topic (plugin, broker switch and topic switch), and every byte accumulates on local disk. The volume fills, and the broker takes the log directory offline β€” which surfaces to clients not as a tidy "disk full" but as produce failures and leadership errors on the affected partitions, since a broker with an offline log dir stops serving them.

using Confluent.Kafka;
using Confluent.Kafka.Admin;

// Verification, not configuration: check what the broker actually reports
// before trusting retention math that depends on offload.
using var admin = new AdminClientBuilder(new AdminClientConfig
{
    BootstrapServers = "localhost:9092"
}).Build();

var results = await admin.DescribeConfigsAsync(new[]
{
    // The subsystem switch is a BROKER config: describe a broker by its node ID.
    new ConfigResource { Type = ResourceType.Broker, Name = "1" },
    // The opt-in and local-retention settings are TOPIC configs.
    new ConfigResource { Type = ResourceType.Topic, Name = "orders" }
});

foreach (var result in results)
{
    foreach (var entry in result.Entries.Values)
    {
        if (entry.Name.Contains("remote", StringComparison.OrdinalIgnoreCase) ||
            entry.Name.StartsWith("local.retention", StringComparison.OrdinalIgnoreCase))
        {
            Console.WriteLine($"{entry.Name} = {entry.Value}");
        }
    }
}

Note DescribeConfigsAsync β€” like DescribeClusterAsync, there is no synchronous overload. The fix for the underlying mistake is procedural: size retention on confirmed broker-side state, not on the theoretical existence of the feature, and treat dev and CI environments β€” which almost always skip object storage setup β€” as local-disk-only unless proven otherwise.

Mistake 5: Treating KRaft as Invisible Plumbing and Skipping Integration Tests ⚠️

The last mistake is the most conceptual, and its failures tend to surface only in production. Because the client-facing API surface barely changes, it's easy to conclude KRaft is purely internal and never needs testing against directly.

❌ Wrong thinking: "We already have integration tests against a Kafka container from before the migration; the client code didn't change, so those tests still validate production behavior." βœ… Correct thinking: Tests last validated against a ZooKeeper-mode cluster can pass while missing real behavioral differences that only appear under KRaft, particularly around AdminClient timing and error codes.

The concrete gap is admin API behaviour. A team that never updated its Testcontainers image from a ZooKeeper-based Kafka image to a KRaft-mode one never exercises KRaft's admin path β€” topic and config changes forwarded by a broker to the active controller, a Controller field that names a random broker β€” until production does.

🎯 Key Principle: If your service calls the AdminClient for anything beyond trivial reads β€” topic creation, partition increases, config updates β€” your integration suite should run against a KRaft-mode cluster like production's (a single combined node covers most admin behaviour; testing controller failover needs three voters), not a leftover ZK-era image kept around out of inertia.

Most of these share a shape: something that was true or harmless in the old setup β€” a health check target, a client library pin, a belief that internals don't matter β€” quietly stops being true after the move to a modern cluster, and nothing forces a team to notice until it fails. (Quorum sizing is the exception: odd-sized quorums mattered for ZooKeeper ensembles too β€” it's just that the quorum now lives inside Kafka. Retention is the reverse case: the old local-disk math still holds until tiered storage is actually switched on.) The fix in every case is the same discipline: audit infrastructure files explicitly rather than assuming they were updated, verify broker-side state rather than trusting configuration intent, and test against the topology you actually run.

Modern Kafka Mental Model: Recap and Where to Go Next

You've walked through the config-level differences between ZooKeeper-era and KRaft-era clusters, spun up a KRaft node from a .NET solution, and hardened a producer and an AdminClient against the new failure modes. The mental model underneath is simple to state and easy to forget under deadline pressure: your .NET code has exactly one door into the cluster, and ZooKeeper was never that door for application traffic β€” it's just gone as an option entirely now.

The One-Door Model

In the ZooKeeper era, a surprising amount of tooling and mental overhead existed because two systems needed reasoning about. Even though .NET producers and consumers never spoke the ZooKeeper protocol, plenty of adjacent work did β€” health checks, migration scripts, diagnostics that shelled out to zkCli.sh or passed --zookeeper. That second system is what's been removed, not a second connection string your client code maintained.

Your .NET service
    ↓
bootstrap.servers  (initial connection, protocol-level only)
    ↓
Kafka protocol requests: produce, fetch, metadata, admin RPCs
    ↓
Brokers (some of which may also serve as controllers)

Every operation β€” producing, consuming, or calling DescribeClusterAsync() β€” travels this same path. No second protocol, no second port, no second client library. If you find yourself reaching for a ZooKeeper client package or a zookeeper.connect setting in a modern Kafka project, that's a signal you're solving a problem that no longer exists.

Quick-Reference Checklist for a .NET Project

βœ… CheckπŸ” What to verifyπŸ› οΈ How to verify it from .NET
πŸ†” cluster.id is the one you provisionedNodes were formatted with your ID β€” not the apache/kafka image's fixed default(await admin.DescribeClusterAsync()).ClusterId
🎭 process.roles correctBroker/controller roles match topology intentDocker/K8s manifest review, not a client call
πŸ“‘ Advertised listeners matchThe address brokers hand back is reachable from the clientConnect and produce one message; if bootstrap succeeds but the produce times out, suspect this first
πŸ“¦ Client version verifiedThe broker supports the API versions your client usesDebug = "broker,feature" trace
⏱️ Resilience timeouts reviewedMessage timeout and metadata refresh set deliberatelyConfig review against librdkafka defaults

If DescribeClusterAsync() ever returns an unexpected cluster ID in staging, that usually means someone pointed a client at a different cluster, or a node was re-formatted with a new ID, rather than a genuine client bug β€” worth ruling out before debugging retry logic.

using Confluent.Kafka;
using Confluent.Kafka.Admin;

var adminConfig = new AdminClientConfig
{
    BootstrapServers = "localhost:9092" // the one door: no zookeeper.connect exists
};

using var admin = new AdminClientBuilder(adminConfig).Build();

// Cluster identity and broker count in one call: a cheap connectivity and
// identity check. It does not prove every admin API is compatible β€” an API-version
// mismatch only shows up on the call that uses that API.
var description = await admin.DescribeClusterAsync(new DescribeClusterOptions
{
    RequestTimeout = TimeSpan.FromSeconds(5)
});

Console.WriteLine($"Cluster ID: {description.ClusterId}");
Console.WriteLine($"Live brokers: {description.Nodes.Count}");

Drop this into a startup routine or CI job so a misconfigured cluster fails fast with a clear message instead of surfacing as a mysterious timeout three services downstream.

Resilience Configuration Is the Default Posture, Not a Special Case

It's tempting to treat retry/backoff and idempotent-producer patterns as extra hardening you add once a service is important enough to justify the effort. That gets the risk backwards. KRaft's failover is fast, but "fast" still means a real window: a crashed broker's partitions wait seconds for a new leader, and admin calls are held β€” or time out if the controller quorum stays leaderless. A service that adds retry logic only after its first outage has, by definition, already paid the cost the default posture was designed to avoid.

// The same defensive shape from "Writing Resilient .NET Clients" β€” repeated to
// underline that it's a baseline, not an advanced option for high-traffic services.
var producerConfig = new ProducerConfig
{
    BootstrapServers = "localhost:9092",
    EnableIdempotence = true,
    Acks = Acks.All,
    MaxInFlight = 5,
    MessageTimeoutMs = 30000,                 // tightened from the 5-minute default
    RetryBackoffMs = 300
};

⚠️ Common Mistake: Copying producer and consumer configuration from an older reference project without revisiting the timeouts. Those values were chosen for a different cluster, client version and failure profile; leaving them unexamined means either a budget too short to ride out a broker failover, or one so long that real failures take minutes to surface.

Where the Deeper Internals Live

This lesson stayed at the boundary between your .NET code and the cluster.

The controller quorum's internals β€” how voters reach consensus, how voters elect a leader while observers only follow the log, how dynamic membership (controller.quorum.bootstrap.servers) is managed with kafka-metadata-quorum.sh β€” are operator territory, documented in the KRaft section of the Kafka operations docs. You don't need them to configure a .NET client correctly, but they become valuable when debugging quorum sizing or failure tolerance during a rolling upgrade.

Tiered storage was covered here only as far as it changes capacity planning; plugin choice, remote-read latency and cost belong with whoever runs the cluster. The connection back to your code is narrow but important: a dev environment with retention sized for a tiered-storage production cluster can fill its disk, because the feature is almost never enabled locally.

🎯 Key Principle: The quick-reference checklist above is a starting heuristic for catching the most common KRaft-related misconfigurations quickly β€” not an exhaustive audit. Cluster-specific topologies (multi-datacenter controller placement, custom authentication, managed-service abstractions) introduce failure modes these items won't catch.

Practical Next Steps

Three actions turn this recap into something you apply. First, audit one existing service's Kafka configuration against the checklist β€” most teams find at least one leftover assumption, whether a stale ZooKeeper health check, an unreviewed timeout, or a client version pinned years ago. Second, if your integration tests still spin up a ZooKeeper container alongside a broker, replace that with a single combined process.roles=broker,controller node β€” a smaller compose file and a faster test startup. Third, before your next deployment touching Kafka, apply the idempotent-producer-plus-retry-wrapper pattern even if the service has never had a Kafka incident; treat it as the default, not the fix.

πŸ’‘ Remember: ZooKeeper's removal simplified the cluster's internals without changing the shape of your application code β€” the door your .NET service walks through was always bootstrap.servers and the Kafka protocol, and KRaft just closed off every other door that used to exist alongside it.