Apache Kafka Explanation¶
What this page covers
How Apache Kafka 4.x works and why it is built this way: the KRaft metadata quorum, the partitioned replicated log, ISR replication and Eligible Leader Replicas, the consumer, share, and streams group protocols, exactly-once transactions, tiered storage, and the security model. Exact config defaults, versions, and CVEs are in Reference. Step-by-step tasks are in How-to Guides.
Related Notes
Overview¶
Apache Kafka is a distributed, partitioned, replicated commit-log service. Since Kafka 4.0 (KRaft only), a cluster has three logical roles:
- Brokers store partition logs and serve produce and fetch requests.
- Controllers form a Raft quorum that owns the cluster metadata.
- Clients (producer, consumer, share consumer, Admin client, Connect workers, Streams applications) talk to brokers over a versioned binary TCP protocol.
Topics are split into partitions that are spread across brokers. Each partition is an append-only log with one leader and zero or more followers. Followers copy the leader by sending Fetch requests. A record is committed once every replica in the in-sync replica set (ISR) has it, and only then does the partition's high watermark move past it.
Component Architecture¶
The diagram shows the main classes inside a 4.x broker and controller and how clients reach them.
graph TB
subgraph Clients["Client Layer"]
KP["KafkaProducer<br/>(idempotent / transactional)"]
KC["KafkaConsumer<br/>(classic or consumer protocol)"]
KSC["KafkaShareConsumer<br/>(share groups)"]
ADM["Admin client"]
STR["KafkaStreams runtime"]
CON["Kafka Connect worker"]
end
subgraph Broker["BrokerServer (process.roles=broker)"]
SR["SocketServer<br/>(acceptor + processor threads)"]
RH["KafkaApis<br/>(request router)"]
RM["ReplicaManager"]
LM["LogManager"]
GC["GroupCoordinator<br/>(consumer, classic, streams groups)"]
SC["ShareCoordinator +<br/>SharePartitionManager"]
TC["TransactionCoordinator"]
RLM["RemoteLogManager<br/>(KIP-405)"]
MDC["MetadataLoader +<br/>BrokerMetadataPublisher"]
end
subgraph Controller["ControllerServer (process.roles=controller)"]
QC["QuorumController"]
RAFT["KafkaRaftClient<br/>(__cluster_metadata)"]
SNAP["Metadata snapshots<br/>(.checkpoint files)"]
end
subgraph Storage["On-Disk Storage"]
SEG["LogSegment files<br/>.log / .index / .timeindex / .txnindex"]
OffStore["__consumer_offsets<br/>(50 partitions, compacted)"]
TxnStore["__transaction_state<br/>(50 partitions, compacted)"]
ShareStore["__share_group_state"]
end
subgraph Remote["Remote Tier"]
RSM["RemoteStorageManager plugin<br/>(S3 / GCS / Azure / HDFS)"]
RLMM["RemoteLogMetadataManager<br/>(__remote_log_metadata)"]
end
KP -->|Produce| SR
KC -->|"Fetch, ConsumerGroupHeartbeat"| SR
KSC -->|"ShareFetch, ShareAcknowledge"| SR
ADM -->|"CreateTopics, IncrementalAlterConfigs"| SR
STR --> KP
STR --> KC
CON --> KP
CON --> KC
SR --> RH
RH --> RM
RH --> GC
RH --> SC
RH --> TC
RM --> LM
LM --> SEG
GC --> OffStore
TC --> TxnStore
SC --> ShareStore
LM --> RLM
RLM --> RSM
RLM --> RLMM
RAFT -->|"TopicRecord, PartitionRecord,<br/>RegisterBrokerRecord"| MDC
MDC --> RM
QC --> RAFT
RAFT --> SNAP
style Controller fill:#1f3a5f,color:#fff
style Remote fill:#3a5a3a,color:#fff
style Broker fill:#4a3a3a,color:#fff
Core Components¶
| Component | Role |
|---|---|
| BrokerServer | The broker process in KRaft mode. It starts SocketServer, KafkaApis, ReplicaManager, LogManager, the coordinators, and RemoteLogManager. (KafkaServer was the ZooKeeper-mode broker and is gone in 4.0.) |
| ControllerServer / QuorumController | Hosts the controller. The active QuorumController processes metadata changes and writes them to the Raft log. |
| KafkaRaftClient | Kafka's Raft implementation. It replicates __cluster_metadata between controllers (voters) and brokers (observers). |
| SocketServer | NIO acceptor and processor threads for TCP connections. Parsed requests go onto a request channel. |
| KafkaApis | Request router. It dispatches each API key (Produce, Fetch, Metadata, OffsetCommit, ShareFetch, and others) to the right subsystem. |
| ReplicaManager | Owns local replica state. It appends produce data and serves Fetch requests from consumers and followers. |
| LogManager | Manages partition logs on disk: segment rolling, retention, compaction scheduling, and recovery on startup. |
| GroupCoordinator | The group coordinator rewritten for 4.0. It runs classic, consumer (KIP-848), and streams (KIP-1071) groups and stores offsets in __consumer_offsets. |
| ShareCoordinator / SharePartitionManager | Tracks per-record delivery state for share groups (KIP-932) and persists it to __share_group_state. |
| TransactionCoordinator | Manages transactional IDs, writes commit and abort markers, and fences zombie producers. |
| RemoteLogManager (RLM) | Copies closed local segments to the remote tier through the configured RemoteStorageManager plugin and serves remote reads. |
| KafkaProducer | Serializes records, batches them per partition, compresses them, attaches idempotence sequence numbers, and optionally joins transactions. |
| KafkaConsumer | Implements group membership (classic or KIP-848 protocol), commits offsets, and supports read_committed isolation. |
| KafkaShareConsumer | Queue-style consumer for share groups with per-record acknowledgement. |
| Admin client | Topic and ACL CRUD, dynamic configs, group describe and reset, partition reassignment, feature upgrades. |
| Kafka Streams | Embedded JVM library implementing the KStream/KTable DSL on top of the producer, the consumer, and RocksDB state stores. |
| Kafka Connect | Distributed worker framework for source and sink connectors (Debezium, JDBC, S3, Iceberg, and others). |
Why Kafka Is Fast¶
Kafka's throughput comes from a few deliberate design choices described in the official design documentation:
- Sequential I/O on an append-only log. Producers only append, and consumers read forward. Disks and SSDs do sequential work far faster than random access.
- The OS page cache instead of a JVM cache. Kafka writes to the filesystem and lets Linux cache recent data. Tail reads by up-to-date consumers are usually page-cache hits, and the JVM heap stays small, which keeps GC pauses short.
- Batching end to end. Producers group records into record batches (see the batch format). The broker stores and forwards the batch as-is, so work is spread over many records.
- Compression per batch. The whole batch is compressed once (gzip, snappy, lz4, or zstd). That gives better ratios than compressing each record, and the broker does not have to recompress.
- Zero-copy transfer. For plaintext listeners the broker uses
FileChannel.transferTo(Linuxsendfile(2)) to move bytes from the page cache to the socket without copying them through user space.
Published throughput numbers and their test conditions are in Reference: Benchmark Results.
KRaft Consensus¶
Since Kafka 4.0 there is no ZooKeeper. A small set of nodes (usually 3 or 5) are controllers (process.roles=controller, or broker,controller in combined mode for small or dev clusters). One controller is the active controller, elected through Raft. The others are hot standbys. Brokers are Raft observers: they follow the metadata log but do not vote.
KRaft is pull-based like the rest of Kafka. Followers and observers send Fetch requests for __cluster_metadata to the leader, instead of the leader pushing AppendEntries as in textbook Raft. The sequence below shows a topic creation.
sequenceDiagram
participant Admin as Admin client
participant B as Broker
participant Active as Active controller
participant Voters as Follower controllers
participant Obs as Other brokers (observers)
Admin->>B: CreateTopics orders, 12 partitions, RF 3
B->>Active: Forward CreateTopics
Active->>Active: Validate and choose replica placement
Active->>Active: Append TopicRecord and PartitionRecords to __cluster_metadata
Voters->>Active: Fetch __cluster_metadata
Active-->>Voters: New records
Note over Active,Voters: Committed once a majority of voters has the records
Active-->>B: CreateTopics result
B-->>Admin: CreateTopicsResponse
Obs->>Active: Fetch __cluster_metadata
Active-->>Obs: New records
Obs->>Obs: MetadataLoader applies records and creates local replicas
Key properties of KRaft:
- Single source of truth. Every metadata change is a record in
__cluster_metadata, so brokers catch up by replaying a log instead of receiving pushed RPCs. - Snapshots. Controllers and brokers periodically write metadata snapshots, so the log does not grow forever and new nodes bootstrap quickly.
- Faster failover. There is no ZooKeeper session timeout. A new active controller already has the full metadata in memory.
- Fewer moving parts. A dev cluster can be one combined process. A production cluster needs 3 controllers plus brokers, and no separate ZooKeeper ensemble.
- Dynamic quorums (KIP-853).
controller.quorum.bootstrap.serversreplaces the staticcontroller.quorum.voterslist. New clusters can use dynamic quorums since 3.9, and static quorums can be upgraded to dynamic ones since 4.1 (kraft.version=1). Controllers can then be added and removed withkafka-metadata-quorum.sh.
A 3-controller quorum tolerates one controller failure. Five controllers tolerate two. Controllers fsync the metadata log on commit, so they need stable, low-latency disks and network.
From ZooKeeper to KRaft¶
| Release | Milestone |
|---|---|
| 2.8 (2021) | KRaft early access |
| 3.3 (2022-10) | KRaft production-ready for new clusters (KIP-833) |
| 3.5 (2023-06) | ZooKeeper mode deprecated |
| 3.6 (2023-10) | ZooKeeper-to-KRaft migration production-ready |
| 3.9 (2024-11) | Last line that supports ZooKeeper, which makes it the bridge release for migrations |
| 4.0 (2025-03) | ZooKeeper mode and migration code removed |
Topic, Partition, and Log¶
A topic is a named stream split into N partitions. Each partition is an append-only sequence of immutable records, identified by 64-bit offsets. Order is strict within a partition. Across partitions, there is no order guarantee. The producer's partitioner picks the partition, by hashing the key when there is one.
The diagram shows how a topic's partitions, segments, and replicas map onto brokers.
graph LR
subgraph Topic["Topic: orders (RF=3, 4 partitions)"]
direction TB
subgraph P0["Partition 0"]
S0a["Segment 00000000000000000000.log"]
S0b["Segment 00000000000000012450.log (active)"]
end
subgraph P1["Partition 1"]
S1a["Segment 00000000000000000000.log"]
S1b["Segment 00000000000000018932.log (active)"]
end
P2["Partition 2"]
P3["Partition 3"]
end
P0 -->|Leader| KSa["Broker 1"]
P0 -->|Follower| KSb["Broker 2"]
P0 -->|Follower| KSc["Broker 3"]
P1 -->|Leader| KSb
P1 -->|Follower| KSa
P1 -->|Follower| KSc
Log-Structured Storage Internals¶
Each partition's log directory holds a series of segments. Files are named by the segment's base offset, zero-padded to 20 digits.
| File | Purpose |
|---|---|
<base-offset>.log |
Record batches in the magic v2 format (layout). |
<base-offset>.index |
Sparse offset-to-file-position index (memory-mapped). |
<base-offset>.timeindex |
Sparse timestamp-to-offset index, used for time-based seeks. |
<base-offset>.txnindex |
Aborted-transaction index, used by read_committed consumers. |
<base-offset>.snapshot |
Producer state snapshot (idempotence sequence numbers). |
leader-epoch-checkpoint |
End offset of each leader epoch, used to truncate safely after a leader change. |
The active segment receives writes. When it reaches segment.bytes (default 1 GiB) or segment.ms, Kafka rolls it and opens a new active segment. Closed segments become eligible for retention deletion or, with tiered storage, for upload.
Compression covers the whole records block of a batch, which gives better ratios than per-record compression. It is also why zero-copy works so well: the broker streams the compressed batch from disk to the consumer socket without decompressing it.
Log Compaction¶
Besides time and size retention (cleanup.policy=delete), Kafka supports cleanup.policy=compact. The log cleaner periodically rewrites segments and keeps only the latest record for each key. That makes compacted topics a good fit for state snapshots: Streams changelog topics, __consumer_offsets, __transaction_state, and configuration-style topics.
Tunables:
min.compaction.lag.ms: minimum age before a record can be compacted.max.compaction.lag.ms: maximum time an uncompacted record can wait.min.cleanable.dirty.ratio: fraction of uncleaned bytes that triggers a cleaning pass.
Tombstones (records with a null value) stay for delete.retention.ms, so that replicas and consumers see the delete before it is purged. Since 4.1, log.cleaner.enable is deprecated (KIP-1148): the cleaner is expected to stay on. Tiered storage does not support compacted topics.
Replication and ISR¶
Each partition has a leader and N-1 followers, where N is replication.factor. Followers send Fetch requests to the leader just as consumers do. A follower stays in sync (in the ISR) while it keeps up with the leader's log end within replica.lag.time.max.ms. The high watermark (HW) is the highest offset that every ISR member has. Consumers can only read up to the HW, so they never see a record that a leader failover could lose.
With acks=all, the leader answers the producer only after all current ISR members have the batch. The write fails with NotEnoughReplicas if the ISR is smaller than min.insync.replicas. Setting RF=3 with min.insync.replicas=2 tolerates one broker loss without losing acknowledged writes.
The sequence shows an acks=all produce with two followers.
sequenceDiagram
participant Producer as KafkaProducer
participant Leader as Broker (leader)
participant F1 as Broker (follower 1)
participant F2 as Broker (follower 2)
Producer->>Leader: ProduceRequest acks=all, batch B at offset N
Leader->>Leader: ReplicaManager appends B to the local log
par Replication
F1->>Leader: FetchRequest from offset N
Leader-->>F1: FetchResponse with B
F1->>F1: Append B, advance log end offset
and
F2->>Leader: FetchRequest from offset N
Leader-->>F2: FetchResponse with B
F2->>F2: Append B, advance log end offset
end
F1->>Leader: FetchRequest from offset N+1
F2->>Leader: FetchRequest from offset N+1
Leader->>Leader: All ISR members past N, high watermark moves to N+1
Leader-->>Producer: ProduceResponse, base offset N
The next Fetch offset doubles as the follower's acknowledgement, so replication needs no extra ack RPC. unclean.leader.election.enable=false (the default) keeps out-of-sync replicas from becoming leader, trading availability for no data loss. The full list of replication settings and defaults is in Reference.
Eligible Leader Replicas (KIP-966)¶
Kafka enforces a "strict min ISR" rule: the high watermark cannot advance while the ISR is smaller than min.insync.replicas. Before ELR, a replica that dropped out of the ISR could not be elected leader without an unclean election, even when it held every committed record.
KIP-966 part 1 adds Eligible Leader Replicas: replicas that left the ISR but are guaranteed to have all data up to the high watermark. The controller stores them in the partition metadata and, when a leader is needed, picks in this order:
- A member of the ISR, if the ISR is not empty.
- An unfenced member of the ELR, if the ELR is not empty.
- The last known leader, if it is unfenced (the pre-4.0 behavior when all replicas were offline).
ELR shipped as a preview in 4.0 (eligible.leader.replicas.version=1) and is on by default for new clusters from 4.1. It does not let acks=all writes succeed below min.insync.replicas. What it does is make leader election after multiple failures safer without an unclean election. With ELR on, min.insync.replicas must be set at the cluster or topic level: broker-level values are removed and cannot be altered.
Consumer Group Protocol¶
Consumers join a group identified by group.id, and the group coordinator spreads the subscribed partitions across members. Each partition goes to exactly one member at a time.
- Classic protocol. Members run a JoinGroup and SyncGroup barrier, and one client-side leader computes the assignment. With the eager assignors, every member gives up all partitions during a rebalance. The cooperative-sticky assignor (KIP-429, 2.4) lessened this but kept the group-wide barrier.
- Consumer protocol (KIP-848, GA in 4.0). The broker computes assignments and sends them to each member through
ConsumerGroupHeartbeat. Changes apply incrementally per member, with no group-wide synchronization barrier. Consumers opt in withgroup.protocol=consumer. The client default is stillclassicin 4.3, but 4.3 logs a recommendation to switch (KIP-1274).
The state diagram shows a member's lifecycle under the consumer protocol.
stateDiagram-v2
[*] --> Joining: subscribe(topics)
Joining --> Stable: Coordinator sends target assignment
Stable --> Reconciling: Membership or subscription changes
Reconciling --> Stable: Revoke then assign only changed partitions
Stable --> Fenced: Heartbeat or poll interval missed
Fenced --> Joining: Rejoin with epoch 0
Stable --> Leaving: close()
Leaving --> [*]
Other group types built on the same coordinator:
- Streams groups (KIP-1071,
group.protocol=streams): broker-side task assignment for Kafka Streams. Early access in 4.1, GA for the core feature set in 4.2. 4.4 adds static membership and custom broker-side assignors. - Share groups (KIP-932): queue semantics. See the next section.
Share Groups (Queues for Kafka)¶
Share groups (KIP-932) add cooperative, queue-style consumption on ordinary topics. They were early access in 4.0, a preview in 4.1, and production-ready in 4.2. How they differ from consumer groups:
- Several consumers can read the same partition at once, so the number of consumers can exceed the number of partitions.
- Records are acknowledged one by one (the client still fetches in batches).
- The broker counts delivery attempts, so it can set aside poison messages automatically.
- There is no ordering guarantee between records handed to different consumers.
When a share consumer fetches, the broker acquires records for it with a time-limited lock (share.record.lock.duration.ms, default 30 s). The consumer then accepts, releases, rejects, or renews each record. The record states are shown below.
stateDiagram-v2
[*] --> Available: Record appended
Available --> Acquired: ShareFetch by a consumer
Acquired --> Acknowledged: ACCEPT
Acquired --> Available: RELEASE or lock timeout
Acquired --> Acquired: RENEW (extends the lock)
Acquired --> Archived: REJECT
Available --> Archived: Delivery count limit reached
Acknowledged --> [*]
Archived --> [*]
The delivery count limit defaults to 5 (share.delivery.count.limit). Share-group state lives in the __share_group_state internal topic. Kafka 4.4 adds a dead-letter queue for rejected or exhausted records (KIP-1191, share.version=2). In 4.2 and 4.3, rejected records are simply archived.
When to use share groups
Use share groups when each record is an independent unit of work, such as jobs or notifications, and you need more consumers than partitions or per-message retry. Keep consumer groups for ordered stream processing. RabbitMQ still offers richer routing, priorities, and per-message TTL.
Exactly-Once Semantics¶
EOS in Kafka rests on two primitives:
- Idempotent producer (on by default since 3.0). The broker gives each producer a producer ID and tracks sequence numbers per partition, so retried batches are dropped as duplicates. Batches that arrive out of order are rejected with
OutOfOrderSequenceException. - Transactions. A producer with a
transactional.idcallsinitTransactions(), then groups sends and offset commits betweenbeginTransaction()andcommitTransaction(). The TransactionCoordinator commits or aborts atomically by writing transaction markers to every partition involved. Consumers withisolation.level=read_committedread only up to the last stable offset and skip aborted records using the.txnindexfiles.
Read-process-write loops, the standard Kafka Streams pattern, get end-to-end exactly-once by adding the consumer's offsets to the producer's transaction with producer.sendOffsetsToTransaction(...). In 4.0, KIP-890 ("transactions server-side defense") bumps the producer epoch on every transaction, which closes a class of hanging and duplicate-transaction bugs.
The sequence shows one transactional read-process-write cycle.
sequenceDiagram
participant Cons as Consumer
participant App as Kafka Streams task
participant Prod as Transactional producer
participant TC as TransactionCoordinator
participant T1 as Topic A (input)
participant T2 as Topic B (output)
participant OS as __consumer_offsets
Cons->>T1: poll() returns records
App->>Prod: beginTransaction()
Prod->>TC: AddPartitionsToTxn for topic B partitions
Prod->>T2: send(transformed records)
Prod->>TC: AddOffsetsToTxn for the group
Prod->>OS: TxnOffsetCommit with consumer offsets
Prod->>TC: EndTxn commit
TC->>T2: Write COMMIT marker
TC->>OS: Write COMMIT marker
Note over Cons,OS: read_committed consumers now see the output and the new offsets together
Kafka Streams, Connect, ksqlDB¶
| Layer | Built on | Purpose |
|---|---|---|
| Kafka Streams | Producer, consumer, and RocksDB local stores | Embedded JVM stream-processing library: KStream/KTable DSL, windows, joins, exactly-once through transactions. |
| Kafka Connect | Distributed worker framework | Runs source connectors (for example Debezium MySQL/Postgres CDC) and sink connectors (S3, Iceberg, Snowflake, Elasticsearch). Keeps its state in three compacted topics: connect-configs, connect-offsets, and connect-status. |
| ksqlDB | Kafka Streams plus a REST API | SQL-like continuous queries over topics. A Confluent product under the Confluent Community License, not part of Apache Kafka. |
Tiered Storage (KIP-405)¶
Tiered storage was early access in 3.6 and became production-ready in 3.9. Brokers keep only recent segments on local disk. Closed segments are copied to a remote store through a pluggable RemoteStorageManager, and segment metadata goes to a RemoteLogMetadataManager, by default the topic-based one backed by __remote_log_metadata.
The flow shows how segments move between tiers and how a consumer reads old data.
flowchart TB
subgraph LocalTier["Local tier (broker disk)"]
Active["Active segment (writes)"]
Hot["Closed local segments<br/>(page cache hits)"]
end
subgraph RemoteTier["Remote tier (object store)"]
Cold["Uploaded segments + indexes"]
end
subgraph Metadata["Remote log metadata"]
RLMM["__remote_log_metadata"]
end
RLM["RemoteLogManager"]
Consumer["Consumer reading old offsets"]
Active -->|roll| Hot
Hot -->|"copy closed segment"| RLM
RLM -->|"RemoteStorageManager put"| Cold
RLM -->|"segment metadata"| RLMM
Hot -->|"delete after local.retention.ms"| Deleted["Local copy removed"]
Consumer -->|"Fetch (offset below local start)"| RLM
RLM -->|"RemoteStorageManager get"| Cold
RLM -->|"stream records"| Consumer
Points that matter in practice:
- Apache Kafka ships no production
RemoteStorageManager. The built-inLocalTieredStorageis a test implementation. Production clusters use a plugin, such as Aiven's open-source tiered-storage plugin for S3, GCS, and Azure, or a vendor distribution (Confluent Platform, Amazon MSK tiered storage). - Configuration. The broker switch is
remote.log.storage.system.enable=true. Per topic you setremote.storage.enable=true,local.retention.ms/local.retention.bytesfor the local tier, andretention.ms/retention.bytesfor the total. See How-to Guides. - Limitations. Compacted topics are not supported. You must disable tiered storage on every topic before disabling it on the broker. Segments without a producer snapshot file (topics created before 2.8) are not supported.
- Recent improvements. 3.9 added per-topic disablement (KIP-950) and quotas (KIP-956). 4.3 added follower bootstrap from the last tiered offset (KIP-1023), and the 4.4 upgrade notes add delayed upload (
remote.copy.lag.ms/remote.copy.lag.bytes, KIP-1241).
Request Flow Examples¶
Produce Path¶
The sequence follows one produce request from send() to the callback.
sequenceDiagram
participant App
participant KP as KafkaProducer
participant Net as Sender thread
participant SS as SocketServer
participant KA as KafkaApis
participant RM as ReplicaManager
participant LM as UnifiedLog
participant Disk as LogSegment
App->>KP: send(ProducerRecord)
KP->>KP: Serialize, partition, append to batch in RecordAccumulator
KP->>KP: Assign producer ID and sequence (idempotence)
Net->>SS: ProduceRequest when batch is full or linger.ms expires
SS->>KA: Dispatch through request channel
KA->>RM: appendRecords(acks=all)
RM->>LM: Append to leader log
LM->>Disk: Write to page cache (fsync per flush policy)
RM-->>KA: Delayed until ISR has replicated (acks=all)
KA-->>Net: ProduceResponse with base offset
Net-->>KP: Complete future
KP-->>App: Callback success
Fetch Path (with Zero-Copy)¶
The sequence shows a consumer fetch served with sendfile.
sequenceDiagram
participant KC as KafkaConsumer
participant SS as SocketServer
participant KA as KafkaApis
participant RM as ReplicaManager
participant LM as UnifiedLog
participant Kernel as Linux kernel sendfile
KC->>SS: FetchRequest (topic, partition, offset)
SS->>KA: Dispatch
KA->>RM: fetchMessages(maxWait, minBytes)
RM->>LM: read(offset, maxBytes)
LM-->>RM: FileRecords slice of the segment file
RM-->>KA: FetchResponse referencing FileRecords
KA->>Kernel: FileChannel.transferTo to the socket
Note over Kernel: Bytes go from page cache to NIC without a user-space copy
Kernel-->>KC: Response bytes
TLS disables zero-copy
On SSL and SASL_SSL listeners the broker must encrypt the data in the JVM, so it reads the segment into user space and cannot use sendfile. Expect higher broker CPU with TLS. That is the cost of encryption in transit, not a reason to skip it.
Security Model¶
Threat Model¶
| Threat vector | Impact | Mitigation |
|---|---|---|
| Broker compromise | Direct read of all topic data on disk. Ability to forge produce and fetch responses | Disk encryption (LUKS, EBS encryption). Restrict OS access. Audit the authorizer log. Network isolation |
| KRaft controller compromise | Forge metadata records (create topics, alter ACLs), rewire partition leadership | Run controllers in isolated mode on hardened hosts. Use TLS on the controller listener. Keep controller.listener.names on internal-only network paths |
| Client impersonation | A producer writes records as another principal, or a consumer reads topics it should not | SASL or mTLS authentication on every listener. ACLs with default deny |
| Data exfiltration by a consumer | An authenticated consumer subscribes to sensitive topics | Topic-scoped ACLs (Read on Topic:payments for named principals only). Per-user consumer_byte_rate quotas. Egress monitoring |
| Man-in-the-middle on the wire | Eavesdropping on produce and fetch traffic, or downgrade to PLAINTEXT | TLS on every listener (SSL or SASL_SSL). No PLAINTEXT listeners in production. TLS 1.2 or later |
| Replay of captured requests | Duplicate writes, or replay of a SCRAM exchange | Idempotent and transactional producers deduplicate and fence. SCRAM only over TLS (CVE-2024-56128) |
| JAAS login-module abuse (CVE-2025-27818, CVE-2025-27819) | A principal with AlterConfigs (or Connect REST access) triggers RCE through LdapLoginModule or JndiLoginModule |
Run 3.9.1/4.0.0 or later, which disable those modules by default. On 4.2+, use org.apache.kafka.allowed.login.modules. Restrict AlterConfigs |
| OAUTHBEARER URL abuse (CVE-2025-27817) | Arbitrary file read or SSRF through sasl.oauthbearer.token.endpoint.url |
Run 3.9.1/4.0.0 or later and set the JVM property org.apache.kafka.sasl.oauthbearer.allowed.urls |
| Unvalidated JWTs (CVE-2026-33557) | On 4.1.0 and 4.1.1 the default broker JWT validator accepts forged tokens | Upgrade to 4.1.2/4.2.0 or set sasl.oauthbearer.jwt.validator.class to BrokerJwtValidator |
| MirrorMaker 2 credential leak | A compromised replicator can read the source and write the target | Separate principals for source and target. Least-privilege ACLs on MM2 internal topics. TLS to both clusters |
| Sensitive logging (CVE-2026-33558) | Credentials and config payloads in DEBUG logs | Keep NetworkClient and the root logger at INFO in production |
Authentication¶
Each listener picks a transport (security.protocol: PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL) and, for SASL, one or more mechanisms. With mTLS (SSL plus ssl.client.auth=required), the certificate's subject DN, optionally rewritten by ssl.principal.mapping.rules, becomes the principal. The mechanisms trade off as follows.
| Mechanism | When to use | Strengths | Weaknesses |
|---|---|---|---|
| PLAIN | Labs, or behind TLS with an external credential store | Simple and widely supported | Sends the password to the broker, so it must run inside TLS |
| SCRAM-SHA-256 / SCRAM-SHA-512 | Username and password without an external identity provider | Challenge-response. Credentials live in the KRaft metadata | Must run over TLS (see CVE-2024-56128) |
| GSSAPI (Kerberos) | Enterprises with Active Directory or an MIT/Heimdal KDC | Strong mutual authentication | Heavy client setup and keytab management |
| OAUTHBEARER | Microservices with a corporate identity provider (Keycloak, Okta, Entra ID) | Short-lived tokens, OIDC integration | The built-in default creates and accepts unsecured JWTs, so it is for development only. Production needs the JWKS-based validator with issuer and audience checks |
Setup steps are in How-to Guides: Secure a Cluster.
Authorization¶
With an authorizer configured, Kafka is default-deny: a resource with no matching ACL is accessible only to super.users, unless you set allow.everyone.if.no.acl.found=true. The KRaft-native authorizer is org.apache.kafka.metadata.authorizer.StandardAuthorizer, which stores ACLs in the metadata log. The ZooKeeper-based AclAuthorizer was removed in 4.0. ACLs combine a principal, host, operation (Read, Write, Describe, Create, Alter, and others), and resource (Topic, Group, Cluster, TransactionalId, DelegationToken), with literal or prefixed resource patterns.
For policy-as-code, community authorizer plugins such as the open-source opa-kafka-plugin wrap the authorizer interface and delegate decisions to Open Policy Agent. That allows attribute-based rules that plain ACLs cannot express. Strimzi deprecated its built-in OPA authorization type in 0.46, so check operator support before adopting it.
Encryption at Rest¶
Apache Kafka has no built-in broker-side record encryption. The usual layers are:
| Layer | Approach | Notes |
|---|---|---|
| Disk | LUKS or cloud volume encryption (EBS, PD) | Transparent to the broker. Standard practice |
| Application | Client-side envelope encryption with a KMS or Vault Transit | The producer encrypts before sending and the consumer decrypts after receiving. Agree on a header for the key ID |
| Proxy | Kroxylicious record-encryption filter | A Kafka-protocol proxy encrypts record values with KMS-managed keys, with no client changes |
| Kafka Connect | Single Message Transforms for field-level encryption | Community SMTs with KMS-backed data keys |
| Confluent Platform | Client-Side Field Level Encryption (CSFLE) | Commercial. Tied to Schema Registry tags |
An end-to-end encryption proposal (KIP-317) was discussed but never implemented. As of 4.3, disk encryption plus client-side or proxy encryption for sensitive fields is the standard layered defense.
Audit Logging¶
Apache Kafka has no dedicated audit subsystem. The de facto audit trail is the kafka.authorizer.logger Log4j2 logger: StandardAuthorizer logs every denied request at INFO and every allowed request at DEBUG (high volume). Ship kafka-authorizer.log to a SIEM and alert on bursts of denials. For record-level audit (who consumed what), use Confluent Platform audit logs or a proxy layer such as Kroxylicious or Conduktor Gateway. The log4j2 configuration is in How-to Guides.
Sources¶
- Apache Kafka design documentation
- Apache Kafka: Eligible Leader Replicas
- Apache Kafka: KRaft operations
- Apache Kafka: Tiered storage
- KIP-405: Tiered Storage
- KIP-848: The Next Generation of the Consumer Rebalance Protocol
- KIP-853: KRaft Controller Membership Changes
- KIP-932: Queues for Kafka
- KIP-966: Eligible Leader Replicas
- KIP-1071: Streams Rebalance Protocol
- Apache Kafka 4.0.0 release announcement
- Apache Kafka 3.9.0 release announcement
- Apache Kafka: Security documentation
- Apache Kafka: Authentication using SASL
- Apache Kafka: Authorization and ACLs
- Apache Kafka CVE list
- Strimzi CHANGELOG (OPA authorization deprecation)