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
- Q469Design a concurrency-safe, generic LRU cache with TTL.
- Q470What is a cache stampede and how does singleflight solve it?
- Q472Fan out N HTTP calls with bounded concurrency, fail fast on first error, and collect results in order.
- Q473How do you unit-test HTTP handlers and HTTP clients with net/http/httptest?
- Q474Implement Server-Sent Events (SSE) in net/http. How do http.ResponseController, flushing and WriteTimeout interact?
- Q475Write an HTTP client retry helper with exponential backoff and jitter. What makes a request safe to retry?