Designing a Pub/Sub Messaging System (Mini-Kafka)
Difficulty: Hard Patterns: Observer, Iterator, Factory Asked at: Uber, LinkedIn, Flipkart, Swiggy, Confluent
Functional Requirements
- Create topics with a configurable number of partitions for parallel message processing
- Publish messages to a topic; messages are routed to partitions via a key-based or round-robin strategy
- Subscribe consumers to one or more topics and consume messages sequentially from assigned partitions
- Consumer groups β multiple consumers form a group; each partition is assigned to exactly one consumer within the group for load balancing
- Offset tracking β each consumer (or consumer group) tracks its read position per partition independently, enabling replay
- At-least-once delivery β messages are not removed after consumption; consumers must explicitly commit offsets to advance
Non-Functional Requirements
- Thread-safety β concurrent producers and consumers must not corrupt shared state (messages, offsets)
- Ordering guarantee within a partition β messages in a single partition are consumed in the exact order they were published
- Independent consumer progress β different consumer groups reading the same partition maintain separate offsets without interference
Core Entities
| Entity | Description |
|---|---|
Topic |
Named channel that holds one or more partitions. Entry point for publishers. |
Partition |
Ordered, append-only log of messages within a topic. Unit of parallelism. |
Message |
Immutable payload with an offset, optional key, value, and timestamp. |
Producer |
Publishes messages to a topic, selecting partition via key hash or round-robin. |
Consumer |
Reads messages from assigned partitions, tracks its own offset. |
ConsumerGroup |
Logical group of consumers; partitions are distributed among members. |
Offset |
Integer position in a partitionβs log, tracked per consumer-group per partition. |
Class Diagram
classDiagram
class Message {
-int offset
-String key
-String value
-long timestamp
}
class Partition {
-int id
-List~Message~ messages
-ReentrantReadWriteLock lock
+append(Message) int
+read(int offset, int count) List~Message~
+getLatestOffset() int
}
class Topic {
-String name
-List~Partition~ partitions
+publish(String key, String value)
+getPartition(int id) Partition
+getPartitionCount() int
}
class Producer {
-String id
+publish(Topic, String key, String value)
}
class Consumer {
-String id
-String groupId
-Map~Partition, Integer~ currentOffsets
+poll(Partition, int maxMessages) List~Message~
+commitOffset(Partition, int offset)
+getOffset(Partition) int
}
class ConsumerGroup {
-String groupId
-List~Consumer~ consumers
-Map~Partition, Consumer~ assignment
-Map~Partition, Integer~ committedOffsets
+addConsumer(Consumer)
+removeConsumer(Consumer)
+rebalance(List~Partition~)
+commitOffset(Partition, int offset)
+getCommittedOffset(Partition) int
}
class PubSubBroker {
-Map~String, Topic~ topics
-Map~String, ConsumerGroup~ groups
-ReentrantLock lock
+createTopic(String name, int partitions) Topic
+getTopic(String name) Topic
+createConsumerGroup(String groupId) ConsumerGroup
+subscribe(String groupId, String topicName)
}
PubSubBroker --> Topic
PubSubBroker --> ConsumerGroup
Topic --> Partition
Partition --> Message
ConsumerGroup --> Consumer
ConsumerGroup --> Partition : assignment
Producer --> Topic : publishes to
Consumer --> Partition : reads from
Design Patterns
| Pattern | Where | Why |
|---|---|---|
| Observer | Consumers subscribe to topics; partitions notify on new messages (optional push mode) | Decouples producers from consumers. Adding new consumers requires zero producer changes. |
| Iterator | Consumerβs poll() method iterates over partition log from current offset |
Clean sequential access without exposing internal list. Offset acts as a cursor. |
| Factory | PubSubBroker.createTopic() and createConsumerGroup() |
Centralizes creation logic, ensures consistent initialization and registration. |
Data Structures
| Component | Structure | Why |
|---|---|---|
| Partition message log | ArrayList<Message> |
Append-only, O(1) append, O(1) random access by offset (offset = index). |
| Topic partition list | ArrayList<Partition> |
Fixed at creation, indexed access for partition selection. |
| Consumer group offsets | ConcurrentHashMap<Partition, AtomicInteger> |
O(1) offset read/write, thread-safe without global lock. |
| Partition assignment | HashMap<Partition, Consumer> |
O(1) lookup of which consumer owns a partition during rebalance. |
| Broker topic registry | ConcurrentHashMap<String, Topic> |
O(1) topic lookup by name, safe for concurrent producer access. |
| Round-robin counter | AtomicInteger |
Lock-free partition rotation for keyless messages. |
How It All Fits Together
Hereβs what happens when a producer publishes a message:
- Producer calls
publish(topic, key, value) - If a key is provided, compute
hash(key) % partitionCountto pick a partition (same key always goes to same partition). If no key, use round-robin via atomic counter. - The selected
Partitionacquires its write lock - A
Messageobject is created with offset = current log size, value, key, and timestamp - Message is appended to the partitionβs log (ArrayList)
- Write lock is released. The message is now durable in memory.
Hereβs what happens when a consumer polls for messages:
- Consumer calls
poll(partition, maxMessages) - Consumer looks up its current offset for that partition (from consumer groupβs committed offset or its local tracking)
- Partition acquires a read lock
- Messages from
offsettomin(offset + maxMessages, log.size())are returned - Read lock is released
- Consumer processes messages, then calls
commitOffset(partition, newOffset)to advance - If the consumer crashes before committing, it will re-read from the last committed offset on restart (at-least-once guarantee)
Complete Code
Message.java
An immutable value object representing a single record in the partition log. The offset is assigned by the partition on append and serves as the position index. Key is optional β when present, it determines partition routing.
package pubsub.model;
public class Message {
private final int offset;
private final String key;
private final String value;
private final long timestamp;
public Message(int offset, String key, String value) {
this.offset = offset;
this.key = key;
this.value = value;
this.timestamp = System.currentTimeMillis();
}
public int getOffset() { return offset; }
public String getKey() { return key; }
public String getValue() { return value; }
public long getTimestamp() { return timestamp; }
@Override
public String toString() {
return "Message{offset=" + offset + ", key='" + key + "', value='" + value + "'}";
}
}
import time
class Message:
def __init__(self, offset: int, key: str | None, value: str):
self._offset = offset
self._key = key
self._value = value
self._timestamp = time.time()
@property
def offset(self) -> int:
return self._offset
@property
def key(self) -> str | None:
return self._key
@property
def value(self) -> str:
return self._value
@property
def timestamp(self) -> float:
return self._timestamp
def __repr__(self) -> str:
return f"Message(offset={self._offset}, key='{self._key}', value='{self._value}')"
#pragma once
#include <string>
#include <chrono>
class Message {
private:
int offset_;
std::string key_;
std::string value_;
long long timestamp_;
public:
Message(int offset, std::string key, std::string value)
: offset_(offset), key_(std::move(key)), value_(std::move(value)) {
timestamp_ = std::chrono::duration_cast<std::chrono::milliseconds>(
std::chrono::system_clock::now().time_since_epoch()).count();
}
int getOffset() const { return offset_; }
const std::string& getKey() const { return key_; }
const std::string& getValue() const { return value_; }
long long getTimestamp() const { return timestamp_; }
std::string toString() const {
return "Message{offset=" + std::to_string(offset_) +
", key='" + key_ + "', value='" + value_ + "'}";
}
};
Partition.java
The core building block β an ordered, append-only log. Uses a ReadWriteLock so multiple consumers can read concurrently while writes are exclusive. The offset is simply the index into the list, giving O(1) random access for consumers polling from any position.
package pubsub.model;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
public class Partition {
private final int id;
private final String topicName;
private final List<Message> messages;
private final ReadWriteLock lock;
public Partition(int id, String topicName) {
this.id = id;
this.topicName = topicName;
this.messages = new ArrayList<>();
this.lock = new ReentrantReadWriteLock();
}
/**
* Append a message to this partition. Assigns offset = current size.
* Thread-safe via write lock.
*/
public int append(String key, String value) {
lock.writeLock().lock();
try {
int offset = messages.size();
Message msg = new Message(offset, key, value);
messages.add(msg);
return offset;
} finally {
lock.writeLock().unlock();
}
}
/**
* Read messages starting from the given offset.
* Returns up to maxCount messages.
*/
public List<Message> read(int fromOffset, int maxCount) {
lock.readLock().lock();
try {
if (fromOffset >= messages.size()) {
return Collections.emptyList();
}
int toIndex = Math.min(fromOffset + maxCount, messages.size());
return new ArrayList<>(messages.subList(fromOffset, toIndex));
} finally {
lock.readLock().unlock();
}
}
public int getLatestOffset() {
lock.readLock().lock();
try {
return messages.size();
} finally {
lock.readLock().unlock();
}
}
public int getId() { return id; }
public String getTopicName() { return topicName; }
@Override
public String toString() {
return topicName + "-partition-" + id;
}
@Override
public int hashCode() {
return 31 * topicName.hashCode() + id;
}
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (!(o instanceof Partition)) return false;
Partition p = (Partition) o;
return id == p.id && topicName.equals(p.topicName);
}
}
import threading
from typing import List
class Partition:
def __init__(self, partition_id: int, topic_name: str):
self._id = partition_id
self._topic_name = topic_name
self._messages: List[Message] = []
self._lock = threading.RWLock() if hasattr(threading, 'RWLock') else threading.Lock()
def append(self, key: str | None, value: str) -> int:
"""Append a message. Returns the assigned offset."""
with self._lock:
offset = len(self._messages)
msg = Message(offset, key, value)
self._messages.append(msg)
return offset
def read(self, from_offset: int, max_count: int) -> List[Message]:
"""Read up to max_count messages starting from from_offset."""
with self._lock:
if from_offset >= len(self._messages):
return []
to_index = min(from_offset + max_count, len(self._messages))
return list(self._messages[from_offset:to_index])
@property
def latest_offset(self) -> int:
with self._lock:
return len(self._messages)
@property
def id(self) -> int:
return self._id
@property
def topic_name(self) -> str:
return self._topic_name
def __repr__(self) -> str:
return f"{self._topic_name}-partition-{self._id}"
def __hash__(self) -> int:
return hash((self._topic_name, self._id))
def __eq__(self, other) -> bool:
if not isinstance(other, Partition):
return False
return self._id == other._id and self._topic_name == other._topic_name
#pragma once
#include <vector>
#include <string>
#include <shared_mutex>
#include <algorithm>
#include "Message.hpp"
class Partition {
private:
int id_;
std::string topicName_;
std::vector<Message> messages_;
mutable std::shared_mutex mutex_;
public:
Partition(int id, std::string topicName)
: id_(id), topicName_(std::move(topicName)) {}
// Non-copyable due to mutex
Partition(const Partition&) = delete;
Partition& operator=(const Partition&) = delete;
int append(const std::string& key, const std::string& value) {
std::unique_lock lock(mutex_);
int offset = static_cast<int>(messages_.size());
messages_.emplace_back(offset, key, value);
return offset;
}
std::vector<Message> read(int fromOffset, int maxCount) const {
std::shared_lock lock(mutex_);
if (fromOffset >= static_cast<int>(messages_.size())) {
return {};
}
int toIndex = std::min(fromOffset + maxCount,
static_cast<int>(messages_.size()));
return std::vector<Message>(
messages_.begin() + fromOffset,
messages_.begin() + toIndex);
}
int getLatestOffset() const {
std::shared_lock lock(mutex_);
return static_cast<int>(messages_.size());
}
int getId() const { return id_; }
const std::string& getTopicName() const { return topicName_; }
bool operator==(const Partition& other) const {
return id_ == other.id_ && topicName_ == other.topicName_;
}
};
Topic.java
A topic owns its partitions and provides the publish entry point. It handles partition selection: key-based hashing for ordering guarantees (same key β same partition), or round-robin for even distribution when no key is specified. The AtomicInteger counter makes round-robin lock-free.
package pubsub.model;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.atomic.AtomicInteger;
public class Topic {
private final String name;
private final List<Partition> partitions;
private final AtomicInteger roundRobinCounter;
public Topic(String name, int partitionCount) {
if (partitionCount <= 0) {
throw new IllegalArgumentException("Partition count must be positive");
}
this.name = name;
this.partitions = new ArrayList<>();
this.roundRobinCounter = new AtomicInteger(0);
for (int i = 0; i < partitionCount; i++) {
partitions.add(new Partition(i, name));
}
}
/**
* Publish a message to this topic.
* If key is non-null, routes to partition via key hash.
* If key is null, uses round-robin.
*/
public void publish(String key, String value) {
int partitionIndex;
if (key != null) {
partitionIndex = Math.abs(key.hashCode() % partitions.size());
} else {
partitionIndex = Math.abs(roundRobinCounter.getAndIncrement() % partitions.size());
}
partitions.get(partitionIndex).append(key, value);
}
public Partition getPartition(int id) {
if (id < 0 || id >= partitions.size()) {
throw new IllegalArgumentException("Invalid partition id: " + id);
}
return partitions.get(id);
}
public List<Partition> getPartitions() {
return Collections.unmodifiableList(partitions);
}
public int getPartitionCount() { return partitions.size(); }
public String getName() { return name; }
@Override
public String toString() {
return "Topic{name='" + name + "', partitions=" + partitions.size() + "}";
}
}
import itertools
from typing import List
class Topic:
def __init__(self, name: str, partition_count: int):
if partition_count <= 0:
raise ValueError("Partition count must be positive")
self._name = name
self._partitions = [Partition(i, name) for i in range(partition_count)]
self._round_robin = itertools.cycle(range(partition_count))
def publish(self, key: str | None, value: str) -> None:
"""
Publish a message. Key-based hash for partition routing,
or round-robin if key is None.
"""
if key is not None:
partition_index = abs(hash(key)) % len(self._partitions)
else:
partition_index = next(self._round_robin)
self._partitions[partition_index].append(key, value)
def get_partition(self, partition_id: int) -> Partition:
if partition_id < 0 or partition_id >= len(self._partitions):
raise ValueError(f"Invalid partition id: {partition_id}")
return self._partitions[partition_id]
@property
def partitions(self) -> List[Partition]:
return list(self._partitions)
@property
def partition_count(self) -> int:
return len(self._partitions)
@property
def name(self) -> str:
return self._name
def __repr__(self) -> str:
return f"Topic(name='{self._name}', partitions={len(self._partitions)})"
#pragma once
#include <vector>
#include <string>
#include <memory>
#include <atomic>
#include <cmath>
#include <stdexcept>
#include <functional>
#include "Partition.hpp"
class Topic {
private:
std::string name_;
std::vector<std::unique_ptr<Partition>> partitions_;
std::atomic<int> roundRobinCounter_{0};
public:
Topic(std::string name, int partitionCount)
: name_(std::move(name)) {
if (partitionCount <= 0) {
throw std::invalid_argument("Partition count must be positive");
}
for (int i = 0; i < partitionCount; i++) {
partitions_.push_back(std::make_unique<Partition>(i, name_));
}
}
void publish(const std::string& key, const std::string& value) {
int partitionIndex;
if (!key.empty()) {
partitionIndex = std::abs(static_cast<int>(
std::hash<std::string>{}(key) % partitions_.size()));
} else {
partitionIndex = std::abs(
roundRobinCounter_.fetch_add(1) % static_cast<int>(partitions_.size()));
}
partitions_[partitionIndex]->append(key, value);
}
Partition* getPartition(int id) {
if (id < 0 || id >= static_cast<int>(partitions_.size())) {
throw std::invalid_argument("Invalid partition id");
}
return partitions_[id].get();
}
int getPartitionCount() const {
return static_cast<int>(partitions_.size());
}
const std::string& getName() const { return name_; }
// Get raw pointers to partitions for iteration
std::vector<Partition*> getPartitions() {
std::vector<Partition*> result;
for (auto& p : partitions_) {
result.push_back(p.get());
}
return result;
}
};
Consumer.java
A consumer belongs to a consumer group and reads from its assigned partitions. It maintains a local view of offsets (what it has read so far) and commits back to the group. The poll + commit two-phase approach enables at-least-once delivery: if the consumer crashes after poll but before commit, messages will be re-delivered.
package pubsub.model;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
public class Consumer {
private final String id;
private final String groupId;
private final Map<Partition, Integer> localOffsets;
public Consumer(String id, String groupId) {
this.id = id;
this.groupId = groupId;
this.localOffsets = new ConcurrentHashMap<>();
}
/**
* Poll messages from a partition starting at the current local offset.
* Does NOT advance committed offset - call commitOffset() after processing.
*/
public List<Message> poll(Partition partition, int maxMessages) {
int currentOffset = localOffsets.getOrDefault(partition, 0);
List<Message> messages = partition.read(currentOffset, maxMessages);
if (!messages.isEmpty()) {
// Advance local offset (not committed yet)
localOffsets.put(partition, currentOffset + messages.size());
}
return messages;
}
/**
* Commit the offset - signals that all messages up to this offset
* have been successfully processed.
*/
public void commitOffset(Partition partition, int offset) {
localOffsets.put(partition, offset);
}
/**
* Reset local offset to a specific position (for replay scenarios).
*/
public void seek(Partition partition, int offset) {
localOffsets.put(partition, offset);
}
public int getOffset(Partition partition) {
return localOffsets.getOrDefault(partition, 0);
}
public String getId() { return id; }
public String getGroupId() { return groupId; }
@Override
public String toString() {
return "Consumer{id='" + id + "', group='" + groupId + "'}";
}
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (!(o instanceof Consumer)) return false;
return id.equals(((Consumer) o).id);
}
@Override
public int hashCode() {
return id.hashCode();
}
}
from typing import List, Dict
class Consumer:
def __init__(self, consumer_id: str, group_id: str):
self._id = consumer_id
self._group_id = group_id
self._local_offsets: Dict[Partition, int] = {}
def poll(self, partition: Partition, max_messages: int) -> List[Message]:
"""
Poll messages from a partition starting at current local offset.
Does NOT commit - call commit_offset() after processing.
"""
current_offset = self._local_offsets.get(partition, 0)
messages = partition.read(current_offset, max_messages)
if messages:
self._local_offsets[partition] = current_offset + len(messages)
return messages
def commit_offset(self, partition: Partition, offset: int) -> None:
"""Commit offset after successful processing."""
self._local_offsets[partition] = offset
def seek(self, partition: Partition, offset: int) -> None:
"""Reset offset to a specific position for replay."""
self._local_offsets[partition] = offset
def get_offset(self, partition: Partition) -> int:
return self._local_offsets.get(partition, 0)
@property
def id(self) -> str:
return self._id
@property
def group_id(self) -> str:
return self._group_id
def __repr__(self) -> str:
return f"Consumer(id='{self._id}', group='{self._group_id}')"
def __eq__(self, other) -> bool:
if not isinstance(other, Consumer):
return False
return self._id == other._id
def __hash__(self) -> int:
return hash(self._id)
#pragma once
#include <string>
#include <unordered_map>
#include <vector>
#include "Partition.hpp"
#include "Message.hpp"
class Consumer {
private:
std::string id_;
std::string groupId_;
std::unordered_map<Partition*, int> localOffsets_;
public:
Consumer(std::string id, std::string groupId)
: id_(std::move(id)), groupId_(std::move(groupId)) {}
std::vector<Message> poll(Partition* partition, int maxMessages) {
int currentOffset = 0;
auto it = localOffsets_.find(partition);
if (it != localOffsets_.end()) {
currentOffset = it->second;
}
auto messages = partition->read(currentOffset, maxMessages);
if (!messages.empty()) {
localOffsets_[partition] = currentOffset + static_cast<int>(messages.size());
}
return messages;
}
void commitOffset(Partition* partition, int offset) {
localOffsets_[partition] = offset;
}
void seek(Partition* partition, int offset) {
localOffsets_[partition] = offset;
}
int getOffset(Partition* partition) const {
auto it = localOffsets_.find(partition);
return it != localOffsets_.end() ? it->second : 0;
}
const std::string& getId() const { return id_; }
const std::string& getGroupId() const { return groupId_; }
bool operator==(const Consumer& other) const {
return id_ == other.id_;
}
};
ConsumerGroup.java
Manages partition-to-consumer assignment and committed offsets. When consumers join or leave, rebalance() redistributes partitions evenly (round-robin assignment). Committed offsets are stored at the group level, not the individual consumer β this way if a consumer dies, its replacement picks up from the last committed position.
package pubsub.model;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.locks.ReentrantLock;
public class ConsumerGroup {
private final String groupId;
private final List<Consumer> consumers;
private final Map<Partition, Consumer> partitionAssignment;
private final Map<Partition, Integer> committedOffsets;
private final ReentrantLock lock;
private List<Partition> subscribedPartitions;
public ConsumerGroup(String groupId) {
this.groupId = groupId;
this.consumers = new ArrayList<>();
this.partitionAssignment = new HashMap<>();
this.committedOffsets = new ConcurrentHashMap<>();
this.lock = new ReentrantLock();
this.subscribedPartitions = new ArrayList<>();
}
public void addConsumer(Consumer consumer) {
lock.lock();
try {
consumers.add(consumer);
rebalance();
} finally {
lock.unlock();
}
}
public void removeConsumer(Consumer consumer) {
lock.lock();
try {
consumers.remove(consumer);
rebalance();
} finally {
lock.unlock();
}
}
public void subscribe(List<Partition> partitions) {
lock.lock();
try {
this.subscribedPartitions = new ArrayList<>(partitions);
rebalance();
} finally {
lock.unlock();
}
}
/**
* Rebalance: assign partitions to consumers using round-robin.
* Each partition goes to exactly one consumer in the group.
*/
private void rebalance() {
partitionAssignment.clear();
if (consumers.isEmpty() || subscribedPartitions.isEmpty()) {
return;
}
for (int i = 0; i < subscribedPartitions.size(); i++) {
Consumer assignedConsumer = consumers.get(i % consumers.size());
partitionAssignment.put(subscribedPartitions.get(i), assignedConsumer);
}
}
/**
* Commit offset for a partition. This is the group-level committed offset.
* On consumer failure, the replacement starts from here.
*/
public void commitOffset(Partition partition, int offset) {
committedOffsets.put(partition, offset);
}
public int getCommittedOffset(Partition partition) {
return committedOffsets.getOrDefault(partition, 0);
}
public Consumer getAssignedConsumer(Partition partition) {
return partitionAssignment.get(partition);
}
/**
* Get partitions assigned to a specific consumer.
*/
public List<Partition> getAssignedPartitions(Consumer consumer) {
List<Partition> result = new ArrayList<>();
for (Map.Entry<Partition, Consumer> entry : partitionAssignment.entrySet()) {
if (entry.getValue().equals(consumer)) {
result.add(entry.getKey());
}
}
return result;
}
public String getGroupId() { return groupId; }
public List<Consumer> getConsumers() { return Collections.unmodifiableList(consumers); }
@Override
public String toString() {
return "ConsumerGroup{id='" + groupId + "', consumers=" + consumers.size() +
", partitions=" + subscribedPartitions.size() + "}";
}
}
import threading
from typing import List, Dict, Optional
class ConsumerGroup:
def __init__(self, group_id: str):
self._group_id = group_id
self._consumers: List[Consumer] = []
self._partition_assignment: Dict[Partition, Consumer] = {}
self._committed_offsets: Dict[Partition, int] = {}
self._subscribed_partitions: List[Partition] = []
self._lock = threading.Lock()
def add_consumer(self, consumer: Consumer) -> None:
with self._lock:
self._consumers.append(consumer)
self._rebalance()
def remove_consumer(self, consumer: Consumer) -> None:
with self._lock:
self._consumers.remove(consumer)
self._rebalance()
def subscribe(self, partitions: List[Partition]) -> None:
with self._lock:
self._subscribed_partitions = list(partitions)
self._rebalance()
def _rebalance(self) -> None:
"""Round-robin partition assignment across consumers."""
self._partition_assignment.clear()
if not self._consumers or not self._subscribed_partitions:
return
for i, partition in enumerate(self._subscribed_partitions):
consumer = self._consumers[i % len(self._consumers)]
self._partition_assignment[partition] = consumer
def commit_offset(self, partition: Partition, offset: int) -> None:
self._committed_offsets[partition] = offset
def get_committed_offset(self, partition: Partition) -> int:
return self._committed_offsets.get(partition, 0)
def get_assigned_consumer(self, partition: Partition) -> Optional[Consumer]:
return self._partition_assignment.get(partition)
def get_assigned_partitions(self, consumer: Consumer) -> List[Partition]:
return [p for p, c in self._partition_assignment.items() if c == consumer]
@property
def group_id(self) -> str:
return self._group_id
@property
def consumers(self) -> List[Consumer]:
return list(self._consumers)
def __repr__(self) -> str:
return (f"ConsumerGroup(id='{self._group_id}', "
f"consumers={len(self._consumers)}, "
f"partitions={len(self._subscribed_partitions)})")
#pragma once
#include <string>
#include <vector>
#include <unordered_map>
#include <mutex>
#include <memory>
#include "Consumer.hpp"
#include "Partition.hpp"
class ConsumerGroup {
private:
std::string groupId_;
std::vector<Consumer*> consumers_;
std::unordered_map<Partition*, Consumer*> partitionAssignment_;
std::unordered_map<Partition*, int> committedOffsets_;
std::vector<Partition*> subscribedPartitions_;
mutable std::mutex mutex_;
void rebalance() {
partitionAssignment_.clear();
if (consumers_.empty() || subscribedPartitions_.empty()) return;
for (size_t i = 0; i < subscribedPartitions_.size(); i++) {
Consumer* assigned = consumers_[i % consumers_.size()];
partitionAssignment_[subscribedPartitions_[i]] = assigned;
}
}
public:
ConsumerGroup(std::string groupId)
: groupId_(std::move(groupId)) {}
void addConsumer(Consumer* consumer) {
std::lock_guard lock(mutex_);
consumers_.push_back(consumer);
rebalance();
}
void removeConsumer(Consumer* consumer) {
std::lock_guard lock(mutex_);
consumers_.erase(
std::remove(consumers_.begin(), consumers_.end(), consumer),
consumers_.end());
rebalance();
}
void subscribe(const std::vector<Partition*>& partitions) {
std::lock_guard lock(mutex_);
subscribedPartitions_ = partitions;
rebalance();
}
void commitOffset(Partition* partition, int offset) {
std::lock_guard lock(mutex_);
committedOffsets_[partition] = offset;
}
int getCommittedOffset(Partition* partition) const {
std::lock_guard lock(mutex_);
auto it = committedOffsets_.find(partition);
return it != committedOffsets_.end() ? it->second : 0;
}
Consumer* getAssignedConsumer(Partition* partition) const {
std::lock_guard lock(mutex_);
auto it = partitionAssignment_.find(partition);
return it != partitionAssignment_.end() ? it->second : nullptr;
}
std::vector<Partition*> getAssignedPartitions(Consumer* consumer) const {
std::lock_guard lock(mutex_);
std::vector<Partition*> result;
for (const auto& [partition, assigned] : partitionAssignment_) {
if (assigned == consumer) {
result.push_back(partition);
}
}
return result;
}
const std::string& getGroupId() const { return groupId_; }
};
Producer.java
A thin client that wraps the publish call. Each producer has an ID for tracing. In a real system this would handle batching, retries, and acknowledgments β here itβs kept simple to focus on the routing logic.
package pubsub.model;
public class Producer {
private final String id;
public Producer(String id) {
this.id = id;
}
/**
* Publish a message to the given topic.
* Key determines partition (null key = round-robin).
*/
public void publish(Topic topic, String key, String value) {
if (topic == null) {
throw new IllegalArgumentException("Topic cannot be null");
}
if (value == null) {
throw new IllegalArgumentException("Value cannot be null");
}
topic.publish(key, value);
}
public String getId() { return id; }
@Override
public String toString() {
return "Producer{id='" + id + "'}";
}
}
class Producer:
def __init__(self, producer_id: str):
self._id = producer_id
def publish(self, topic: Topic, key: str | None, value: str) -> None:
"""Publish a message to the given topic."""
if topic is None:
raise ValueError("Topic cannot be None")
if value is None:
raise ValueError("Value cannot be None")
topic.publish(key, value)
@property
def id(self) -> str:
return self._id
def __repr__(self) -> str:
return f"Producer(id='{self._id}')"
#pragma once
#include <string>
#include <stdexcept>
#include "Topic.hpp"
class Producer {
private:
std::string id_;
public:
Producer(std::string id) : id_(std::move(id)) {}
void publish(Topic* topic, const std::string& key, const std::string& value) {
if (!topic) {
throw std::invalid_argument("Topic cannot be null");
}
if (value.empty()) {
throw std::invalid_argument("Value cannot be empty");
}
topic->publish(key, value);
}
const std::string& getId() const { return id_; }
};
PubSubBroker.java
The central coordinator β the βmini-Kafka broker.β It manages topic creation, consumer group registration, and subscription wiring. Thread-safe via a ReentrantLock for topic/group creation (infrequent operations). Once created, topics and groups are accessed lock-free via ConcurrentHashMap.
package pubsub;
import pubsub.model.*;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.locks.ReentrantLock;
public class PubSubBroker {
private final Map<String, Topic> topics;
private final Map<String, ConsumerGroup> consumerGroups;
private final ReentrantLock lock;
public PubSubBroker() {
this.topics = new ConcurrentHashMap<>();
this.consumerGroups = new ConcurrentHashMap<>();
this.lock = new ReentrantLock();
}
/**
* Create a new topic with the specified number of partitions.
*/
public Topic createTopic(String name, int partitionCount) {
lock.lock();
try {
if (topics.containsKey(name)) {
throw new IllegalArgumentException("Topic '" + name + "' already exists");
}
Topic topic = new Topic(name, partitionCount);
topics.put(name, topic);
return topic;
} finally {
lock.unlock();
}
}
public Topic getTopic(String name) {
Topic topic = topics.get(name);
if (topic == null) {
throw new IllegalArgumentException("Topic '" + name + "' not found");
}
return topic;
}
/**
* Create a consumer group. Consumers in the same group share partitions.
*/
public ConsumerGroup createConsumerGroup(String groupId) {
lock.lock();
try {
if (consumerGroups.containsKey(groupId)) {
throw new IllegalArgumentException("Group '" + groupId + "' already exists");
}
ConsumerGroup group = new ConsumerGroup(groupId);
consumerGroups.put(groupId, group);
return group;
} finally {
lock.unlock();
}
}
public ConsumerGroup getConsumerGroup(String groupId) {
ConsumerGroup group = consumerGroups.get(groupId);
if (group == null) {
throw new IllegalArgumentException("Group '" + groupId + "' not found");
}
return group;
}
/**
* Subscribe a consumer group to a topic.
* Triggers partition assignment among group members.
*/
public void subscribe(String groupId, String topicName) {
ConsumerGroup group = getConsumerGroup(groupId);
Topic topic = getTopic(topicName);
group.subscribe(topic.getPartitions());
}
/**
* Register a consumer into a group. Triggers rebalance.
*/
public void registerConsumer(String groupId, Consumer consumer) {
ConsumerGroup group = getConsumerGroup(groupId);
group.addConsumer(consumer);
}
/**
* Deregister a consumer from its group. Triggers rebalance.
*/
public void deregisterConsumer(String groupId, Consumer consumer) {
ConsumerGroup group = getConsumerGroup(groupId);
group.removeConsumer(consumer);
}
public boolean hasTopic(String name) {
return topics.containsKey(name);
}
@Override
public String toString() {
return "PubSubBroker{topics=" + topics.size() +
", groups=" + consumerGroups.size() + "}";
}
}
import threading
from typing import Dict
class PubSubBroker:
def __init__(self):
self._topics: Dict[str, Topic] = {}
self._consumer_groups: Dict[str, ConsumerGroup] = {}
self._lock = threading.Lock()
def create_topic(self, name: str, partition_count: int) -> Topic:
"""Create a new topic with the specified number of partitions."""
with self._lock:
if name in self._topics:
raise ValueError(f"Topic '{name}' already exists")
topic = Topic(name, partition_count)
self._topics[name] = topic
return topic
def get_topic(self, name: str) -> Topic:
if name not in self._topics:
raise ValueError(f"Topic '{name}' not found")
return self._topics[name]
def create_consumer_group(self, group_id: str) -> ConsumerGroup:
"""Create a consumer group."""
with self._lock:
if group_id in self._consumer_groups:
raise ValueError(f"Group '{group_id}' already exists")
group = ConsumerGroup(group_id)
self._consumer_groups[group_id] = group
return group
def get_consumer_group(self, group_id: str) -> ConsumerGroup:
if group_id not in self._consumer_groups:
raise ValueError(f"Group '{group_id}' not found")
return self._consumer_groups[group_id]
def subscribe(self, group_id: str, topic_name: str) -> None:
"""Subscribe a consumer group to a topic."""
group = self.get_consumer_group(group_id)
topic = self.get_topic(topic_name)
group.subscribe(topic.partitions)
def register_consumer(self, group_id: str, consumer: Consumer) -> None:
"""Register a consumer into a group. Triggers rebalance."""
group = self.get_consumer_group(group_id)
group.add_consumer(consumer)
def deregister_consumer(self, group_id: str, consumer: Consumer) -> None:
"""Deregister a consumer from its group. Triggers rebalance."""
group = self.get_consumer_group(group_id)
group.remove_consumer(consumer)
def has_topic(self, name: str) -> bool:
return name in self._topics
def __repr__(self) -> str:
return (f"PubSubBroker(topics={len(self._topics)}, "
f"groups={len(self._consumer_groups)})")
#pragma once
#include <string>
#include <unordered_map>
#include <memory>
#include <mutex>
#include <stdexcept>
#include "Topic.hpp"
#include "ConsumerGroup.hpp"
#include "Consumer.hpp"
class PubSubBroker {
private:
std::unordered_map<std::string, std::unique_ptr<Topic>> topics_;
std::unordered_map<std::string, std::unique_ptr<ConsumerGroup>> consumerGroups_;
mutable std::mutex mutex_;
public:
Topic* createTopic(const std::string& name, int partitionCount) {
std::lock_guard lock(mutex_);
if (topics_.count(name)) {
throw std::invalid_argument("Topic '" + name + "' already exists");
}
topics_[name] = std::make_unique<Topic>(name, partitionCount);
return topics_[name].get();
}
Topic* getTopic(const std::string& name) {
auto it = topics_.find(name);
if (it == topics_.end()) {
throw std::invalid_argument("Topic '" + name + "' not found");
}
return it->second.get();
}
ConsumerGroup* createConsumerGroup(const std::string& groupId) {
std::lock_guard lock(mutex_);
if (consumerGroups_.count(groupId)) {
throw std::invalid_argument("Group '" + groupId + "' already exists");
}
consumerGroups_[groupId] = std::make_unique<ConsumerGroup>(groupId);
return consumerGroups_[groupId].get();
}
ConsumerGroup* getConsumerGroup(const std::string& groupId) {
auto it = consumerGroups_.find(groupId);
if (it == consumerGroups_.end()) {
throw std::invalid_argument("Group '" + groupId + "' not found");
}
return it->second.get();
}
void subscribe(const std::string& groupId, const std::string& topicName) {
auto* group = getConsumerGroup(groupId);
auto* topic = getTopic(topicName);
group->subscribe(topic->getPartitions());
}
void registerConsumer(const std::string& groupId, Consumer* consumer) {
auto* group = getConsumerGroup(groupId);
group->addConsumer(consumer);
}
void deregisterConsumer(const std::string& groupId, Consumer* consumer) {
auto* group = getConsumerGroup(groupId);
group->removeConsumer(consumer);
}
bool hasTopic(const std::string& name) const {
return topics_.count(name) > 0;
}
};
Main.java (Demo)
A runnable demo that wires everything together: creates a topic with 3 partitions, two consumer groups, publishes messages with keys (ensuring ordering within a partition), and demonstrates independent consumption across groups.
package pubsub;
import pubsub.model.*;
import java.util.List;
public class Main {
public static void main(String[] args) {
// 1. Create the broker
PubSubBroker broker = new PubSubBroker();
// 2. Create a topic with 3 partitions
Topic ordersTopic = broker.createTopic("orders", 3);
System.out.println("Created: " + ordersTopic);
// 3. Create two consumer groups (independent consumption)
ConsumerGroup analyticsGroup = broker.createConsumerGroup("analytics-group");
ConsumerGroup notificationGroup = broker.createConsumerGroup("notification-group");
// 4. Create consumers and register them
Consumer analyticsConsumer1 = new Consumer("analytics-1", "analytics-group");
Consumer analyticsConsumer2 = new Consumer("analytics-2", "analytics-group");
Consumer notifConsumer = new Consumer("notif-1", "notification-group");
broker.registerConsumer("analytics-group", analyticsConsumer1);
broker.registerConsumer("analytics-group", analyticsConsumer2);
broker.registerConsumer("notification-group", notifConsumer);
// 5. Subscribe groups to the topic
broker.subscribe("analytics-group", "orders");
broker.subscribe("notification-group", "orders");
// 6. Produce messages with keys (key ensures same-partition routing)
Producer producer = new Producer("order-service");
producer.publish(ordersTopic, "user-123", "Order #1001 placed");
producer.publish(ordersTopic, "user-456", "Order #1002 placed");
producer.publish(ordersTopic, "user-123", "Order #1001 shipped"); // same partition as #1001
producer.publish(ordersTopic, "user-789", "Order #1003 placed");
producer.publish(ordersTopic, null, "Broadcast: flash sale started"); // round-robin
System.out.println("\nPublished 5 messages to 'orders' topic");
// 7. Analytics group - consumer 1 reads its assigned partitions
System.out.println("\n--- Analytics Group ---");
List<Partition> assigned1 = analyticsGroup.getAssignedPartitions(analyticsConsumer1);
for (Partition p : assigned1) {
List<Message> msgs = analyticsConsumer1.poll(p, 10);
System.out.println(" " + analyticsConsumer1.getId() + " from " + p + ": " + msgs);
}
List<Partition> assigned2 = analyticsGroup.getAssignedPartitions(analyticsConsumer2);
for (Partition p : assigned2) {
List<Message> msgs = analyticsConsumer2.poll(p, 10);
System.out.println(" " + analyticsConsumer2.getId() + " from " + p + ": " + msgs);
}
// 8. Notification group - single consumer reads ALL partitions
System.out.println("\n--- Notification Group (independent progress) ---");
List<Partition> notifAssigned = notificationGroup.getAssignedPartitions(notifConsumer);
for (Partition p : notifAssigned) {
List<Message> msgs = notifConsumer.poll(p, 10);
System.out.println(" " + notifConsumer.getId() + " from " + p + ": " + msgs);
// Commit offset after processing
notificationGroup.commitOffset(p, notifConsumer.getOffset(p));
}
// 9. Demonstrate offset independence - notification group re-polls (gets nothing new)
System.out.println("\n--- Re-poll notification group (should be empty - already consumed) ---");
for (Partition p : notifAssigned) {
List<Message> msgs = notifConsumer.poll(p, 10);
System.out.println(" " + notifConsumer.getId() + " from " + p + ": " + msgs);
}
// 10. Demonstrate seek/replay
System.out.println("\n--- Replay: seek notification consumer to offset 0 ---");
for (Partition p : notifAssigned) {
notifConsumer.seek(p, 0);
List<Message> msgs = notifConsumer.poll(p, 10);
System.out.println(" REPLAY " + p + ": " + msgs);
}
}
}
"""
Demo: Mini Pub/Sub Messaging System (Mini-Kafka)
"""
import time
import threading
def main():
# 1. Create the broker
broker = PubSubBroker()
# 2. Create a topic with 3 partitions
orders_topic = broker.create_topic("orders", 3)
print(f"Created: {orders_topic}")
# 3. Create two consumer groups (independent consumption)
analytics_group = broker.create_consumer_group("analytics-group")
notification_group = broker.create_consumer_group("notification-group")
# 4. Create consumers and register them
analytics_consumer1 = Consumer("analytics-1", "analytics-group")
analytics_consumer2 = Consumer("analytics-2", "analytics-group")
notif_consumer = Consumer("notif-1", "notification-group")
broker.register_consumer("analytics-group", analytics_consumer1)
broker.register_consumer("analytics-group", analytics_consumer2)
broker.register_consumer("notification-group", notif_consumer)
# 5. Subscribe groups to the topic
broker.subscribe("analytics-group", "orders")
broker.subscribe("notification-group", "orders")
# 6. Produce messages with keys
producer = Producer("order-service")
producer.publish(orders_topic, "user-123", "Order #1001 placed")
producer.publish(orders_topic, "user-456", "Order #1002 placed")
producer.publish(orders_topic, "user-123", "Order #1001 shipped")
producer.publish(orders_topic, "user-789", "Order #1003 placed")
producer.publish(orders_topic, None, "Broadcast: flash sale started")
print("\nPublished 5 messages to 'orders' topic")
# 7. Analytics group consumption
print("\n--- Analytics Group ---")
for p in analytics_group.get_assigned_partitions(analytics_consumer1):
msgs = analytics_consumer1.poll(p, 10)
print(f" {analytics_consumer1.id} from {p}: {msgs}")
for p in analytics_group.get_assigned_partitions(analytics_consumer2):
msgs = analytics_consumer2.poll(p, 10)
print(f" {analytics_consumer2.id} from {p}: {msgs}")
# 8. Notification group consumption
print("\n--- Notification Group (independent progress) ---")
for p in notification_group.get_assigned_partitions(notif_consumer):
msgs = notif_consumer.poll(p, 10)
print(f" {notif_consumer.id} from {p}: {msgs}")
notification_group.commit_offset(p, notif_consumer.get_offset(p))
# 9. Re-poll (should be empty)
print("\n--- Re-poll notification group (should be empty) ---")
for p in notification_group.get_assigned_partitions(notif_consumer):
msgs = notif_consumer.poll(p, 10)
print(f" {notif_consumer.id} from {p}: {msgs}")
# 10. Replay via seek
print("\n--- Replay: seek notification consumer to offset 0 ---")
for p in notification_group.get_assigned_partitions(notif_consumer):
notif_consumer.seek(p, 0)
msgs = notif_consumer.poll(p, 10)
print(f" REPLAY {p}: {msgs}")
if __name__ == "__main__":
main()
#include <iostream>
#include <string>
#include "PubSubBroker.hpp"
#include "Producer.hpp"
int main() {
// 1. Create the broker
PubSubBroker broker;
// 2. Create a topic with 3 partitions
Topic* ordersTopic = broker.createTopic("orders", 3);
std::cout << "Created topic 'orders' with 3 partitions\n";
// 3. Create two consumer groups
ConsumerGroup* analyticsGroup = broker.createConsumerGroup("analytics-group");
ConsumerGroup* notifGroup = broker.createConsumerGroup("notification-group");
// 4. Create consumers and register them
Consumer analyticsConsumer1("analytics-1", "analytics-group");
Consumer analyticsConsumer2("analytics-2", "analytics-group");
Consumer notifConsumer("notif-1", "notification-group");
broker.registerConsumer("analytics-group", &analyticsConsumer1);
broker.registerConsumer("analytics-group", &analyticsConsumer2);
broker.registerConsumer("notification-group", ¬ifConsumer);
// 5. Subscribe groups to the topic
broker.subscribe("analytics-group", "orders");
broker.subscribe("notification-group", "orders");
// 6. Produce messages
Producer producer("order-service");
producer.publish(ordersTopic, "user-123", "Order #1001 placed");
producer.publish(ordersTopic, "user-456", "Order #1002 placed");
producer.publish(ordersTopic, "user-123", "Order #1001 shipped");
producer.publish(ordersTopic, "user-789", "Order #1003 placed");
producer.publish(ordersTopic, "", "Broadcast: flash sale started");
std::cout << "\nPublished 5 messages to 'orders' topic\n";
// 7. Analytics group - each consumer reads assigned partitions
std::cout << "\n--- Analytics Group ---\n";
for (auto* p : analyticsGroup->getAssignedPartitions(&analyticsConsumer1)) {
auto msgs = analyticsConsumer1.poll(p, 10);
std::cout << " " << analyticsConsumer1.getId() << " from partition-"
<< p->getId() << ": ";
for (const auto& m : msgs) std::cout << m.toString() << " ";
std::cout << "\n";
}
for (auto* p : analyticsGroup->getAssignedPartitions(&analyticsConsumer2)) {
auto msgs = analyticsConsumer2.poll(p, 10);
std::cout << " " << analyticsConsumer2.getId() << " from partition-"
<< p->getId() << ": ";
for (const auto& m : msgs) std::cout << m.toString() << " ";
std::cout << "\n";
}
// 8. Notification group
std::cout << "\n--- Notification Group ---\n";
for (auto* p : notifGroup->getAssignedPartitions(¬ifConsumer)) {
auto msgs = notifConsumer.poll(p, 10);
std::cout << " " << notifConsumer.getId() << " from partition-"
<< p->getId() << ": ";
for (const auto& m : msgs) std::cout << m.toString() << " ";
std::cout << "\n";
notifGroup->commitOffset(p, notifConsumer.getOffset(p));
}
// 9. Re-poll (empty)
std::cout << "\n--- Re-poll (should be empty) ---\n";
for (auto* p : notifGroup->getAssignedPartitions(¬ifConsumer)) {
auto msgs = notifConsumer.poll(p, 10);
std::cout << " " << notifConsumer.getId() << " from partition-"
<< p->getId() << ": " << msgs.size() << " messages\n";
}
// 10. Replay
std::cout << "\n--- Replay via seek ---\n";
for (auto* p : notifGroup->getAssignedPartitions(¬ifConsumer)) {
notifConsumer.seek(p, 0);
auto msgs = notifConsumer.poll(p, 10);
std::cout << " REPLAY partition-" << p->getId() << ": ";
for (const auto& m : msgs) std::cout << m.toString() << " ";
std::cout << "\n";
}
return 0;
}
Related Concepts
Scale this design past a single process and these are the concepts it runs into:
- Message Queues β β Kafka, SQS and RabbitMQ: this design is the in-memory version of a real broker
- Dead Letter Queue β β where a message lands when a consumer can never successfully process it
- Idempotency β β at-least-once delivery guarantees duplicates, so consumers have to dedupe
- Consistent Hashing β β key-based partition routing, and what rebalancing costs when partitions or consumers change
- Leader Election β β a real broker elects a leader per partition so there is one authority on write ordering
- Batch vs Stream β β committed offsets and replay are what let one log feed both a streaming job and a batch backfill
Discussion
Newest first