Fanning Out Events with a Go Hub Goroutine Permalink to this section
Part of Go Streaming Patterns, under Backend Stream Generation & Connection Management.
Every Go SSE server eventually needs the same component: something that knows which streams are open and delivers each event to the right ones. The idiomatic design is a hub goroutine that owns the subscriber set outright, receiving register, unregister and publish requests over channels. Nothing else touches the set, so there are no locks to get wrong. This guide builds that hub, makes broadcast non-blocking so one slow client cannot stall the others, and extends it to per-topic routing.
Symptom & Developer Intent Permalink to this section
Hubs written without a clear ownership model show familiar problems:
fatal error: concurrent map writesunder load, or a data race reported bygo test -race.- A single stalled client freezes delivery to every other client.
- Goroutines leak:
runtime.NumGoroutine()climbs with every connection and never falls. - A panic “send on closed channel” when a client disconnects during a broadcast.
- Broadcast latency grows linearly with the number of subscribers.
The intent is a hub where delivery to fast clients is unaffected by slow ones, every subscriber’s resources are released on disconnect, and the code is provably race-free.
Root Cause Analysis Permalink to this section
A shared map guarded by a mutex works until someone performs a blocking channel send while holding the lock; then one slow client holds the lock and every publish waits. Closing a client’s channel from the handler while the hub may be sending to it causes the “send on closed channel” panic. Both come from shared ownership of the same state.
The fix is to give ownership to one goroutine. The hub is the only reader and writer of the subscriber map and the only party that closes subscriber channels. Handlers ask it to register and unregister; publishers ask it to broadcast. Broadcast uses a non-blocking send into each subscriber’s buffered channel, and a subscriber whose buffer is full is evicted — its channel closed by the hub, its handler returning, its client reconnecting and replaying.
Step-by-Step Resolution Permalink to this section
Step 1 — Define the hub and its messages Permalink to this section
package hub
type Event struct {
ID string
Topic string
Frame []byte // pre-serialised SSE frame: "id: …\nevent: …\ndata: …\n\n"
}
type Sub struct {
Topic string
C chan []byte
}
type Hub struct {
register chan *Sub
unregister chan *Sub
publish chan Event
subs map[string]map[*Sub]struct{} // topic → set, owned by run()
}
func New() *Hub {
h := &Hub{
register: make(chan *Sub),
unregister: make(chan *Sub),
publish: make(chan Event, 1024),
subs: make(map[string]map[*Sub]struct{}),
}
go h.run()
return h
}
Step 2 — Run the loop that owns everything Permalink to this section
func (h *Hub) run() {
for {
select {
case s := <-h.register:
if h.subs[s.Topic] == nil {
h.subs[s.Topic] = make(map[*Sub]struct{})
}
h.subs[s.Topic][s] = struct{}{}
case s := <-h.unregister:
h.remove(s)
case e := <-h.publish:
for s := range h.subs[e.Topic] {
select {
case s.C <- e.Frame: // delivered into the client's buffer
default: // buffer full: this client is too slow
h.remove(s)
}
}
}
}
}
func (h *Hub) remove(s *Sub) {
if set, ok := h.subs[s.Topic]; ok {
if _, ok := set[s]; ok {
delete(set, s)
close(s.C) // only the hub closes, so no send-on-closed panic
if len(set) == 0 {
delete(h.subs, s.Topic)
}
}
}
}
func (h *Hub) Publish(e Event) { h.publish <- e }
func (h *Hub) Subscribe(topic string) *Sub {
s := &Sub{Topic: topic, C: make(chan []byte, 64)}
h.register <- s
return s
}
func (h *Hub) Unsubscribe(s *Sub) { h.unregister <- s }
remove is idempotent, which matters: a slow client can be evicted by a broadcast and then unregistered by its handler moments later.
Step 3 — Serve a stream from the hub Permalink to this section
func (a *App) stream(w http.ResponseWriter, r *http.Request) {
flusher, ok := w.(http.Flusher)
if !ok {
http.Error(w, "streaming unsupported", http.StatusInternalServerError)
return
}
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("Cache-Control", "no-cache")
w.Header().Set("X-Accel-Buffering", "no")
sub := a.hub.Subscribe("user:" + userID(r))
defer a.hub.Unsubscribe(sub) // runs on disconnect or eviction
hb := time.NewTicker(15 * time.Second)
defer hb.Stop()
for {
select {
case frame, ok := <-sub.C:
if !ok {
return // evicted by the hub: the client will reconnect and replay
}
if _, err := w.Write(frame); err != nil {
return
}
flusher.Flush()
case <-hb.C:
if _, err := io.WriteString(w, ": hb\n\n"); err != nil {
return
}
flusher.Flush()
case <-r.Context().Done():
return // client disconnected
}
}
}
The handler owns its goroutine and returns on any of three conditions; the deferred Unsubscribe then asks the hub to release the subscription. The basic handler shape is covered in implementing SSE with Go channels and Flusher.
Step 4 — Serialise once, before publishing Permalink to this section
Build the frame bytes once per event and share the slice with every subscriber. Formatting inside the handler would repeat the work per client:
func frame(id, event string, v any) []byte {
b, _ := json.Marshal(v)
return []byte(fmt.Sprintf("id: %s\nevent: %s\ndata: %s\n\n", id, event, b))
}
h.Publish(hub.Event{ID: id, Topic: "user:42", Frame: frame(id, "notification", n)})
The byte slice is never mutated after publishing, so sharing it across goroutines is safe.
Step 5 — Scale the hub Permalink to this section
One hub goroutine handles hundreds of thousands of sends per second, because each non-blocking send is a few dozen nanoseconds. If profiling shows the hub loop saturated, shard: run N hubs and route each topic to hubs[hash(topic)%N]. Topics stay single-owner, and throughput scales with cores.
Validation & Monitoring Permalink to this section
go test -race ./hub/... # the ownership model should make this trivially clean
# Goroutines must return to baseline after clients leave.
for i in $(seq 1 1000); do curl -sN localhost:8080/events > /dev/null & done; sleep 2
curl -s localhost:8080/debug/vars | jq .goroutines
kill $(jobs -p); sleep 2
curl -s localhost:8080/debug/vars | jq .goroutines
Export three expvar or Prometheus metrics from the hub: subscribers per topic (summed), evictions per minute, and the length of the publish channel. A publish channel that stays near its capacity means the hub loop is the bottleneck and should be sharded.
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
Why not a sync.Map or RWMutex instead of a hub goroutine?
They can work, but only if no blocking operation ever happens under the lock and channel closing is carefully coordinated. The hub goroutine makes those guarantees structural rather than a matter of discipline.
What buffer size should subscriber channels have?
Enough to absorb normal bursts — 32 to 256 frames is typical. The buffer is not a queue for slow clients; it is a shock absorber. A client that fills it is evicted and recovers through replay.
Does eviction lose events?
The evicted client misses the live frames sent after eviction, but its reconnect sends Last-Event-ID and the server replays from the store. For state-shaped data, send a fresh snapshot on reconnect instead.
How do I publish from another process?
Run a subscriber goroutine per process that reads from Redis, NATS or Kafka and calls hub.Publish. The hub itself stays in-process.