Introduction
Apache Kafka is a distributed streaming platform that has become the de facto standard for building real-time data pipelines and streaming applications. Its architecture is designed to provide high throughput, low latency, fault tolerance, and horizontal scalability. Understanding Kafka's core components and the data flow between them is essential for anyone tasked with operating, monitoring, or troubleshooting Kafka clusters in production.
This guide explains Kafka architecture through the lens of a technical operator. Each section pairs a core concept with practical, hands-on commands and configuration examples. You will learn how to inventory your Kafka environment, safely change configurations, verify system health, diagnose common failures, and follow an operations checklist that prioritizes safety and reversibility.
We will use a fictional company, StreamCo, to illustrate real-world scenarios. StreamCo runs an e-commerce platform and uses Kafka to handle order events, inventory updates, and customer notifications. Their Kafka cluster is version 3.4.0, deployed across three data centers with three brokers per data center.
Before diving into the details, let's establish the foundational components of Kafka architecture.
Core Components of Kafka Architecture
Kafka's architecture is built around a few key components that work together to provide a scalable and fault-tolerant messaging system.
Brokers and Clusters
A Kafka broker is a server process that stores data and serves client requests. A Kafka cluster consists of one or more brokers. In a production environment, you typically have at least three brokers to ensure availability and fault tolerance.
Each broker is identified by a unique integer ID. For example, in StreamCo's environment, the broker IDs are 0, 1, and 2 in the primary data center.
Brokers receive messages from producers, assign them offsets, and commit them to storage on disk. They also serve consumer requests by delivering messages from committed log segments.
Practical command: To list the active brokers in a cluster, you can use the kafka-broker-api-versions.sh utility that ships with Kafka:
kafka-broker-api-versions.sh --bootstrap-server broker1.streamco.example:9092
This command queries the cluster metadata and returns a list of brokers with their supported API versions. It is a read-only operation that does not alter cluster state.
Topics and Partitions
A topic is a logical channel to which producers publish messages and from which consumers read. Topics are split into partitions for scalability and parallelism. Each partition is an ordered, immutable sequence of records that is continually appended to a commit log.
Partitions allow Kafka to distribute data across multiple brokers. For example, a topic orders might have 6 partitions. These partitions can be spread across the three brokers in the primary data center, with each broker hosting two partitions.
The number of partitions for a topic is a critical design decision. More partitions allow higher parallelism in production and consumption but increase overhead on the cluster. A common guideline is to set partitions based on the expected throughput and the number of consumers in a consumer group.
Practical command: To create a topic with a specified number of partitions and replication factor, use kafka-topics.sh:
kafka-topics.sh --create \
--bootstrap-server broker1.streamco.example:9092 \
--topic orders \
--partitions 6 \
--replication-factor 3
This command creates a topic named orders with 6 partitions and a replication factor of 3, meaning each partition is replicated across three brokers.
Producers and Consumers
A producer is an application that publishes (writes) messages to Kafka topics. A consumer is an application that subscribes to topics and processes the published messages.
Producers can choose which partition to send a message to. By default, if no key is specified, the producer uses a round-robin strategy to distribute messages evenly across all partitions of the topic. If a key is specified, Kafka uses the key's hash to determine the partition, ensuring that all messages with the same key go to the same partition, which preserves ordering per key.
Consumers can form consumer groups. In a consumer group, each partition is consumed by exactly one consumer. This allows for parallel consumption and load balancing. If a consumer fails, another consumer in the group takes over its partitions.
Practical configuration snippet: Here is a minimal Java producer configuration using the Kafka client library:
Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
This example uses string serializers for keys and values. In production, you might use Avro or Protobuf for schema management.
The Role of ZooKeeper
Historically, Kafka relied on Apache ZooKeeper for cluster coordination, leader election, and metadata management. ZooKeeper maintains the list of brokers, topic configurations, and partition leadership information.
As of Kafka 3.3, a new KRaft mode (Kafka Raft metadata mode) is available, which replaces ZooKeeper with a built-in consensus protocol. KRaft simplifies deployment and improves scalability. However, many production clusters still run with ZooKeeper.
For StreamCo, their Kafka 3.4 cluster still uses ZooKeeper, so it is important to monitor ZooKeeper health.
Practical command: To check the status of ZooKeeper, you can use the echo srvr command via nc:
echo srvr | nc zookeeper1.streamco.example 2181
This returns a line like Mode: follower or Mode: leader, indicating the ZooKeeper node's role. A cluster with an odd number of nodes ensures a quorum.
Data Flow in Kafka: From Producer to Consumer
Understanding the path a message takes from producer to consumer is crucial for diagnosing latency and throughput issues.
- Producer sends a message: The producer serializes the message and sends it to the leader partition of the appropriate topic partition. The leader is the broker that handles all reads and writes for that partition.
- Leader appends to log: The leader appends the message to its local log and assigns a sequential offset to it.
- Replication to followers: The leader replicates the message to follower replicas (if replication factor > 1). Once the message is replicated to the required number of replicas (as per
ackssetting), the leader acknowledges the producer. - Consumer reads: Consumers poll for messages from the leader. They track their offset in the partition and can commit offsets to Kafka or externally.
Example with acks settings:
acks=0: Producer does not wait for any acknowledgment. Fastest but no delivery guarantee.acks=1: Producer waits for leader acknowledgment only. The leader writes to its local log but does not wait for followers. This is the default.acks=all(oracks=-1): Producer waits for all in-sync replicas to acknowledge. This provides the strongest durability guarantee.
For StreamCo's order handling, they set acks=all to ensure no order loss even if a broker fails.
Version and Environment Inventory
Before making any changes to a Kafka cluster, you must understand the exact versions, deployment topology, and current state. This is the first step in safe operations.
Identify Installed Version
Use the following command to find the Kafka broker version:
kafka-broker-api-versions.sh --bootstrap-server broker1.streamco.example:9092 | grep -i version
For StreamCo, this returns something like broker1:9092 (id: 0 rack: null) -> 3.4.0. Note the broker ID and the rack (if configured). The rack information is useful for ensuring replicas are spread across different failure domains.
Additionally, the Kafka distribution includes a kafka-configs.sh script that can retrieve broker configuration. For example, to check the log.retention.hours setting:
kafka-configs.sh --bootstrap-server broker1.streamco.example:9092 \
--entity-type brokers \
--entity-name 0 \
--describe
This returns all broker configurations for broker 0. It is a read-only operation.
Understand Deployment Topology
Map out your cluster: how many data centers, how many brokers per data center, how many partitions per topic, and the replication factor. For StreamCo:
- Primary data center: 3 brokers (IDs 0,1,2)
- Secondary data center: 3 brokers (IDs 3,4,5)
- Disaster recovery data center: 3 brokers (IDs 6,7,8)
They use a multi-region replication setup with custom partition assignment to ensure replicas are in different data centers.
Capture Current State
Before any change, capture a snapshot of relevant metrics and configurations. For example, to get the current partition leadership distribution:
kafka-topics.sh --describe --bootstrap-server broker1.streamco.example:9092 --topic orders
Expected output snippet:
Topic: orders Partition: 0 Leader: 0 Replicas: 0,3,6 Isr: 0,3,6
Topic: orders Partition: 1 Leader: 1 Replicas: 1,4,7 Isr: 1,4,7
...
This shows which broker is the leader for each partition, the replicas, and the in-sync replicas (ISR). If a replica is not in the ISR, it is lagging or offline.
Safe Configuration Path
Changing Kafka configurations can have significant impact. Follow a safe path: observe, make a minimal change, verify, and have a rollback plan.
Observation First
Always start with read-only commands. For example, to check the current log retention setting for a topic:
kafka-configs.sh --bootstrap-server broker1.streamco.example:9092 \
--entity-type topics \
--entity-name orders \
--describe
This returns all topic-level configurations. Suppose the output shows retention.ms=604800000 (7 days). StreamCo wants to increase retention to 14 days to allow offline analysis.
Minimal Justified Change
Use a dynamic configuration change to update the topic's retention without restarting brokers:
kafka-configs.sh --bootstrap-server broker1.streamco.example:9092 \
--entity-type topics \
--entity-name orders \
--alter \
--add-config retention.ms=1209600000
This command changes the retention to 14 days and applies it dynamically. The change does not require a broker restart.
Verification and Rollback
After the change, verify that the new configuration is active:
kafka-configs.sh --bootstrap-server broker1.streamco.example:9092 \
--entity-type topics \
--entity-name orders \
--describe | grep retention.ms
Expected output: retention.ms=1209600000.
If the change causes unexpected disk usage, you can revert to the original value by repeating the --alter command with retention.ms=604800000.
Blast radius: This change only affects the orders topic. Other topics retain their original retention settings.
Verification and Diagnostics
Effective verification involves checking cluster health, broker status, and consumer lag.
Cluster Health Check
Use the kafka-broker-api-versions.sh command to quickly verify that all brokers are reachable:
kafka-broker-api-versions.sh --bootstrap-server broker1.streamco.example:9092
If a broker is down, you'll see a connection error for that broker. Additionally, you can use Kafka's JMX metrics to monitor broker health. For example, the metric kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions indicates the number of partitions that are under-replicated. This should ideally be zero.
Using JMX with jconsole:
jconsole broker1.streamco.example:9999
Navigate to MBeans and expand kafka.server -> ReplicaManager -> UnderReplicatedPartitions.
Consumer Lag Monitoring
Consumer lag is the difference between the latest offset in a partition and the consumer's current offset. High lag indicates consumers are falling behind.
Use the kafka-consumer-groups.sh tool:
kafka-consumer-groups.sh --bootstrap-server broker1.streamco.example:9092 \
--group order-processor \
--describe
Output:
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID
order-processor orders 0 10234 11200 966 consumer-1-a1c /192.168.1.10 consumer-1
order-processor orders 1 10450 11500 1050 consumer-1-b2 /192.168.1.10 consumer-1
Lag values above a threshold (e.g., 1000) may trigger alerts. StreamCo monitors lag with a custom Prometheus exporter.
Failure Modes and Recovery
Kafka is designed for fault tolerance, but failures still occur. Here are common failure scenarios and how to recover.
Broker Failure
Symptom: A broker goes offline. Partitions that had leadership on that broker will elect new leaders from the ISR. If the ISR contains other replicas, there may be brief unavailability during leader election.
Detection: Use kafka-broker-api-versions.sh to see if a broker is unreachable. Check the broker logs for FATAL or ERROR messages.
Recovery:
- Determine the cause (hardware failure, network partition, etc.).
- If the broker can be restarted, restart it and watch for it to rejoin the cluster. It will catch up on any missed data.
- If the broker cannot be restarted, replace it. Add a new broker with a different ID, and reassign partitions if necessary.
For StreamCo, when broker 2 in the primary data center failed due to a disk issue, they used the following command to see partition status:
kafka-topics.sh --describe --bootstrap-server broker1.streamco.example:9092 --topic orders
They noticed partition 2 had leader -1 (no leader) because all replicas were on the failed broker. Since the replication factor was 3, the other replicas on different brokers should have taken over, but in this case, the ISR was empty because the followers were also on the same rack and had lost connectivity. After restarting the broker, the ISR repopulated and leadership was restored.
Under-Replicated Partitions
Symptom: The metric UnderReplicatedPartitions is non-zero for an extended period. This means some partitions have fewer replicas than the configured replication factor.
Cause: A broker is down or a follower is lagging too far behind.
Recovery:
- If a broker is down, bring it back up.
- If a follower is lagging, check network connectivity and disk I/O on that broker.
- You can also use the
kafka-reassign-partitions.shtool to move partitions to healthy brokers.
Example reassignment:
- Create a JSON file
reassignment.jsonwith new assignments:
{
"version": 1,
"partitions": [
{"topic": "orders", "partition": 0, "replicas": [1,4,7]}
]
}
- Run the reassignment:
kafka-reassign-partitions.sh --bootstrap-server broker1.streamco.example:9092 \
--reassignment-json-file reassignment.json \
--execute
- Verify with the
--verifyflag.
Consumer Group Rebalance Issues
Symptom: Consumers in a group are constantly rebalancing, causing processing delays.
Cause: This often occurs when consumers take too long to process messages, exceeding max.poll.interval.ms, or when there are frequent membership changes.
Recovery:
- Increase
max.poll.interval.msandsession.timeout.msappropriately. - Ensure consumer processing time is within the poll interval.
- Check for frequent consumer restarts due to exceptions.
Example consumer configuration fix:
max.poll.interval.ms=600000
session.timeout.ms=45000
heartbeat.interval.ms=15000
Operations Checklist
Use this checklist for routine operations and when making changes to your Kafka cluster.
| Check Item | Owner | Frequency | Command / Metric | Expected Result |
|---|---|---|---|---|
| Verify all brokers are up | Priya Shah, Kafka Admin | Daily | kafka-broker-api-versions.sh --bootstrap-server broker1:9092 | All broker IDs listed without errors |
| Check under-replicated partitions | Priya Shah | Hourly | JMX UnderReplicatedPartitions | 0 |
| Monitor consumer lag | Alex Chen, Data Engineer | Hourly | kafka-consumer-groups.sh --group order-processor --describe | Lag < 1000 for all partitions |
| Review disk usage | Priya Shah | Weekly | df -h on broker hosts | < 70% disk usage |
| Backup ZooKeeper data | Ravi Patel, DevOps | Daily | Use ZooKeeper snapshot export | Successful export |
| Test failover | Priya Shah and Ravi Patel | Quarterly | Manually stop a broker and observe | Leadership transfers within 30 seconds; no data loss |
| Review topic configurations | Priya Shah | Monthly | kafka-configs.sh --entity-type topics --describe | Configs align with business requirements |
Owner accountability: Priya Shah, as the Kafka Admin, is the primary accountable owner for cluster health. She reviews the checklist outcomes weekly with the engineering team and adjusts monitoring thresholds as needed.
Revisit frequency: The checklist itself is reviewed quarterly to incorporate new learnings and changes in the system.
Common Pitfalls and How to Avoid Them
1. Insufficient Partition Count
Problem: Creating topics with too few partitions limits parallelism and throughput. For example, if a topic has only one partition, only one consumer in a group can consume it, even if there are many consumers.
How to recognize: High consumer lag even with many consumers, or producer throughput plateaus.
How to avoid: Before creating a topic, estimate throughput and consumer parallelism. Use a partition count that is a multiple of the number of brokers and allows for future growth. For StreamCo, they calculated that the orders topic needed 6 partitions to handle 10,000 messages/sec with 3 consumers, allowing 3x parallelism.
Recovery: You can increase partitions later with kafka-topics.sh --alter --partitions N, but be aware that increasing partitions changes key partitioning and may affect ordering guarantees for existing keys.
2. Misconfigured Replication Factor
Problem: Setting replication.factor=1 means data is lost if the broker fails. This is a common mistake in development that accidentally gets carried to production.
How to recognize: Check with kafka-topics.sh --describe and see ReplicationFactor: 1.
How to avoid: Always set replication factor to at least 3 in production, and ensure brokers are spread across different racks or availability zones.
Recovery: If you have a topic with RF=1, you can use kafka-reassign-partitions.sh to increase the replication factor, but you need to ensure the cluster has enough brokers to host the replicas.
3. Ignoring ZooKeeper Health
Problem: Even with KRaft, many clusters still depend on ZooKeeper. If ZooKeeper goes down, Kafka brokers may lose coordination and fail to elect leaders.
How to recognize: Brokers start showing errors like Unable to connect to zookeeper in logs. Brokers may shut down if connection is lost for too long.
How to avoid: Monitor ZooKeeper with echo srvr | nc zookeeper1 2181 and set up alerts. Ensure ZooKeeper cluster has an odd number of nodes (3 or 5) for quorum.
Recovery: Restart ZooKeeper nodes one at a time to restore quorum. Then restart brokers if necessary.
4. Not Monitoring Consumer Lag
Problem: Consumer lag can grow silently, leading to processing delays that affect downstream systems.
How to recognize: Downstream applications report stale data. Use kafka-consumer-groups.sh to see high lag values.
How to avoid: Set up lag monitoring with Prometheus and Grafana. Alert when lag exceeds a threshold for more than a few minutes.
Recovery: Scale up consumers by adding more instances to the group (if partitions allow) or optimize consumer processing logic.
5. Dynamic Config Changes Without Verification
Problem: Changing configurations like retention.ms without verifying the effect can lead to unintended data loss or disk full errors.
How to recognize: After a change, unexpected alerts or user reports of missing data.
How to avoid: Always follow the safe configuration path: observe, change minimally, verify, and have a rollback plan.
Recovery: Revert the configuration change immediately and investigate the root cause.
Conclusion
Kafka's architecture is robust but requires careful operation to maintain reliability and performance. By understanding the core components, data flow, and operational best practices, you can safely manage Kafka clusters in production.
Start by inventorying your environment and mastering read-only diagnostic commands. When changes are necessary, follow a safe path: observe, make a minimal change, verify, and have a rollback plan. Use the operations checklist to assign accountability and ensure regular monitoring. Avoid common pitfalls by designing topics with appropriate partition counts and replication factors, and by monitoring consumer lag and cluster health diligently.
For StreamCo, implementing these practices reduced unexpected downtime by 40% and improved the team's ability to respond to incidents quickly. Apply these principles to your own Kafka environment to achieve similar results.
As a next step, choose one low-risk verification from this guide, such as checking under-replicated partitions, and integrate it into your daily routine. Build from there.