Home/Learn/Low Level Design/Design a Pub-Sub System

Design a Pub-Sub System

Advanced
LLD Interview Problems
Java Source Code

Build an in-process publish-subscribe system with topics, async delivery via ExecutorService, subscription management, and configurable delivery guarantees.

Overview

A Pub-Sub System decouples Publishers from Subscribers through Topics. Publishers emit Messages to a Topic without knowing who is subscribed. The EventBus routes each Message to all active Subscribers for that Topic and delivers asynchronously using a thread pool. Subscription represents the link between a Subscriber and a Topic; cancelling it stops further delivery. DeliveryGuarantee enum models AT_MOST_ONCE, AT_LEAST_ONCE, and EXACTLY_ONCE semantics. This design is asked at companies building internal event buses, microservice communication layers, or platform-level messaging abstractions.

Requirements Analysis

Functional: create topics, subscribe to a topic, publish a message, receive messages asynchronously, unsubscribe, support multiple topics with independent subscriber lists. Non-functional: decoupled publisher and subscriber (neither holds a reference to the other), async delivery via thread pool, subscription cancellation thread-safe.

Requirements
// Entities : PubSubSystem, Topic, Message, Publisher, Subscriber, Subscription, EventBus
// Patterns : Observer (pub-sub), Command (messages), Strategy (delivery guarantee)

Core Classes & Relationships

Message holds payload, topic name, and timestamp. Subscriber is a FunctionalInterface with onMessage(Message). Subscription links a subscriber to a topic and has a cancel() method. DeliveryGuarantee enum: AT_MOST_ONCE, AT_LEAST_ONCE, EXACTLY_ONCE. Topic manages its subscriber list. EventBus is the central registry of topics.

Java — enums & interfaces
import java.util.concurrent.*;
import java.util.concurrent.CopyOnWriteArrayList;

public class Message {
    private final String id;
    private final String topicName;
    private final Object payload;
    private final long timestamp;

    public Message(String topicName, Object payload) {
        this.id = java.util.UUID.randomUUID().toString();
        this.topicName = topicName; this.payload = payload;
        this.timestamp = System.currentTimeMillis();
    }
    public String getTopicName() { return topicName; }
    public Object getPayload()   { return payload; }
    public String getId()        { return id; }
}

@FunctionalInterface
public interface Subscriber {
    void onMessage(Message message);
}

public enum DeliveryGuarantee { AT_MOST_ONCE, AT_LEAST_ONCE, EXACTLY_ONCE }

public class Subscription {
    private final String subscriptionId;
    private final Topic topic;
    private final Subscriber subscriber;
    private volatile boolean active = true;

    public Subscription(Topic topic, Subscriber subscriber) {
        this.subscriptionId = java.util.UUID.randomUUID().toString();
        this.topic = topic; this.subscriber = subscriber;
    }
    public void cancel()            { active = false; topic.removeSubscriber(subscriber); }
    public boolean isActive()       { return active; }
    public Subscriber getSubscriber(){ return subscriber; }
    public String getId()           { return subscriptionId; }
}

Java Implementation

Topic uses CopyOnWriteArrayList for thread-safe subscriber management without locking on publish. EventBus manages all topics and owns a shared ExecutorService for async delivery. Publisher.publish() calls EventBus.publish(). Each subscriber delivery is submitted as a Callable so failures in one subscriber do not affect others.

Java — core classes
public class Topic {
    private final String name;
    private final CopyOnWriteArrayList<Subscriber> subscribers = new CopyOnWriteArrayList<>();

    public Topic(String name) { this.name = name; }
    public String getName()   { return name; }

    public Subscription subscribe(Subscriber subscriber) {
        subscribers.add(subscriber);
        return new Subscription(this, subscriber);
    }

    public void removeSubscriber(Subscriber subscriber) {
        subscribers.remove(subscriber);
    }

    public CopyOnWriteArrayList<Subscriber> getSubscribers() { return subscribers; }
}

public class EventBus {
    private final ConcurrentHashMap<String, Topic> topics = new ConcurrentHashMap<>();
    private final ExecutorService executor;

    public EventBus(int threadPoolSize) {
        this.executor = Executors.newFixedThreadPool(threadPoolSize);
    }

    public Topic getOrCreateTopic(String name) {
        return topics.computeIfAbsent(name, Topic::new);
    }

    public Subscription subscribe(String topicName, Subscriber subscriber) {
        Topic topic = getOrCreateTopic(topicName);
        return topic.subscribe(subscriber);
    }

    public void publish(String topicName, Object payload) {
        Topic topic = topics.get(topicName);
        if (topic == null) return;
        Message message = new Message(topicName, payload);
        for (Subscriber subscriber : topic.getSubscribers()) {
            executor.submit(() -> {
                try {
                    subscriber.onMessage(message);
                } catch (Exception e) {
                    System.err.println("Subscriber error for topic " + topicName + ": " + e.getMessage());
                }
            });
        }
    }

    public void shutdown() { executor.shutdown(); }
}

public class Publisher {
    private final EventBus eventBus;
    private final String sourceName;

    public Publisher(EventBus eventBus, String sourceName) {
        this.eventBus = eventBus; this.sourceName = sourceName;
    }

    public void publish(String topicName, Object payload) {
        System.out.println("[" + sourceName + "] Publishing to " + topicName);
        eventBus.publish(topicName, payload);
    }
}

// Usage
EventBus bus = new EventBus(4);
Subscription s1 = bus.subscribe("orders", msg -> System.out.println("Email service: " + msg.getPayload()));
Subscription s2 = bus.subscribe("orders", msg -> System.out.println("Inventory: " + msg.getPayload()));

Publisher orderService = new Publisher(bus, "OrderService");
orderService.publish("orders", "Order #1234 placed");

s1.cancel(); // email unsubscribes
orderService.publish("orders", "Order #1235 placed"); // only inventory gets this
bus.shutdown();

Key Points to Remember

  • 1CopyOnWriteArrayList for subscribers allows safe iteration during publish while cancel() modifies the list concurrently.
  • 2ExecutorService for async delivery isolates subscriber failures — a slow subscriber does not block the publisher or other subscribers.
  • 3Subscription.cancel() removes the subscriber from the topic atomically; volatile active flag prevents delivery after cancellation.
  • 4AT_LEAST_ONCE delivery requires retry on subscriber failure; EXACTLY_ONCE requires distributed coordination (idempotency keys, transactions).

Interview Questions

Sign in to ask Aria
1

How would you implement message persistence so subscribers that are offline receive messages on reconnect?

HardFlipkart
2

What is the difference between a Pub-Sub system and a Message Queue? When would you use each?

MediumAmazon
3

How would you implement exactly-once delivery semantics in this design?

HardGoogle

Ask Aria about Design a Pub-Sub System

Your personal AI tutor — ask anything about this concept

Revision Status

Personal Notes

Sign in to save personal notes for this topic.

Discussion

Sign in to join the discussion.

Loading discussion…