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.
// 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.
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.
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 AriaHow would you implement message persistence so subscribers that are offline receive messages on reconnect?
What is the difference between a Pub-Sub system and a Message Queue? When would you use each?
How would you implement exactly-once delivery semantics in this design?
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.