PrepZone Logo
PrepZone

Distributed Message Queue

Build a queue service with partitions, consumer groups and at-least-once delivery.

Read these first

Design a queue service for StreamHub's internal teams: publish messages to named topics, consume in consumer groups with offset tracking, retain messages for replay, and scale to 1M messages/sec.

Requirements

Functional requirements

  • produce(topic, message) — append message to topic partition.
  • consume(topic, group, offset) — read messages from assigned partitions.
  • Consumer groups: competing consumers share partition assignment.
  • Message retention by time (7 days) and size (1 TB per topic).
  • Dead letter topic for messages exceeding max retry count.

Non-functional requirements

  • Throughput: 1M messages/sec aggregate across topics.
  • Latency: p99 produce under 10 ms; p99 consume under 5 ms (after publish).
  • Durability: Replicate each message to 3 brokers before ack.
  • Availability: Tolerate broker failure without message loss.
  • Ordering: Per-partition ordering guaranteed; cross-partition unordered.

Estimation

  • 1M msg/sec × 1 KB avg = 1 GB/sec ingest → 86 TB/day raw → replication factor 3 = 258 TB/day.
  • 100 topics × 24 partitions × 3 replicas = 7200 partition-replicas — ~50 broker cluster.

Amazon MSK event pipeline

produceCOMPUTE
EKS APIorder placed
INTEGRATION
Amazon MSKorders.placed.v1
COMPUTE
InventoryEKS worker
COMPUTE
Email svcEKS worker
ANALYTICS
AnalyticsFlink / EMR
API publishes events; worker fleets scale independently on consumer lag.

API design

Java
# Produce
POST /v1/topics/{topic}/messages:
  headers: { X-Partition-Key: "streamer_4420" }
  body: { "payload": { "event": "viewer.joined", "count": 1523 } }
  response: { "offset": 918273, "partition": 7 }

# Consume (long-polling)
GET /v1/topics/{topic}/messages:
  query: { group: "analytics-workers", max_messages: 100, timeout_ms: 5000 }
  response:
    messages: [{ "offset": 918274, "partition": 7, "payload": {...} }]
    next_offset: 918275

# Commit offset
POST /v1/consumer-groups/{group}/offsets:
  body: { "topic": "viewer.events", "partition": 7, "offset": 918275 }

Topic and partition design

Partitioning rules

  • Partition key (e.g., streamer_id) hashes to partition — all events for one streamer ordered.
  • More partitions = more parallelism; but also more overhead (文件句柄, replication).
  • Rule of thumb: partitions ≥ max desired consumer count in a group.
  • Rebalance: when consumers join/leave, partitions reassigned (sticky assignment minimises movement).
Java
# Topic configuration (pseudo)
topic: viewer.events
partitions: 48
replication_factor: 3
retention_ms: 604800000        # 7 days
min_insync_replicas: 2         # ack after 2 of 3 replicas
compression: lz4

Storage architecture

Broker internals

  • Log segments: Partition stored as sequence of segment files (1 GB each); active segment accepts writes.
  • Index files: Sparse offset and timestamp indexes for O(1) lookup within segments.
  • Replication: Leader handles writes; followers fetch from leader; ISR (in-sync replicas) ack produce.
  • Zero-copy transfer: sendfile() from disk to network socket — no user-space copy.

Consumer groups

AspectConsumer group (competing)Independent consumer (broadcast)
Partition assignmentEach partition → one consumer in groupAll partitions → every consumer
Use caseScale processing — N consumers share loadMultiple independent pipelines on same topic
Offset trackingGroup commits shared offset per partitionEach consumer tracks own offset
StreamHubAnalytics workers (group: analytics)Audit logger + analytics both read viewer.events
  • Partition assignment

    Consumer group (competing)Each partition → one consumer in group
    Independent consumer (broadcast)All partitions → every consumer
  • Use case

    Consumer group (competing)Scale processing — N consumers share load
    Independent consumer (broadcast)Multiple independent pipelines on same topic
  • Offset tracking

    Consumer group (competing)Group commits shared offset per partition
    Independent consumer (broadcast)Each consumer tracks own offset
  • StreamHub

    Consumer group (competing)Analytics workers (group: analytics)
    Independent consumer (broadcast)Audit logger + analytics both read viewer.events

Quick recall

Everything you need if you only revisit this box.

  • Topics split into partitions; partition key controls ordering and placement.
  • Consumer groups: one consumer per partition in group; scale by adding partitions + consumers.
  • Replicate to 3 brokers; ack after min.insync.replicas for durability without blocking on slow replica.
  • Log-structured storage: append-only segments with sparse indexes; zero-copy to network.
  • Retention by time/size enables replay — reset offset and re-read history.
  • Rebalance on consumer join/leave; use sticky assignment to minimise partition movement.

Test yourself

Answer these before moving on — recall is what makes it stick.