Go

Design an in-process pub/sub broker. How do you handle slow subscribers?

Question 471HardGo 1.22 to 1.25

Requirements: topics, concurrent publish, subscribe/unsubscribe at any time, and a publisher must never block on a slow consumer. Each subscriber gets a buffered channel; on full buffer choose a policy: drop (with metrics), disconnect the subscriber, or block with timeout. Unsubscribe must be idempotent and must not close a channel while a publisher may be sending (send on closed channel panics) — doing both under the same lock avoids this.

type Broker[T any] struct {
	mu     sync.RWMutex
	subs   map[string]map[*sub[T]]struct{}
	closed bool
}

type sub[T any] struct {
	ch      chan T
	once    sync.Once
	dropped atomic.Int64
}

func NewBroker[T any]() *Broker[T] {
	return &Broker[T]{subs: make(map[string]map[*sub[T]]struct{})}
}

// Subscribe returns a receive channel and an unsubscribe func.
func (b *Broker[T]) Subscribe(topic string, buf int) (<-chan T, func()) {
	s := &sub[T]{ch: make(chan T, buf)}
	b.mu.Lock()
	defer b.mu.Unlock()
	if b.closed {
		close(s.ch)
		return s.ch, func() {}
	}
	if b.subs[topic] == nil {
		b.subs[topic] = make(map[*sub[T]]struct{})
	}
	b.subs[topic][s] = struct{}{}
	return s.ch, func() {
		b.mu.Lock()
		defer b.mu.Unlock()
		delete(b.subs[topic], s)
		s.once.Do(func() { close(s.ch) })
	}
}

// Publish never blocks: full subscriber buffers drop the message.
func (b *Broker[T]) Publish(topic string, msg T) {
	b.mu.RLock()
	defer b.mu.RUnlock()
	for s := range b.subs[topic] {
		select {
		case s.ch <- msg:
		default:
			s.dropped.Add(1)
		}
	}
}

func (b *Broker[T]) Close() {
	b.mu.Lock()
	defer b.mu.Unlock()
	b.closed = true
	for _, set := range b.subs {
		for s := range set {
			s.once.Do(func() { close(s.ch) })
		}
	}
	clear(b.subs)
}

Why it is safe: sends happen under RLock, closes under Lock, so they never overlap. Discuss trade-offs: at-most-once delivery, no ordering across topics, and for durability/multi-process fan-out you'd move to NATS, Kafka or Redis Streams.

More on Standard Library, HTTP & Systems Design in Go

All 35 Standard Library, HTTP & Systems Design in Go questions