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.

August 16, 2026·12 min read

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#

EntityResponsibility
MessageImmutable value object wrapping the message content string
PublisherInterfaceContract: one method publish(topicName, message)
PublisherHolds a name and a reference to TopicRegistry; delegates publish calls to the registry
SubscriberInterfaceContract: update(topicName, message), subscribe(registry, topic), unsubscribe(registry, topic)
SubscriberReceives messages via update(); manages its own subscriptions through the registry
TopicMaintains a thread-safe set of subscribers; calls update() on each during broadcast()
TopicRegistrySingleton — 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.

java
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().

java
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.

java
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.

java
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.

java
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#

java
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.