Skip to content
LowLevelDesign Mastery

Kafka Deep Dive

Distributed streaming platform for high-throughput event processing

Apache Kafka is a distributed streaming platform designed for building real-time data pipelines and streaming applications. Originally developed by LinkedIn and open-sourced in 2011, Kafka has become the de facto standard for event streaming in modern distributed systems.

Kafka operates as a distributed commit log - a persistent, append-only data structure that stores streams of records (events) in topics. Unlike traditional message queues that delete messages after consumption, Kafka retains messages for a configurable retention period, allowing multiple consumers to read the same messages at different times and speeds.

This design makes Kafka ideal for scenarios where you need to:

  • Decouple producers and consumers (they don’t need to be active simultaneously)
  • Replay events (consumers can read historical data)
  • Scale horizontally (add partitions and consumers independently)
  • Guarantee ordering within partitions
  • Ensure durability and fault tolerance through replication
Producers publishing to a Kafka cluster of three brokers that replicate data, with consumers subscribing to read events
  1. Distributed - Runs on multiple servers (brokers)
  2. Fault-tolerant - Replicates data across brokers
  3. High-throughput - Millions of messages per second
  4. Scalable - Add brokers/partitions to scale
  5. Durable - Messages persisted to disk
  6. Real-time - Low latency streaming

A topic is a category or stream of messages in Kafka. Think of it as a named channel where producers publish messages and consumers subscribe to read them. Topics are similar to tables in a database or folders in a filesystem - they provide logical organization for related messages.

Key Characteristics:

Topics in Kafka are immutable logs - messages are appended to the end and never modified or deleted (until retention expires). This append-only design provides several benefits: it enables efficient sequential disk I/O, allows multiple consumers to read the same messages independently, and supports event replay for debugging or reprocessing.

Each topic is partitioned for scalability and replicated across multiple brokers for fault tolerance. Messages within a partition are strictly ordered, but ordering across partitions is not guaranteed unless you use a partitioning key.

Kafka topics user-events, orders and payments, with producers publishing to and consumers subscribing to the user-events topic

Characteristics:

  • Immutable log - Messages appended, never modified
  • Partitioned - Split into partitions for parallelism
  • Replicated - Copies across brokers for reliability
  • Ordered - Messages ordered within partition

A partition is an ordered, immutable sequence of messages within a topic. Each topic is divided into one or more partitions, which enables Kafka to scale horizontally and process messages in parallel.

Why Partitions Matter:

Partitions are Kafka’s fundamental unit of parallelism. Without partitions, a topic would be processed by only one consumer at a time, creating a bottleneck. With partitions, different consumers can process different partitions simultaneously, dramatically increasing throughput.

Partitioning Strategy:

When a producer sends a message, Kafka determines which partition to use based on:

  • Key-based partitioning: If a message has a key, Kafka uses a hash of the key to determine the partition. This ensures all messages with the same key go to the same partition, maintaining ordering for that key.
  • Round-robin partitioning: If no key is provided, Kafka distributes messages evenly across partitions in round-robin fashion.

Real-World Example: In an e-commerce system, you might partition the orders topic by customer_id. This ensures all orders from the same customer are processed in order, while orders from different customers can be processed in parallel.

Kafka topic user-events split into three ordered partitions, each read in parallel by its own consumer

Why Partitions?

  • Parallelism - Multiple consumers process different partitions
  • Scalability - Add partitions to scale throughput
  • Ordering - Messages ordered within partition (not globally)

Partitioning Strategy:

  • Key-based - Same key → same partition (ensures ordering)
  • Round-robin - Distribute evenly (no key)

A consumer group is a set of consumers that work together to consume messages from one or more topics. Kafka automatically distributes partitions across consumers in the same group, ensuring each partition is consumed by exactly one consumer in the group.

How Consumer Groups Work:

When consumers in a group subscribe to a topic, Kafka performs partition assignment - it divides the topic’s partitions among the available consumers. If you have 3 partitions and 2 consumers, Kafka might assign partitions 0 and 1 to consumer 1, and partition 2 to consumer 2.

Rebalancing:

When consumers join or leave a consumer group, Kafka automatically rebalances - it redistributes partitions among the remaining consumers. During rebalancing, consumers stop processing messages, which can cause a brief pause. This is why it’s important to design consumers to handle rebalancing gracefully and commit offsets frequently.

Scaling Pattern:

  • Scale consumption: Add more consumers to the group (up to the number of partitions)
  • Scale production: Add more partitions to the topic (requires careful planning as partitions can’t be reduced)
Kafka consumer group order-processors with three consumers, each assigned one partition of the three-partition orders topic

Key Rules:

  • One partition → one consumer (in same group)
  • One consumer → multiple partitions (can handle multiple)
  • Rebalancing - When consumer joins/leaves, partitions redistributed

Example:

  • Topic has 3 partitions
  • Consumer group has 2 consumers
  • Consumer 1 gets partitions 0, 1
  • Consumer 2 gets partition 2

An offset is a sequential number that uniquely identifies the position of a message within a partition. Offsets start at 0 and increment for each message. They are immutable - once assigned, an offset never changes.

Offset Management:

Consumers track their current offset - the position of the last message they’ve successfully processed. After processing a message, the consumer commits the offset, which tells Kafka “I’ve processed up to this point.” If the consumer crashes and restarts, it resumes from the last committed offset, ensuring no messages are lost and no messages are processed twice (assuming proper offset management).

Offset Commit Strategies:

  • Automatic commit: Kafka commits offsets periodically (default: every 5 seconds). Simple but can lead to duplicate processing if consumer crashes between processing and commit.
  • Manual commit: Consumer explicitly commits offsets after processing. More control, but requires careful implementation to avoid duplicates or lost messages.

Offset Storage:

Kafka stores committed offsets in a special internal topic called __consumer_offsets. This allows Kafka to track consumer progress even if consumers restart or rebalance.

Kafka partition with messages at offsets 0 to 3 and a consumer at offset 1 reading the next offset and committing its progress

Offset Management:

  • Consumer tracks current offset
  • After processing, commits offset
  • On restart, resumes from committed offset

Kafka cluster of three brokers holding leader and follower partition replicas, coordinated by Zookeeper or KRaft for metadata

Broker = Kafka server. Cluster = Multiple brokers.

Replication:

  • Leader - Handles reads/writes for partition
  • Followers - Replicate leader’s data
  • ISR (In-Sync Replicas) - Followers in sync with leader

💡 Tip: Click dropdown to switch between languages
"kafka_producer.py
from kafka import KafkaProducer
import json
class KafkaEventProducer:
"""Kafka producer for events"""
def __init__(self, bootstrap_servers: list):
self.producer = KafkaProducer(
bootstrap_servers=bootstrap_servers,
value_serializer=lambda v: json.dumps(v).encode('utf-8'),
key_serializer=lambda k: k.encode('utf-8') if k else None,
# Idempotent producer (exactly-once)
enable_idempotence=True,
acks='all', # Wait for all replicas
retries=3,
max_in_flight_requests_per_connection=1
)
def send_event(self, topic: str, event: dict, key: str = None):
"""Send event to topic"""
future = self.producer.send(
topic,
key=key,
value=event
)
# Wait for acknowledgment
try:
record_metadata = future.get(timeout=10)
print(f"Event sent to {record_metadata.topic} "
f"partition {record_metadata.partition} "
f"offset {record_metadata.offset}")
return record_metadata
except Exception as e:
print(f"Error sending event: {e}")
raise
def send_with_callback(self, topic: str, event: dict, key: str = None):
"""Send event with callback"""
def on_send_success(record_metadata):
print(f"Event sent: {record_metadata.topic}/"
f"{record_metadata.partition}/"
f"{record_metadata.offset}")
def on_send_error(exception):
print(f"Error sending event: {exception}")
self.producer.send(
topic,
key=key,
value=event
).add_callback(on_send_success).add_errback(on_send_error)
def flush(self):
"""Flush pending messages"""
self.producer.flush()
def close(self):
"""Close producer"""
self.producer.close()
# Usage
producer = KafkaEventProducer(['localhost:9092'])
# Send event
producer.send_event('user-events', {
'event_type': 'user.created',
'user_id': 123,
'email': '[email protected]',
'timestamp': '2024-01-01T10:00:00Z'
}, key='123') # Key ensures same user goes to same partition
producer.flush()
producer.close()

💡 Tip: Click dropdown to switch between languages
"kafka_consumer.py
from kafka import KafkaConsumer
from kafka.errors import KafkaError
import json
from typing import Callable
class KafkaEventConsumer:
"""Kafka consumer for events"""
def __init__(self, bootstrap_servers: list, group_id: str):
self.consumer = KafkaConsumer(
bootstrap_servers=bootstrap_servers,
group_id=group_id,
value_deserializer=lambda m: json.loads(m.decode('utf-8')),
key_deserializer=lambda k: k.decode('utf-8') if k else None,
enable_auto_commit=False, # Manual offset commit
auto_offset_reset='earliest', # Start from beginning if no offset
max_poll_records=100 # Batch size
)
def subscribe(self, topics: list):
"""Subscribe to topics"""
self.consumer.subscribe(topics)
print(f"Subscribed to topics: {topics}")
def consume(self, handler: Callable):
"""Consume messages"""
try:
while True:
# Poll for messages (batch)
message_batch = self.consumer.poll(timeout_ms=1000)
for topic_partition, messages in message_batch.items():
for message in messages:
try:
# Process message
handler(message.value, message.key, message.offset)
# Commit offset after processing
self.consumer.commit()
except Exception as e:
print(f"Error processing message: {e}")
# Don't commit - will retry
except KeyboardInterrupt:
print("Stopping consumer...")
finally:
self.consumer.close()
def consume_with_manual_commit(self, handler: Callable):
"""Consume with manual offset commit"""
try:
while True:
message_batch = self.consumer.poll(timeout_ms=1000)
offsets_to_commit = {}
for topic_partition, messages in message_batch.items():
for message in messages:
try:
# Process message
handler(message.value, message.key, message.offset)
# Track offset for commit
offsets_to_commit[topic_partition] = \
OffsetAndMetadata(message.offset + 1, None)
except Exception as e:
print(f"Error processing message: {e}")
# Don't commit failed messages
break
# Commit all processed offsets
if offsets_to_commit:
self.consumer.commit(offsets_to_commit)
except KeyboardInterrupt:
print("Stopping consumer...")
finally:
self.consumer.close()
# Usage
def handle_event(event, key, offset):
print(f"Processing event: {event['event_type']} "
f"key: {key} offset: {offset}")
# Process event...
consumer = KafkaEventConsumer(
['localhost:9092'],
group_id='event-processors'
)
consumer.subscribe(['user-events', 'order-events'])
consumer.consume(handle_event)

Exactly-once semantics ensures that each message is processed exactly once, with no duplicates and no lost messages. This is the strongest delivery guarantee but also the most complex to implement.

In distributed systems, achieving exactly-once processing is challenging because:

  1. Network failures can cause retries, leading to duplicate messages
  2. Consumer failures can cause rebalancing, leading to duplicate processing
  3. Producer retries can send the same message multiple times
  4. Offset commits can fail, causing messages to be reprocessed

Achieving exactly-once semantics requires coordination across three components:

  1. Idempotent Producer

    • Prevents duplicate messages from retries
    • Uses producer ID and sequence numbers
  2. Transactional Producer

    • Atomic writes across partitions
    • Uses transactions
  3. Idempotent Consumer

    • Tracks processed offsets
    • Deduplicates messages
💡 Tip: Click dropdown to switch between languages
"exactly_once.py
from kafka import KafkaProducer, KafkaConsumer
from kafka.errors import KafkaError
import json
from typing import Set
class ExactlyOnceProcessor:
"""Exactly-once message processing"""
def __init__(self, bootstrap_servers: list, group_id: str):
# Idempotent producer
self.producer = KafkaProducer(
bootstrap_servers=bootstrap_servers,
enable_idempotence=True,
acks='all',
transactional_id='exactly-once-producer'
)
# Consumer with manual commit
self.consumer = KafkaConsumer(
bootstrap_servers=bootstrap_servers,
group_id=group_id,
enable_auto_commit=False,
isolation_level='read_committed' # Only read committed messages
)
# Track processed offsets (for idempotency)
self.processed_offsets: Set[tuple] = set()
def process_exactly_once(self, topic: str, handler):
"""Process messages exactly once"""
self.consumer.subscribe([topic])
try:
while True:
message_batch = self.consumer.poll(timeout_ms=1000)
for topic_partition, messages in message_batch.items():
# Begin transaction
self.producer.begin_transaction()
try:
for message in messages:
# Check if already processed (idempotency)
offset_key = (topic_partition.topic,
topic_partition.partition,
message.offset)
if offset_key in self.processed_offsets:
print(f"Skipping duplicate: {offset_key}")
continue
# Process message
result = handler(message.value)
# Send result to output topic (in transaction)
if result:
self.producer.send('output-topic',
value=json.dumps(result))
# Mark as processed
self.processed_offsets.add(offset_key)
# Commit transaction (atomic)
self.producer.commit_transaction()
# Commit consumer offset
self.consumer.commit()
except Exception as e:
# Abort transaction on error
self.producer.abort_transaction()
print(f"Error processing: {e}")
finally:
self.producer.close()
self.consumer.close()

Partitions Enable Parallelism

Partitions allow multiple consumers to process topic in parallel. Ordering guaranteed per partition.

Consumer Groups Scale

Consumer groups distribute partitions across consumers. Add consumers to scale throughput.

Offsets Track Progress

Offsets track consumer position. Commit offsets to resume after restart. Critical for reliability.

Exactly-Once is Complex

Exactly-once requires idempotent producer, transactions, and idempotent consumer. Use when duplicates are critical.