System Design Problem

Design a Distributed Message Broker (Kafka-style)

Commonly Asked By:LinkedInConfluentUberNetflix

Interview Setup

Interview Prompt

Design a distributed message broker (like Apache Kafka) that stores messages in an append only log, supports multiple consumer groups reading independently, allows replay from any offset, and provides both time-based retention and log compaction for keyed topics.

Clarifying Questions (ask before designing)

QuestionWhy it matters
Pull or push consumption model?Pull architecture in Kafka gives consumer-controlled pace, whereas push delivery in RabbitMQ enforces broker-controlled delivery rates.
Retention policy: time-based, size-based, or compaction?An event log with 7-day retention requires different storage engine semantics than a compacted changelog that only preserves the latest value per key.
Ordering guarantees needed?Per-partition ordering is native to the log design, whereas global ordering requires a single partition bottleneck.
Expected throughput and retention period?1M msg/sec x 7 days drives disk sizing and partition count.

Scope

In scope

  • Append-only commit log architecture
  • Partitions and offsets
  • Consumer groups with independent offset tracking
  • Pull-based consumption
  • Message replay from arbitrary offset
  • Log compaction for keyed topics

Out of scope (state explicitly)

  • Deep Kafka Streams / ksqlDB application design
  • Deep ZooKeeper/KRaft controller implementation internals
  • Exactly-once distributed transactions across systems

Functional Requirements

Start by framing this as a durable append only log, not a task queue such as the Distributed Worker Queue. Confirm publish and subscribe capabilities with partitions, consumer groups, offset tracking, replay, delivery semantics, and the retention window with your interviewer.

In the room: delivery semantics (at-most, at-least, or exactly once) and when you commit offsets are the first deep probes.

  • Publish: Producers publish messages to named topics.
  • Subscribe: Consumers subscribe to topics and receive messages in strict partition-level order.
  • Persistence: Messages durably stored on disk for a configurable retention period (time or size-based).
  • Consumer Groups: Multiple consumers in a group share processing load, with each partition assigned to exactly one group member.
  • Ordering: Messages within a partition are strictly ordered (FIFO) via monotonically increasing offsets.
  • Delivery Semantics: Support at most once and at least once delivery, plus exactly once processing for supported Kafka transactional pipelines.
  • Replay: Consumers can re-read historic messages by seeking to any offset or timestamp.
  • Partitioning: Topics split into partitions for parallelism, load distribution, and horizontal scaling.
  • Log Compaction: Retain only the latest value per key, essential for changelogs and CDC use cases.
  • Schema Evolution: Support schema validation and evolutionary compatibility checks via an external Schema Registry.

Non-Functional Requirements

Your interviewer will push on throughput, durability with acks=all and ISR replication, and failover behavior during broker failure. Sampling and backpressure matter less here than on distributed tracing, so focus on the trade-offs between zero message loss and availability.

  • High Throughput: 1M messages/sec system-wide peak across the broker cluster, with multiple GB/sec of aggregate throughput.
  • Low Latency: Target producer acknowledgement latency below 10 ms p99 and low broker fetch latency. End-to-end latency also depends on consumer processing and the configured fetch wait.
  • Durability: Protect acknowledged messages against the stated single broker failure model using acks=all, min.insync.replicas=2, and replication.factor=3. Catastrophic loss of all replicas is outside this assumption.
  • Horizontal Scalability: Add brokers and rebalance partitions to increase storage and throughput when partition distribution is sufficiently even. Capacity grows with additional balanced partitions, but not every workload scales linearly because hot keys, uneven leadership, network limits, and reassignment overhead can become limiting factors.
  • Fault Tolerance: Survive broker, rack, and even Availability Zone (AZ) failures with rack-aware replication.
  • Backpressure: Bounded pull model naturally protects slow consumers from getting overwhelmed.
  • Multi-Tenancy: Quotas (bandwidth, rate limits) per client-id to prevent noisy-neighbor congestion.

Capacity Estimations

Partition count and retention drive storage, so calculate these numbers before committing to broker sizing. Peak write bandwidth and days of retention determine how many brokers and how much disk are required.

MetricCalculationValue
Messages / sec (system-wide peak)Given1 Million
Avg message sizeGiven2 KB
Throughput (write)1M x 2 KB2 GB/sec
Throughput (read)3 consumer groups x 2 GB/s6 GB/sec
Retention PeriodGiven7 Days
Storage per day (raw)2 GB/s x 86400~172.8 TB
Storage for 7 days172.8 TB x 7~1.21 PB
Replication factor 31.21 PB x 3~3.63 PB total storage
Brokers (12 TB usable each, storage minimum)3.63 PB / 12 TB~303 brokers
Topics countGiven10,000
Partitions (total)Given500,000
Replica assignments (RF=3)500,000 x 31.5 Million
Planning brokers for ~4K replica assignments each1.5M / 4,000~375 brokers
Average network per broker at storage minimum(2 GB/s + 4 GB/s + 6 GB/s) ÷ 303~39.6 MB/s
Average replica assignments / storage-minimum broker1.5M ÷ 303~4,950
Average replica assignments / planning broker1.5M ÷ 375~4,000
10x stress throughput10M x 2 KB20 GB/sec
Producer network: 2 GB/s ingress across cluster
Replication network: 2 GB/s x 2 (RF=3 means 2 follower copies) = 4 GB/s intra-cluster
Consumer network: 6 GB/s egress for 3 full-stream consumer groups
Total cluster I/O: ~12 GB/s (~96 Gb/s aggregate). Size broker links and network fabric for peak traffic, using 25GbE or 100GbE NICs as standard starting points

Per broker at storage-minimum sizing (~303 brokers):
  Write: ~6.6 MB/s per broker
  Replicate: ~13.2 MB/s
  Serve reads: ~19.8 MB/s
  Disk: sequential writes at ~200 MB/s per SSD → comfortable headroom
Capacity uses uncompressed payload size. Actual network and disk usage can be lower after compression, while production disk planning should also reserve space for indexes, metadata, recovery headroom, and operational slack.

Planning note:
  ~375 brokers brings average replica assignments to ~4,000 per broker
  and leaves storage headroom above the theoretical ~303-broker minimum

10x stress case:
  20 GB/s write ingress
  40 GB/s replication traffic
  60 GB/s consumer egress
  ~120 GB/s total cluster I/O
  If sustained for the full 7-day retention window:
    ~12.1 PB raw storage and ~36.3 PB at RF=3
    The baseline ~303-broker cluster would therefore need roughly 10x storage capacity
  Sustaining the full stress case at the same per-broker load would require roughly 3030 brokers

Architecture Diagram

Walk the diagram as a log pipeline where producers append to partition leaders, followers replicate via ISR, and consumers pull at their own offset. We omit an edge gateway layer because producers talk directly to brokers once they receive metadata from KRaft.

Unlike a Distributed Worker Queue, messages are retained durably for days and can be replayed by resetting consumer offsets. Partition count sets the maximum parallelism per consumer group. The combination of acks=all and min.insync.replicas trades write availability for durability when brokers fail.

In the room: when asked about exactly once processing, describe the combination of an idempotent producer and transactional consume-process-produce loops, and distinguish Kafka to Kafka processing from external side effects.

Loading...

Component Deep Dives

1. Topics and Partitions ⭐

Next we examine the architecture partition by partition. Start with how topics split for parallelism, then cover the append only storage engine, which forms Kafka's core interview story. Partitions serve as the fundamental unit of ordering, parallelism, and horizontal scale-out. Explain partition keys before diving into low-level storage mechanics. Topics are divided into partitions so that different partitions are spread across different brokers.

Topic: user-events (3 partitions)

Partition 0: [msg0, msg1, msg2, msg3, ...]  → Offset 0, 1, 2, 3
Partition 1: [msg0, msg1, msg2, ...]         → Offset 0, 1, 2
Partition 2: [msg0, msg1, ...]               → Offset 0, 1

Key insight: ordering is ONLY within a partition, NOT across partitions.
- All events for user_id=123 go to same partition → ordered.
- Events across user_id=123 and user_id=456 have NO ordering guarantee.

2. Broker Storage Engine: Why Append-Only is Fast ⭐

The storage engine explains how Kafka reaches high throughput through sequential writes, batching, page cache reads, and zero copy transfer. Kafka stores each partition as an append only sequence of segment files. The following JVM and disk values are illustrative sizing assumptions, not universal broker defaults.

  • Sequential I/O: Sequential writes (600 MB/s on SSD) are 10-100x faster than random I/O.
  • OS Page Cache: The JVM heap is kept small (~6 GB). The rest of RAM acts as the OS Page Cache. Read paths hit memory instead of disk.
  • Zero Copy: Transfers file data through the kernel to the network path with fewer user space copies, reducing CPU overhead for reads.
  • Batched I/O: Groups multiple messages to amortize per request and storage scheduling overhead. Kafka does not require an fsync for every batch.
/data/user-events-0/
  00000000000000000000.log      (segment 1: offsets 0 – 999,999)
  00000000000000000000.index    (sparse index: maps offset → byte position)
  00000000000000000000.timeindex
  00000000000001000000.log      (segment 2: offsets 1,000,000 – 1,999,999)
  00000000000001000000.index
  leader-epoch-checkpoint        (tracks leader changes for truncation)

Message Lookup by Offset Walkthrough

Consumer requests offset 1,500,042:
1. Binary search segment files by base offset → find segment starting at 1,000,000
2. Open 00000000000001000000.index → binary search for ≤ 1,500,042
   → Index entry: offset 1,500,000 → file position 483,200
3. Seek to position 483,200 in .log file
4. Scan forward 42 messages → found offset 1,500,042
Total: index lookup + positioned log read + short scan → typically < 1ms from page cache

3. ZooKeeper vs KRaft (Metadata Management) ⭐

Kafka 4.0 and later use KRaft rather than ZooKeeper. The older ZooKeeper model remains relevant when explaining migrations from pre 4.0 clusters.

ZooKeeper (legacy for pre Kafka 4.0 deployments):
- Stores: broker registry, topic configs, partition to leader mapping, ACLs
- Provides the older metadata coordination and controller model
- Adds a separate operational dependency for the Kafka cluster

KRaft (Kafka Raft, required in Kafka 4.0+):
- Kafka's own Raft based consensus layer for cluster metadata
- Metadata is maintained in the internal __cluster_metadata log
- Controller quorum: typically 3 or 5 controller nodes. Use dedicated controller roles in critical production deployments, reserving combined broker and controller mode for development or small noncritical environments
  Active controller = Raft leader → handles metadata changes
  Standby controllers = Raft followers → take over on failure

Benefits over ZooKeeper:
  1. No external dependency → simpler operations
  2. Metadata is replicated through Kafka's own quorum mechanism
  3. Designed for large partition counts without a separate ZooKeeper ensemble
  4. Metadata is log based → supports snapshot and replay

4. Replication Consensus (In-Sync Replicas) ⭐

Replication provides high availability and protects acknowledged records under the configured replication and failure assumptions.

  • Leader: Serves all producer writes and consumer reads.
  • Followers: Continuously pull (fetch) from the leader to replicate logs.
  • ISR (In-Sync Replicas): Followers that remain sufficiently caught up within replica.lag.time.max.ms (default 30s).
  • High Watermark (HW): The replication watermark up to which data is committed for normal consumption. With read_committed, the Last Stable Offset can be lower when open transactions exist.
  • Eligible Leader Replicas (Kafka 4.0+): When ELR is enabled, Kafka can track replicas outside ISR that are still safe to become leader without data loss. When applicable, these replicas can participate in safe leader failover alongside ISR members.

ISR Commit Offset Example

Partition 0 (Topic: order-events):
  Leader:    Broker 1 (Rack A)  LEO: 5,000,151   ← Serves all clients
  Follower:  Broker 4 (Rack B)  LEO: 5,000,149   ← In ISR
  Follower:  Broker 7 (Rack C)  LEO: 5,000,151   ← In ISR

  ISR = {Broker 1, Broker 4, Broker 7}

  High Watermark (HW) = min(ISR LEOs) = 5,000,149
  → Record offsets below 5,000,149 are committed and visible through offset 5,000,148
  → Offset 5,000,149 and later are not yet committed by the replication watermark

  Note: With read_committed consumers, the Last Stable Offset (LSO) can be lower than HW when open transactions exist.

5. Rack-Aware Replication ⭐

Rack-aware replica allocation protects against datacenter and rack-level network or power failures.

broker.rack=us-east-1a  (Broker 1, 4)
broker.rack=us-east-1b  (Broker 2, 5)
broker.rack=us-east-1c  (Broker 3, 6)

Partition 0 replicas placed on: Broker 1 (1a), Broker 2 (1b), Broker 3 (1c)
→ Survives any single AZ failure

Without rack-awareness: replicas might all land on the same rack
→ Rack failure can make all replicas unavailable and may cause data loss if no other recovery copy exists

6. Producer Write Path ⭐

Loading...
  • acks=0: Fire-and-forget (fastest, high data-loss risk).
  • acks=1: Await leader confirmation (vulnerable if leader dies before replicating).
  • acks=all: Await acknowledgements from every replica currently in ISR. With min.insync.replicas=2 and RF=3, this protects acknowledged records against the stated single broker failure model.

Producer Partitioner Strategies

1. Key-based routing when a key is present: conceptually hash(key) → partition
   ✅ Same key → same partition while the partition mapping is stable → per-key ordering
   ❌ Hot key → hot partition (e.g. viral user's writes saturate one partition)
   ⚠️ Increasing partition count can remap future records for an existing key

2. Round-robin: distribute events evenly across partitions when entity ordering is unnecessary
   ✅ Even load distribution
   ❌ No ordering guarantees for any entity

3. Sticky partitioning for records without keys: batch to one partition until the batch is ready to send
   ✅ Maximizes batch size → better compression, fewer network requests
   ❌ No entity-based ordering guarantee when records have no key

4. Custom partitioner: application-specific routing logic
   - Example: route by tenant_id to dedicated isolated partitions to enforce SLA bounds.

7. Consumer Read Path ⭐

Loading...

Consumers poll brokers for data and maintain committed offsets in the internal __consumer_offsets topic. The committed offset records the recovery point, not proof that an external side effect was performed exactly once.

8. Consumer Group Rebalancing ⭐

  • Eager Rebalance (Legacy): Pauses the group, revokes all partitions, and assigns them from scratch.
  • Cooperative Incremental (Classic Protocol): Revokes and reassigns affected partitions in stages so unaffected consumers can continue processing.
  • Kafka 4.x Consumer Protocol: When group.protocol=consumer is enabled, assignment is coordinated by the broker and the rebalance protocol is incrementally designed without a global synchronization barrier. The consumer must opt in, while the server supports the protocol.
  • Static Membership (group.instance.id): Identifies a stable consumer instance so a brief restart within the configured session timeout can avoid unnecessary reassignment when using the classic protocol.

API Design

Sketch producer and consumer API contracts to show you understand offsets and commit timing, which makes delivery semantics concrete.

Producer API Interface

TYPESCRIPT
interface ProducerRecord<K = string, V = string> {
  topic: string;
  key: K;
  value: V;
  partition?: number;
  timestamp?: number;
  headers?: Record<string, string>;
}

interface RecordMetadata {
  topic: string;
  partition: number;
  offset: bigint;
  timestamp: number;
}

interface TopicPartition {
  topic: string;
  partition: number;
}

interface OffsetAndMetadata {
  offset: bigint;
  metadata?: string;
}

interface Producer {
  // Synchronous send that blocks until broker acknowledgment is confirmed
  sendSync(record: ProducerRecord): Promise<RecordMetadata>;

  // Asynchronous non-blocking send with result callback
  send(
    record: ProducerRecord,
    callback: (meta: RecordMetadata | null, err: Error | null) => void
  ): void;

  // Transactional publish supporting atomic cross-topic writes and offset commits
  beginTransaction(): Promise<void>;
  sendOffsetsToTransaction(
    offsets: Map<TopicPartition, OffsetAndMetadata>,
    consumerGroupId: string
  ): Promise<void>;
  commitTransaction(): Promise<void>;
  abortTransaction(): Promise<void>;
}

Consumer API Interface

TYPESCRIPT
interface TopicPartition {
  topic: string;
  partition: number;
}

interface ConsumerRecord<K = string, V = string> {
  topic: string;
  partition: number;
  offset: bigint;
  key: K;
  value: V;
  timestamp: number;
}

interface OffsetAndTimestamp {
  offset: bigint;
  timestamp: number;
}

interface OffsetAndMetadata {
  offset: bigint;
  metadata?: string;
}

interface Consumer {
  // Subscribe to topics dynamically with automatic group partition assignment
  subscribe(topics: string[]): void;

  // Manually assign specific static partitions to bypass rebalancing
  assign(partitions: TopicPartition[]): void;

  // Fetch batches of records with a bounded wait timeout
  poll(timeoutMs: number): Promise<ConsumerRecord[]>;

  // Commit offsets synchronously after message processing succeeds
  commitSync(offsets?: Map<TopicPartition, OffsetAndMetadata>): Promise<void>;

  // Seek to an absolute offset for historical replay or error recovery
  seek(partition: TopicPartition, offset: bigint): void;

  // Look up partition offsets corresponding to a specific target timestamp
  offsetsForTimes(
    timestamps: Map<TopicPartition, number>
  ): Promise<Map<TopicPartition, OffsetAndTimestamp | null>>;
}

Admin API Interface

TYPESCRIPT
interface NewTopicConfig {
  name: string;
  numPartitions: number;
  replicationFactor: number;
  configs?: {
    "retention.ms"?: string; // e.g. "604800000" for 7 days
    "min.insync.replicas"?: string; // e.g. "2"
    "compression.type"?: string; // e.g. "lz4"
  };
}

interface PartitionReassignment {
  topic: string;
  partition: number;
  replicas: number[]; // Broker IDs for target replica placement
}

interface ClusterDescription {
  clusterId: string;
  controllerId: number;
  nodes: Array<{ id: number; host: string; port: number; rack?: string }>;
}

interface AdminClient {
  // Create topics with explicit rack-aware replica configurations
  createTopics(topics: NewTopicConfig[]): Promise<void>;

  // Reassign partitions during broker decommissioning or cluster rebalancing
  alterPartitionReassignments(reassignments: PartitionReassignment[]): Promise<void>;

  // Describe broker nodes, active KRaft controller, and cluster state
  describeCluster(): Promise<ClusterDescription>;
}

Data Model

Record Batch Format (On-Disk Segment Layout)

Producers serialize and pack records into a compact batch layout to maximize disk and network compression.

Loading...

Sparse Segment Index Layout

Offset      Position (file byte offset)
0           0
4096        32768        ← Entry every ~4 KB of log data
8192        65536
12288       98304

Lookup offset 10000:
1. Binary search index → nearest ≤ entry = 8192 at position 65536
2. Seek to 65536 in log, scan forward to offset 10000.

Consumer Offset Storage

Topic: __consumer_offsets (50 partitions, compacted)
Partition key: hash(group_id) % 50

Key:   {group_id: "order-processor", topic: "orders", partition: 7}
Value: {offset: 5000150, metadata: "", commit_timestamp: 1710403200000}

Compacted: compaction eventually retains the latest offset record for each (group, topic, partition) key in the log.

Fault Tolerance

Failure CaseSystem Solution Design
Broker Leader CrashKRaft selects a safe replacement leader from eligible replicas, normally the ISR and, when ELR is enabled in Kafka 4.0+ where applicable, Eligible Leader Replicas. Clients refresh metadata automatically.
Message LossEnforce acks=all + min.insync.replicas=2 + replication.factor=3 for durability.
Consumer FailureGroup Coordinator triggers partition rebalancing to re-assign tasks to active consumers.
AZ or Rack FailureRack-aware partition replica allocation (broker.rack) distributes data copies across physical racks/AZs.
Split-Brain ControllersKRaft consensus maintains a single active metadata leader, and epochs fence stale controller or broker generations.
Duplicate RetriesIdempotent producer validates Producer ID (PID) + Sequence Number to filter duplicates.
Disk FailureWith RF=3 and rack aware placement, data has two other replicas on distinct racks or AZs. Operators replace the failed disk, and the broker rebuilds missing segments from surviving replicas.
Network PartitionReplicas that fall behind leave the ISR, while leader epochs fence stale leaders. With min.insync.replicas enforced, the cluster stops acknowledging new writes when the ISR falls below the safety threshold, while committed records remain readable when a healthy leader is available.
Unclean Leader ElectionSet unclean.leader.election.enable=false to prevent out-of-sync replicas from becoming leader.
Slow ConsumerEvent lag grows without directly blocking producers, but sustained lag can increase retention pressure and disk usage if consumers remain behind.

1. Exactly-Once Semantics (EOS) Framework ⭐

Kafka separates producer durability from consumer processing semantics. At most once and at least once describe common consume and commit patterns, while exactly once applies to supported Kafka processing topologies that use transactions and read committed consumption.

  • At Most Once: Commit the consumer offset BEFORE processing. A crash after the commit can lose the record because it will not be replayed. Producer publish durability is a separate concern.
  • At Least Once: Process the record BEFORE committing its offset. A crash after processing but before the commit can cause the record to be processed again, so downstream effects must be idempotent. Producer acknowledgement settings are a separate publishing durability concern, and producers should normally use acks=all for important data.
  • Exactly-Once: For Kafka to Kafka processing, combine an idempotent producer with a transaction that atomically publishes output records and commits consumed offsets. Consumers that should ignore aborted transactional records use isolation.level=read_committed. External side effects still require an idempotent or transactionally integrated destination.

2. Race Condition: Stale Leader split-brain (Fencing via Epochs) ⭐

If an isolated leader continues accepting writes after being replaced, divergence occurs. Kafka prevents this using Leader Epochs.

T=0:   B1 is leader, leader_epoch=5.
T=10s: Network isolates B1. Controller elects B2 as leader, leader_epoch=6.
T=11s: B1's network recovers.
       - Client metadata may still refer to the old leader epoch.
       - Requests sent using stale leader metadata are rejected or fenced as applicable, and clients refresh metadata.
       - B1 learns epoch=6, truncates any uncommitted tail that is no longer valid, and remains a follower.

3. Race Condition: Consumer Offset Commit Gap ⭐

If a consumer crashes after processing but before committing its offset, it can process the same record again after restart. Offset commits alone therefore do not provide exactly once side effects. An idempotent sink or Kafka transaction is required when the application needs stronger guarantees.

Solutions:

  • Idempotent Consumer: Make downstream actions idempotent, such as using an UPSERT with a unique key constraint (e.g. INSERT INTO results (id, data) VALUES ('orders-7-50', ...) ON CONFLICT DO NOTHING).
  • Kafka Transactions (consume-transform-produce): Commit offsets and write output atomically in a single transaction.

4. ISR Shrinkage & The Data Loss Window ⭐

ISR (In-Sync Replicas) are followers caught up with the leader within replica.lag.time.max.ms (default 30s).

Config: acks=all, min.insync.replicas=2, replication.factor=3

T=0:   All healthy. ISR = {B1, B2, B3}. acks=all means all 3 ACK.
T=10s: Broker3 hits GC pause, falls behind > 30 seconds.
T=10s: Controller removes B3 from ISR. ISR = {B1, B2}.
T=10s: acks=all now means acks={B1, B2} only (reduced redundancy, still safe).
T=15s: Broker1 (leader) crashes. B2 has all committed data → becomes leader.
T=15s: B3 recovers → truncates uncommitted data to High Watermark and re-fetches.

BUT if min.insync.replicas=1 (DANGEROUS CONFIG):
- ISR shrinks to {B1} (just the leader).
- acks=all effectively behaves as acks=1.
- If B1 crashes, un-replicated messages are lost forever.

Durability Recommendation:
Set: acks=all + min.insync.replicas=2 + replication.factor=3

5. Follower Fetch Protocol ⭐

Followers pull data from the partition leader using a continuous loop:

TYPESCRIPT
while (true) {
  FetchRequest to leader: {partition: 0, fetch_offset: 5000148, max_bytes: 1MB}
  Leader responds: {messages from offset 5000148..5000200, leader_epoch: 5, HW: 5000145}
  Follower appends messages to local log
  Follower advances its local High Watermark toward the leader watermark after replication
  Repeat continuously to minimize replication lag
}

Why Pull over Push?
1. Follower controls pace → natural backpressure
2. Simpler: leader does not track per-follower push state
3. Follower crash → it stops fetching while the leader continues serving
4. Follower recovery → it resumes from its last offset and catches up

Specs: same data-center replication is commonly low latency, but actual lag depends on hardware, load, network, and broker configuration.

6. Controlled Shutdown vs Uncontrolled Shutdown

  • Controlled: Admin requests shutdown. Leadership is moved off the broker before it stops, which can avoid client visible interruption when healthy replicas are available.
  • Uncontrolled: Broker crashes. The cluster detects the failure and elects a new leader, which can cause temporary partition unavailability during detection and election.

7. Broker Decommission and Partition Reassignment

To safely remove a broker (e.g., hardware refresh):

  1. Generate a partition replica reassignment plan via the Admin tool.
  2. Execute reassignment: New replicas fetch full logs from partition leaders in the background.
  3. Apply bandwidth throttling (limit.bytes.per.second) to prevent cluster network starvation during transfer.
  4. Once caught up, the new node enters the ISR set, and the decommissioned broker's replicas are removed.
  5. Safely shut down the old broker (now containing 0 active replicas).

Additional Considerations

1. Schema Registry & Compatibility Modes ⭐

Decoupling producer and consumer deployments requires strict schema validation to prevent structural breaking changes.

Compatibility ModeEvolution RuleDeployment Order
BACKWARD (Default)New schema can read data written by old schemas. Add optional fields with defaults or remove fields that the new reader no longer needs.Upgrade Consumers first, then upgrade Producers.
FORWARDOld schema can read data written by new schemas. Add fields only when older readers can safely ignore them, and remove fields only when compatibility permits it.Upgrade Producers first, then upgrade Consumers.
FULLSchemas remain compatible in both directions. Changes should use optional fields with defaults or other evolution patterns supported by the registry.Deploy independently after compatibility validation.

Serialization Format Comparison

Avro vs Protobuf vs JSON Schema:
- Avro: Compact binary, schema evolution support, and widely used in Kafka ecosystems.
- Protobuf: Compact binary, strong typing, and common in gRPC microservices ecosystems.
- JSON Schema: Human-readable, larger payload footprint, and easy to inspect during debugging.

Recommendation:
Use Avro for general Kafka environments due to Confluent's extensive native tooling, or Protobuf if already using a gRPC microservices ecosystem.

Wire Message Binary Layout

Loading...

2. Tiered Storage (Solving Cost at Petabyte Scale)

Holding petabytes of historical logs on fast local NVMe SSDs is cost-prohibitive. Tiered storage separates computing and storage.

  • Hot Tier (Local Disk): Holds last 4-24 hours for real-time consumers (reads hit page cache / SSDs).
  • Cold Tier (S3 / Object Store): Moves closed segment data to a configured remote storage implementation such as S3 or GCS for long-term retention. Kafka provides the tiered storage framework, while the remote storage implementation must be configured separately. This can reduce local storage costs substantially. Current Kafka tiered storage does not support compacted topics, so use it with eligible retained event topics rather than compacted changelog topics.

3. Kafka Connect: CDC and Database Integrations

Instead of writing custom services to move data in/out of Kafka, Kafka Connect provides a robust, distributed connector runtime.

Loading...
 
Common Kafka Connect connector plugins:
- Source connectors:
  * Debezium for MySQL/PostgreSQL CDC log tailing to Kafka ⭐
  * JDBC Source for polling relational databases
  * Other vendor and community plugins for S3, Kinesis, ActiveMQ, and RabbitMQ
- Sink connectors:
  * Elasticsearch/OpenSearch for search indexing
  * S3 / GCS / Azure Blob for long-term data lake archival
  * Snowflake, BigQuery, and Redshift for data warehouse sinks
  * JDBC for relational database sinks

Why use Kafka Connect instead of a custom connector runtime?
- Connect handles offset management, task parallelism, worker load balancing, fault tolerance, monitoring,
  and schema registry integration.
- Kafka Connect can provide exactly once semantics for connector and framework paths that explicitly support them and are configured for them.
- Debezium avoids the operational complexity of building and maintaining a full CDC pipeline from scratch.

4. Kafka Streams vs. Apache Flink

For real-time transformations, aggregations, and streaming joins, choose the right processor engine.

TYPESCRIPT
interface KStream<K, V> {
  filter(predicate: (key: K, value: V) => boolean): KStream<K, V>;
  groupByKey(): KGroupedStream<K, V>;
  to(topic: string): void;
}

interface KGroupedStream<K, V> {
  windowedBy(windowSizeMs: number): TimeWindowedKStream<K, V>;
}

interface Windowed<K> {
  key: K;
  startMs: number;
  endMs: number;
}

interface TimeWindowedKStream<K, V> {
  count(): KTable<Windowed<K>, number>;
}

interface KTable<K, V> {
  toStream(): KStream<K, V>;
}

interface StreamBuilder {
  stream<K, V>(topic: string): KStream<K, V>;
}

// Processing pipeline: 5-minute windowed aggregation of high-value orders
const builder: StreamBuilder = getStreamBuilder();
builder
  .stream<string, { amount: number }>("orders")
  .filter((key, order) => order.amount > 100)
  .groupByKey()
  .windowedBy(5 * 60 * 1000)
  .count()
  .toStream()
  .to("high-value-order-counts");

// Key Benefits:
// - Embeds as a lightweight application library without requiring a separate cluster
// - Supports exactly once processing when configured with an EOS processing guarantee such as exactly_once_v2
// - Uses local RocksDB state stores backed by internal changelog topics for fault recovery
  • Kafka Streams: Run as a lightweight Java library in your app. Excellent for pure Kafka-to-Kafka streaming and simple microservice transforms.
  • Apache Flink: Separate dedicated compute cluster. Well suited for multi-source joins, out-of-order events with watermarks, and complex event processing at scale.

5. Security and Multi-Tenancy

A production broker must isolate tenants and protect client connections and topic data.

  • Encryption in Transit: Require TLS for producer, consumer, and inter broker connections.
  • Authentication: Use SASL mechanisms such as SCRAM or OAuth based on the identity platform.
  • Authorization: Enforce topic and consumer group ACLs so clients can access only permitted resources.
  • Quotas: Apply client and user quotas for request rate and bandwidth to prevent noisy neighbors from exhausting broker capacity.
  • Encryption at Rest: Use encrypted broker volumes and protect credentials through the platform secret store.

6. Monitoring & Alerting

Keep a close eye on the health of your message broker using these alerts:

  • Critical Alerts:
    • Under-replicated partitions > 0 (Data redundancy at risk)
    • Offline partitions > 0 (Topic partition unavailable)
    • ISR shrink rate > threshold (Replicas falling out of sync)
    • Consumer lag > threshold (Consumers falling behind)
  • Warning Alerts:
    • Disk usage > 80% on any broker node
    • p99 request latency > 100ms
    • Request handler idle ratio < 30% (Broker CPU overloaded)
    • Unclean leader election count > 0

7. Performance Tuning Cheat Sheet (Staff Level)

Producer tuning:
  batch.size=65536 (64 KB)     → larger batches = higher throughput
  linger.ms=10                  → wait up to 10ms to fill a batch
  compression.type=lz4          → strong throughput/compression trade-off
  buffer.memory=67108864 (64MB) → buffer for async sends
  acks=all                      → durability for important data

Consumer tuning for the classic protocol:
  fetch.min.bytes=1048576 (1 MB)     → batch more data per fetch when available
  fetch.max.wait.ms=50               → bound intentional long-poll wait for latency sensitive workloads
  max.poll.records=500               → messages per poll() call
  max.poll.interval.ms=300000 (5min) → max processing interval between poll() calls before membership can be revoked
  Kafka 4.x: group.protocol=consumer  → opt in to the broker-coordinated Consumer rebalance protocol, where classic heartbeat and session settings become broker-managed

Broker tuning:
  num.io.threads=16                  → disk I/O threads
  num.network.threads=8              → network threads
  socket.send.buffer.bytes=1048576   → 1 MB socket buffer
  log.segment.bytes=1073741824 (1GB) → segment size
  log.retention.hours=168 (7 days)   → retention period
  num.partitions=12                  → illustrative topic setting, not Kafka's broker default

Planning heuristics:
- Keep replica assignments per broker within a tested operating range, with ~4000 as a planning heuristic unless the cluster is explicitly sized for higher metadata overhead
- Partition count = max(
    producer_throughput / single_partition_throughput,
    maximum_consumer_count_per_group
  )
- Adding partitions can change hash(key) → partition mapping, so continuous per-key ordering across the resize is not automatic

Interview Time Budgets (25 / 50 / 75 min)

25 min: log, partitions, consumer groups, and at least once vs exactly once semantics.50 min: add ISR, min.insync.replicas, rebalancing, compaction, and tiered storage.75 min: EOS transactions, controller failover, rebalance storms, and multi region mirroring.

Interview Walkthrough

  • 25-minute cut

    Skip KRaft/ZooKeeper migration unless staff.

    • Log + partitions + key routing (5 min)
    • consumer groups + rebalance (5 min)
    • at least once vs exactly once (7 min)
    • ISR + min.insync.replicas (5 min)
    • retention/replay use case (3 min)
  • Start with the log abstraction by presenting an append only, partitioned commit log where producers append records and consumers track independent offsets.
  • Explain partitions by highlighting that key based routing guarantees ordering per key, and that adding partitions increases parallelism at the expense of higher consumer group coordination.
  • Cover consumer groups by explaining that each partition maps to exactly one consumer within a group. Explain classic eager and cooperative rebalancing, then mention the Kafka 4.x consumer protocol when the interview reaches current large cluster behavior.
  • Compare delivery semantics by detailing at most once commits before processing, at least once commits after processing with possible duplicates, and exactly once Kafka processing using idempotent producers plus transactions and read committed consumers. Clarify that external side effects still require destination level idempotency or transactional integration.
  • Discuss retention and replay capabilities where brokers store messages for days, allowing consumers to rewind committed offsets to reprocess historical event streams.
  • Address partition replication with leader and follower roles, the In-Sync Replicas set, and min.insync.replicas thresholds balancing durability against availability.
  • Highlight the common pitfall of assuming global ordering across a topic, emphasizing that ordering is strictly guaranteed per partition. Global ordering requires a single partition or an application level sequencing and resequencing mechanism.

Architectural Context and Related Systems

A distributed message broker functions as the central asynchronous backbone within event driven architectures. Understanding where it fits alongside adjacent distributed systems strengthens design interview discussions:

Engineering Trade-offs

Close with broker comparisons and partition-count trade-offs. Your interviewer may ask why not RabbitMQ for this workload.

1. Broker Technology: Kafka vs RabbitMQ vs Pulsar

Choose the broker model that best matches the processing pipeline and its durability, replay, and latency requirements:

DimensionKafkaRabbitMQApache Pulsar
Storage ModelImmutable append only logDurable transient queue (delete on ACK)Segment storage with object-store offload
Consumer ModelPull (consumer tracks offsets)Push (broker tracks delivery state)Hybrid push/pull (cursor-per-consumer)
Message Replay✅ Yes (offset seek)❌ No (deleted on ACK)✅ Yes (cursor reset)
Multi-tenancyWeak (quotas per client)Fair (virtual hosts)Strong (native namespaces + quotas)
Geo-replicationMirrorMaker2 (external tool)Shovel/Federation pluginsBuilt-in (native out of the box)
LatencyLow, workload and batching dependentVery low for queue workloadsLow, workload and batching dependent

2. Consumer Group Rebalancing: The Throughput Killer

Rebalancing can pause consumer processing and create latency spikes, especially under classic eager assignment. Kafka 4.x also provides the newer consumer rebalance protocol for incrementally coordinated groups.

CLASSIC PROTOCOL: EAGER REBALANCING (legacy stop-the-world pattern):
  T=0:    Consumer B crashes
  T=10s:  Illustrative session timeout expires and coordinator marks B dead
  T=10s:  ALL consumers in the group revoke their assignments
  T=10.5s: Consumers rejoin and coordinator assigns partitions
  T=11s:  Consumers resume
  ⚠️ TOTAL PAUSE: roughly 11 seconds in this illustrative configuration

COOPERATIVE INCREMENTAL REBALANCING (modern and preferred approach):
  T=0:    Consumer B crashes
  T=10s:  Illustrative session timeout expires
  T=10s:  Only B's affected partitions are revoked
  T=10s:  Consumers A and C continue processing their unaffected partitions
  T=10.5s: B's partitions are reassigned to active consumers
  ⚠️ Only affected partitions pause during reassignment
  Config: partition.assignment.strategy=CooperativeStickyAssignor

STATIC GROUP MEMBERSHIP under the classic protocol (group.instance.id) ⭐:
  Consumer B restarts within the configured session timeout
  → Coordinator can preserve its assignment instead of triggering a full rebalance
  → Useful for Kubernetes rolling deployments

3. Choosing the Right Partition Count ⭐

Partition count is an expansion sensitive decision. You can increase partitions, but Kafka cannot decrease them in place, and hash based partitioning can remap future records for existing keys after an expansion.

TYPESCRIPT
partitions = max(
  target_throughput / throughput_per_partition,
  maximum_consumer_count_per_group
)

Example:
  Target: 100 MB/s for topic. Single partition: ~10 MB/s throughput.
  → 10 partitions minimum for throughput.
  But we want 20 consumers for processing → need 20 partitions.
  → Choose 20 partitions.

Important:
  Partitions can be increased, but Kafka cannot decrease them in place.
  If routing uses hash(key) % partition_count, increasing the partition count can move future records for a key to a different partition, so continuous per-key ordering across the resize is not automatic.
  • Too Few: Limits consumer scaling and throughput while creating single-broker hotspots.
  • Too Many: Exhausts file descriptors (each segment has .log, .index, .timeindex files), delays controller elections on crash, and slows down consumer rebalances.

4. Message Ordering Guarantees: Per-Partition vs Global

Kafka guarantees ordering strictly within a partition. The broker's log order is authoritative for record order within that partition and does not imply wall clock ordering across producers.

  • Global Ordering: Configure a topic with exactly 1 partition. This provides total order but limits each consumer group to one active consumer for that topic and to the throughput of that single partition.
  • Per-Key Ordering (Default): Partition by entity ID (e.g., user_id). All events for a specific user land in the same partition and are processed in order.
  • Common Pitfalls:
    • 1. Producer Retries without Idempotency:

      If a producer sends a message but the network acknowledgment is lost, a retry can create a duplicate record and, with multiple in flight requests, can also change append order. For example, if message 1 times out and message 2 is committed before the retry for message 1, message 1 may be appended after message 2. To resolve this, use enable.idempotence=true so producer IDs and sequence numbers let the broker detect retried duplicates.

    • 2. Multiple Producers writing to the same partition:

      Independent producers can arrive at a partition in an order that differs from their event timestamps. Kafka preserves append order in the partition, not wall clock order across producers. If business logic requires event time ordering, include application sequence information or use a bounded consumer-side reordering strategy.

    • 3. Concurrent Consumer processing (multi-threading):

      When a consumer delegates records from one partition to an asynchronous thread pool, later records can finish before earlier records. Idempotency prevents duplicate side effects but does not restore ordering. If processing order matters, enforce serial execution per partition or per key.

💬Review

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...