Implement a pub/sub broadcaster where one slow subscriber cannot block the others.
Question 258HardGo 1.22 to 1.25
A channel delivers each value to one receiver, so broadcasting needs a channel per subscriber. The hard parts are slow consumers, unsubscribing without a send-on-closed-channel panic, and not leaking goroutines.
type Broker[T any] struct {
mu sync.Mutex
subs map[chan T]struct{}
}
func NewBroker[T any]() *Broker[T] {
return &Broker[T]{subs: make(map[chan T]struct{})}
}
// Subscribe returns a receive-only channel and an idempotent cancel func.
func (b *Broker[T]) Subscribe(buf int) (<-chan T, func()) {
ch := make(chan T, buf)
b.mu.Lock()
b.subs[ch] = struct{}{}
b.mu.Unlock()
var once sync.Once
return ch, func() {
once.Do(func() {
b.mu.Lock()
delete(b.subs, ch) // no Publish can see ch after this
b.mu.Unlock()
close(ch) // the broker is the only sender, so closing is safe
})
}
}
// Publish never blocks: a subscriber whose buffer is full misses the value.
func (b *Broker[T]) Publish(v T) (dropped int) {
b.mu.Lock()
defer b.mu.Unlock()
for ch := range b.subs {
select {
case ch <- v:
default:
dropped++
}
}
return dropped
}
Why it is correct:
- Sends happen only while holding
mu, and cancel removes the channel under the same lock before closing it, so nothing ever sends on a closed channel. - The non-blocking
selectwithdefaultmeans holding the lock while sending cannot stall publishers. - Subscribers get a
<-chan T, so they cannot close or send on it.
Design choices to discuss:
- Slow-consumer policy: drop newest (as above), drop oldest (a ring buffer per subscriber), block with a timeout, or disconnect the subscriber. Always export a dropped-messages metric.
- Latest-value-only subscribers, such as config updates, can use a buffer of 1 and replace the pending value.
- Tie cancel to the subscriber's context, e.g.
stop := context.AfterFunc(ctx, cancel), so an abandoned subscriber cannot leak. - With many subscribers, publishing under one mutex becomes the bottleneck. Copy-on-write subscriber lists in an
atomic.Pointerlet publishes run without the lock, but cancel then needs its own coordination with in-flight sends.
More on Concurrency Patterns & sync
- Q256What is false sharing, how does it hurt concurrent Go code, and how do you fix it?
- Q257How is GOMAXPROCS chosen in containers, and what changed in Go 1.25?