Functional Requirements:
Non-Functional Requirements:
Offer both gRPC and HTTP/REST endpoints for admin and simple produce; clients primarily use native language clients (high-performance TCP binary protocol over TLS).
{
"ProducerAPI": {
"publishSingle": {
"method": "POST",
"url": "/topics/{topic_name}/publish",
"request": {
"key": "string (optional)",
"value": "JSON object",
"timestamp": "ISO8601 (optional)",
"headers": "JSON object (optional)"
},
"response": {
"status": "success/failure",
"topic": "string",
"partition": "int",
"offset": "int",
"timestamp": "ISO8601"
}
},
"publishBatch": {
"method": "POST",
"url": "/topics/{topic_name}/publish/batch",
"request": {
"messages": [
{ "key": "string (optional)", "value": "JSON object" }
]
},
"response": {
"status": "success/failure",
"results": [
{ "partition": "int", "offset": "int" }
]
}
}
},
"ConsumerAPI": {
"subscribe": {
"method": "POST",
"url": "/consumer-groups/{group_id}/subscribe",
"request": {
"topics": ["string"]
},
"response": {
"status": "subscribed",
"consumer_id": "string",
"assigned_partitions": { "topic_name": ["partition_numbers"] }
}
},
"fetchMessages": {
"method": "GET",
"url": "/consumer-groups/{group_id}/fetch",
"queryParams": {
"timeout": "ms (optional)",
"max_records": "int (optional)"
},
"requestBody": {
"topics": ["string (optional)"],
"partitions": ["int (optional)"]
},
"response": {
"messages": [
{
"topic": "string",
"partition": "int",
"offset": "int",
"key": "string",
"value": "JSON object",
"timestamp": "ISO8601"
}
]
}
},
"commitOffsets": {
"method": "POST",
"url": "/consumer-groups/{group_id}/commit",
"request": {
"offsets": [
{ "topic": "string", "partition": "int", "offset": "int" }
]
},
"response": {
"status": "success/failure",
"committed_offsets": [
{ "topic": "string", "partition": "int", "offset": "int" }
]
}
},
"seekOffset": {
"method": "POST",
"url": "/consumer-groups/{group_id}/seek",
"request": {
"topic": "string",
"partition": "int",
"offset": "int"
},
"response": {
"status": "success",
"message": "Consumer will read from specified offset"
}
}
},
"AdminAPI": {
"createTopic": {
"method": "POST",
"url": "/topics",
"request": {
"name": "string",
"partitions": "int",
"replication_factor": "int",
"retention_ms": "int"
},
"response": {
"status": "success/failure",
"topic": "string",
"partitions": "int"
}
},
"listTopics": {
"method": "GET",
"url": "/topics",
"response": {
"topics": [
{ "name": "string", "partitions": "int" }
]
}
},
"getConsumerGroupOffsets": {
"method": "GET",
"url": "/consumer-groups/{group_id}/offsets",
"response": {
"consumer_group": "string",
"offsets": [
{ "topic": "string", "partition": "int", "offset": "int" }
]
}
},
"getClusterStatus": {
"method": "GET",
"url": "/cluster/status",
"response": {
"brokers": [
{ "id": "string", "host": "string", "status": "ACTIVE/DOWN" }
],
"topics": [
{ "name": "string", "partitions": "int" }
]
}
}
}
}
Producers — publish messages to topics/queues.
- Connects to broker leader for target partition
- Supports batching, retries, idempotence
- Assigns keys for partitioning and ordering
- Uses acknowledgments to confirm durable write
Brokers — cluster of servers responsible for storing and serving messages (log-based storage).
- Maintains commit logs (append-only segments)
- Handles replication to follower brokers
- Routes messages to consumers
- Tracks offsets and consumer group states
- Performs compaction and retention cleanup
Metadata/Coordination service — leader election, cluster membership, topic metadata (can be consensus-based: Raft/ZooKeeper/etcd).
- Leader election for topic partitions
- Tracks broker liveness (heartbeats)
- Manages consumer group membership
- Stores configuration (topic metadata, offsets)
- Detects and triggers failover
Consumers — pull or receive messages; support consumer groups (competing consumers).
- Subscribes to topics (via group)
- Fetches from assigned partitions
- Acknowledges or commits offsets
- Handles redelivery and recovery
- Supports different ack modes (auto/manual)
Coordinator/Controller — manages partitions, rebalances consumers, enforces quotas, controller leader.
Storage layer — append-only log files per partition with segmenting, compaction, retention policies.
- Stores messages in partitioned commit logs on disk
- Handles log segments, indexing, compaction
- Replicates data across brokers for durability
Replicator — intra-cluster replication (sync/async) for high availability.
Delivery & offset manager — track per-consumer offsets/ACKs.
Admin API & Management UI — create topics, set retention, view status.
Observability stack — metrics, tracing, logging, alerting.
Clients SDKs — for major languages, supporting batching, compression, retries, transactions.
Storage model:
Consensus/replication:
Fault tolerance:
Sequence (produce/consume):
Responsibilities:
Brokers are the heart of the system. Each broker stores partitions, handles replication, and serves consumers.
Responsibilities:
Consumers read messages from brokers and acknowledge progress.
Responsibilities:
Purpose: Maintain cluster state and enable consistent partition leadership.
Responsibilities:
Typical Implementations:
Data Stored:
Purpose: Efficiently persist messages with durability and fast access.
Design:
Responsibilities:
reply_to topic + correlation_id. Server consumes, processes, produces response to reply_to topic with same correlation_id. Client consumes responses filtered by correlation id (or uses temporary exclusive reply queue). Provide timeouts and idempotence on server.Config knobs:
acks (0/1/all), retries, enable.idempotence, transactional.id, min.insync.replicas, replication.factor, retention.ms, cleanup.policy.batch.size, linger.ms, compression (lz4/snappy), acks=all for durability.Capacity planning formula (simplified):
Replication & ISR:
min.insync.replicas=2 ensures writes succeed only if at least 2 replicas in sync.Hardware:
trace_id header for distributed tracing. Integrate with Jaeger/Zipkin.Recovery flows:
Auditing:
Database primarily needs to store:
Schema:
"Topics": {
"topic_id": "UUID",
"name": "VARCHAR",
"partitions": "INT",
"replication_factor": "INT",
"retention_ms": "BIGINT",
"created_at": "TIMESTAMP"
},
"Partitions": {
"partition_id": "UUID",
"topic_id": "UUID",
"partition_index": "INT",
"leader_broker_id": "UUID",
"replicas": "ARRAY[UUID]",
"created_at": "TIMESTAMP"
},
"Brokers": {
"broker_id": "UUID",
"host": "VARCHAR",
"port": "INT",
"status": "ENUM('ACTIVE','DOWN')",
"last_heartbeat": "TIMESTAMP"
},
"ConsumerGroups": {
"consumer_group_id": "UUID",
"name": "VARCHAR",
"created_at": "TIMESTAMP"
},
"Consumers": {
"consumer_id": "UUID",
"consumer_group_id": "UUID",
"host": "VARCHAR",
"status": "ENUM('ACTIVE','INACTIVE')",
"last_heartbeat": "TIMESTAMP"
},
"Offsets": {
"offset_id": "UUID",
"consumer_group_id": "UUID",
"topic_id": "UUID",
"partition_id": "UUID",
"offset": "BIGINT",
"updated_at": "TIMESTAMP"
},
"Messages": {
"message_id": "UUID",
"topic_id": "UUID",
"partition_id": "UUID",
"key": "VARCHAR",
"payload": "BLOB/JSON",
"created_at": "TIMESTAMP",
"delivered": "BOOLEAN"
},
Messages (the actual payloads) - store in append-only log (custom file storage) or distributed commit log (Cassandra).
Metadata - store in relational DB (PostgreSQL, MySQL) or Key-Value store (etcd, Consul, ZooKeeper).
Offsets / Consumer Tracking - store in Key-Value store (Redis).
Monitoring & Logs - Time-series DB (Prometheus, InfluxDB).