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
API design
# 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).
# 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
| Aspect | Consumer group (competing) | Independent consumer (broadcast) |
|---|---|---|
| Partition assignment | Each partition → one consumer in group | All partitions → every consumer |
| Use case | Scale processing — N consumers share load | Multiple independent pipelines on same topic |
| Offset tracking | Group commits shared offset per partition | Each consumer tracks own offset |
| StreamHub | Analytics workers (group: analytics) | Audit logger + analytics both read viewer.events |
Partition assignment
Consumer group (competing)Each partition → one consumer in groupIndependent consumer (broadcast)All partitions → every consumerUse case
Consumer group (competing)Scale processing — N consumers share loadIndependent consumer (broadcast)Multiple independent pipelines on same topicOffset tracking
Consumer group (competing)Group commits shared offset per partitionIndependent consumer (broadcast)Each consumer tracks own offsetStreamHub
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.