Low Level Design
Design a Pub-Sub Model
Build a Publisher-Subscriber messaging system from scratch — decoupled components, dynamic topic management, and thread-safe broadcast using the Observer pattern and Singleton registry.
Problem Description#
Design a Publish-Subscribe (Pub-Sub) messaging system that allows producers to send messages to named topics and consumers to receive those messages asynchronously.
The Pub-Sub model is the backbone of modern event-driven architectures — used in systems like Kafka, Google Pub/Sub, and Redis. The key challenge is decoupling: a publisher should know nothing about who is listening, and a subscriber should know nothing about who is producing. A central registry manages this coordination.
This is a classic LLD problem because it naturally exercises the Observer pattern, the Singleton pattern, interface design, and thread-safety — all within a small, cohesive system.
Clarify Requirements#
Before designing, ask these questions in an interview:
Functional
- Should publishers create topics, or should that be a separate admin operation?
- Can a subscriber listen to multiple topics simultaneously?
- Should a subscriber be able to unsubscribe at runtime?
- When a message is published, should all current subscribers receive it synchronously?
- Should the system guard against duplicate subscriptions to the same topic?
Non-functional
- Do we need thread safety for concurrent publish/subscribe operations?
- Should messages be persisted or delivered in-memory only?
- Do we need delivery guarantees (at-most-once vs at-least-once)?
- Is horizontal scalability (multiple registry nodes) in scope?
Final Requirements#
After clarification, here's what we'll build:
- A TopicRegistry Singleton manages all topics centrally — one shared instance for the entire application
- Publishers publish messages to a named topic via the registry; they never talk to subscribers directly
- Subscribers can subscribe() and unsubscribe() from any topic at runtime
- Topic maintains a thread-safe Set<Subscriber> and broadcasts to all of them when a message arrives
- Messages are delivered synchronously to all current subscribers of a topic
- Duplicate subscriptions to the same topic are silently ignored (set semantics)
- Thread safety is achieved via ConcurrentHashMap in both Topic and TopicRegistry
Core Entities#
| Entity | Responsibility |
|---|---|
| Message | Immutable value object wrapping the message content string |
| PublisherInterface | Contract: one method publish(topicName, message) |
| Publisher | Holds a name and a reference to TopicRegistry; delegates publish calls to the registry |
| SubscriberInterface | Contract: update(topicName, message), subscribe(registry, topic), unsubscribe(registry, topic) |
| Subscriber | Receives messages via update(); manages its own subscriptions through the registry |
| Topic | Maintains a thread-safe set of subscribers; calls update() on each during broadcast() |
| TopicRegistry | Singleton — owns the Map<String, Topic>; exposes create, subscribe, unsubscribe, publish, and display operations |
Patterns Used#
1. Observer — the core delivery mechanism#
Topic is the Subject and Subscriber is the Observer. When broadcast(message) is called on a Topic, it iterates its subscriber set and calls update() on each — the classic one-to-many notification. Publishers never hold references to subscribers; the topic is the only intermediary.
This means adding a new subscriber type (e.g., a logging subscriber, a metrics collector) requires zero changes to Topic or Publisher.
2. Singleton — TopicRegistry#
There must be exactly one TopicRegistry so all publishers and subscribers agree on the same topic map. The registry uses lazy initialization guarded by synchronized:
TopicRegistry.getInstance()
└─ creates the single instance on first call, thread-safely
All topic operations (createTopic, subscribe, unsubscribe, publish) go through this single instance.
3. Interface Segregation#
PublisherInterface and SubscriberInterface define separate, minimal contracts. A class that only needs to publish depends solely on PublisherInterface; one that only subscribes depends solely on SubscriberInterface. Neither interface bleeds responsibilities from the other.
4. Encapsulation — subscribers own their subscriptions#
Rather than requiring a central manager to wire up subscribers, each Subscriber instance calls subscribe() and unsubscribe() on itself, passing the registry and topic name. This keeps subscription logic co-located with the subscriber and makes each Subscriber self-contained.
Code#
Core Value Object#
Message is an immutable wrapper around the content string — passed between publishers, the registry, topics, and subscribers without mutation.
public class Message {
private String content;
public Message(String content) {
this.content = content;
}
public String getContent() {
return content;
}
}
Publisher#
Publisher holds a name and a reference to the shared registry. It prints what it's publishing and delegates immediately to topicRegistry.publish().
interface PublisherInterface {
void publish(String topicName, Message message);
}
public class Publisher implements PublisherInterface {
private String publisherName;
private TopicRegistry topicRegistry;
public Publisher(String name, TopicRegistry topicRegistry) {
this.publisherName = name;
this.topicRegistry = topicRegistry;
}
public void publish(String topicName, Message message) {
System.out.println("[Publisher " + publisherName + "] Publishing to '" + topicName + "': " + message.getContent());
topicRegistry.publish(topicName, message);
}
}
Subscriber#
Subscriber implements the update callback and manages its own lifecycle — calling subscribe / unsubscribe on the registry on behalf of itself.
interface SubscriberInterface {
void update(String topicName, Message message);
void subscribe(TopicRegistry topicRegistry, String topicName);
void unsubscribe(TopicRegistry topicRegistry, String topicName);
}
public class Subscriber implements SubscriberInterface {
private final String subscriberName;
public Subscriber(String name) {
this.subscriberName = name;
}
public String getName() {
return subscriberName;
}
public void update(String topicName, Message message) {
System.out.println("[" + subscriberName + "] received message on '" + topicName + "': " + message.getContent());
}
public void subscribe(TopicRegistry topicRegistry, String topicName) {
System.out.println("[" + subscriberName + "] subscribed '" + topicName + "'");
topicRegistry.subscribe(topicName, this);
}
public void unsubscribe(TopicRegistry topicRegistry, String topicName) {
System.out.println("[" + subscriberName + "] unsubscribed '" + topicName + "'");
topicRegistry.unsubscribe(topicName, this);
}
}
Topic#
Topic maintains a thread-safe subscriber set backed by ConcurrentHashMap.newKeySet(). broadcast() fans the message out to every current subscriber.
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
public class Topic {
private String topicName;
private final Set<Subscriber> subscribers = ConcurrentHashMap.newKeySet();
public Topic(String name) {
topicName = name;
}
public void subscribe(Subscriber subscriber) {
subscribers.add(subscriber);
}
public void unsubscribe(Subscriber subscriber) {
subscribers.remove(subscriber);
}
public void broadcast(Message message) {
for (Subscriber subscriber : subscribers) {
subscriber.update(topicName, message);
}
}
}
TopicRegistry — Singleton#
The central coordinator. Uses ConcurrentHashMap for thread-safe topic management. computeIfAbsent auto-creates a topic on first subscribe if it doesn't exist yet.
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
public class TopicRegistry {
private static TopicRegistry instance;
private Map<String, Topic> topics = new ConcurrentHashMap<>();
public static synchronized TopicRegistry getInstance() {
if (instance == null)
instance = new TopicRegistry();
return instance;
}
public void createTopic(String topicName) {
System.out.println("Topic Created: [" + topicName + "]");
topics.putIfAbsent(topicName, new Topic(topicName));
}
public void subscribe(String topicName, Subscriber subscriber) {
topics.computeIfAbsent(topicName, Topic::new).subscribe(subscriber);
}
public void unsubscribe(String topicName, Subscriber subscriber) {
if (topics.containsKey(topicName)) {
topics.get(topicName).unsubscribe(subscriber);
}
}
public void publish(String topicName, Message message) {
if (topics.containsKey(topicName)) {
topics.get(topicName).broadcast(message);
}
}
public void displayAllTopics() {
System.out.println("Available Topics: " + topics.keySet());
}
}
Demo#
public class pubSubDemo {
public static void main(String[] args) {
TopicRegistry topicRegistry = TopicRegistry.getInstance();
topicRegistry.createTopic("Computers");
topicRegistry.createTopic("Biology");
// Create subscribers
Subscriber subscriber1 = new Subscriber("subscriber1");
Subscriber subscriber2 = new Subscriber("subscriber2");
Subscriber subscriber3 = new Subscriber("subscriber3");
subscriber1.subscribe(topicRegistry, "Computers");
subscriber2.subscribe(topicRegistry, "Biology");
subscriber2.subscribe(topicRegistry, "Computers");
subscriber3.subscribe(topicRegistry, "Computers");
// Publishers broadcast messages
Publisher publisher1 = new Publisher("DailyCode", topicRegistry);
publisher1.publish("Computers", new Message("New OOPs language is in town."));
Publisher publisher2 = new Publisher("DailyBio", topicRegistry);
publisher2.publish("Biology", new Message("New neurology study is in town."));
// subscriber2 leaves Computers; next publish only reaches subscriber1 and subscriber3
subscriber2.unsubscribe(topicRegistry, "Computers");
publisher1.publish("Computers", new Message("New ReactJs version available. Learn today!"));
}
}
Class Diagram#
Extendible — Follow Ups#
1. Asynchronous message delivery#
Replace the synchronous broadcast() loop with an ExecutorService. Each subscriber's update() is submitted as a Callable — the publisher returns immediately and subscribers are notified concurrently. This is how real brokers like Kafka achieve high throughput without blocking producers.
2. Message filtering / content-based routing#
Add a Predicate<Message> to each subscription. Topic.broadcast() only calls update() when the predicate passes. Subscribers can express interest in subsets — e.g., only messages whose content contains a keyword — without modifying Topic or TopicRegistry.
3. Message persistence and replay#
Introduce a MessageStore inside Topic that appends every published message to a log. New subscribers can request replayFrom(offset) to receive historical messages — the foundation of Kafka's offset-based consumer model.
4. Dead-letter queue#
Wrap each subscriber.update() call in a try-catch. If delivery fails (exception thrown), enqueue the message and subscriber into a DeadLetterQueue. A background worker retries failed deliveries with exponential backoff, matching the reliability guarantees of production messaging systems.
5. Topic hierarchy and wildcards#
Extend topic names to support dot-separated paths (news.tech.java). Allow subscribers to use wildcard patterns (news.tech.* or news.#). TopicRegistry.publish() matches the published topic against all registered patterns — the approach MQTT uses for its topic tree.
6. Publisher and subscriber metrics#
Wrap TopicRegistry with a decorator that increments counters on every publish, subscribe, and unsubscribe call. Expose metrics via a MetricsReporter interface. Swap in a Prometheus-compatible reporter without touching any core logic.