Go SSE服务器推送:EventSource实现
摘要: 本篇讲解Go语言SSE服务器推送实现,基于text/event-stream协议格式编写EventSource服务端,使用http.Flusher实时刷新事件,实现多客户端广播与心跳保活,对比SSE与WebSocket的技术差异和选型建议,分享Nginx反向代理缓冲导致事件延迟到达的踩坑经验。
开篇故事
我们有个运维告警系统,用户在页面上实时看告警推送。最初用WebSocket做的,功能正常但代码量大,心跳重连连接管理一堆事。后来评估发现告警推送是单向的,服务端推客户端收,客户端不需要往服务端发消息。这种场景用SSE更合适,协议简单,浏览器原生支持EventSource,服务端几十行代码搞定。
迁移到SSE后开发环境一切正常,事件秒到。上了生产环境用户反馈告警有时延迟30秒到1分钟才显示。排查发现是Nginx反向代理开启了缓冲,事件攒了一批再转发,实时性全丢了。关掉代理缓冲后事件恢复秒到。
这个坑让我把SSE从协议层到部署层完整梳理了一遍。
一、text/event-stream协议格式
SSE的协议非常简单。HTTP响应头Content-Type设为text/event-stream,响应体每条消息用固定格式组织,用空行分隔。客户端用EventSource对象接收。
// 单条事件的格式 data: {"alert":"CPU超过90%"}\n\n // 带事件类型和ID event: alert\n data: {"level":"critical","msg":"内存不足"}\n id: 12345\n\n // 注释行(用于心跳保活) : heartbeat\n\n每行用\n结尾,一条事件用\n\n空行结束。data:后面跟数据,event:指定事件类型,id:设置事件ID用于断线续传。以:开头的是注释行,浏览器不处理但能保持连接活跃。
二、服务端实现与Flusher刷新
Go实现SSE服务端的关键是http.Flusher。HTTP默认会把响应缓冲起来攒够一批再发送,SSE要求每条事件立即推给客户端,必须手动调用Flush。
packagemainimport("encoding/json""fmt""log""net/http""time")// AlertEvent 告警事件结构typeAlertEventstruct{Levelstring`json:"level"`// 告警级别: info/warning/criticalMessagestring`json:"message"`// 告警内容Timeint64`json:"time"`// 时间戳}funcmain(){http.HandleFunc("/events",sseHandler)log.Println("SSE服务启动在 :8080")log.Fatal(http.ListenAndServe(":8080",nil))}// sseHandler 处理SSE连接funcsseHandler(w http.ResponseWriter,r*http.Request){// 检查是否支持Flush// SSE的核心: 每条事件必须立即刷新到客户端flusher,ok:=w.(http.Flusher)if!ok{http.Error(w,"不支持流式推送",http.StatusInternalServerError)return}// 设置SSE响应头w.Header().Set("Content-Type","text/event-stream")w.Header().Set("Cache-Control","no-cache")// 禁用缓存w.Header().Set("Connection","keep-alive")// 保持连接w.Header().Set("Access-Control-Allow-Origin","*")// 跨域支持// 模拟每2秒推送一条告警ticker:=time.NewTicker(2*time.Second)deferticker.Stop()// 心跳定时器,每15秒发一条注释保活// 防止代理或浏览器因空闲超时断开连接heartbeat:=time.NewTicker(15*time.Second)deferheartbeat.Stop()ctx:=r.Context()// 请求上下文,客户端断开时ctx.Done()for{select{case<-ctx.Done():// 客户端断开连接log.Println("客户端断开")returncase<-heartbeat.C:// 发送心跳注释,保持连接活跃// 注释行不会被EventSource解析为事件fmt.Fprintf(w,": heartbeat\n\n")flusher.Flush()// 立即推送case<-ticker.C:// 构造告警事件event:=AlertEvent{Level:"warning",Message:fmt.Sprintf("告警 #%d: CPU使用率超过80%%",time.Now().Unix()),Time:time.Now().Unix(),}// 序列化为JSONdata,_:=json.Marshal(event)// 写入SSE格式: data: {json}\n\nfmt.Fprintf(w,"data: %s\n\n",data)flusher.Flush()// 关键: 立即刷新到网络}}}flusher.Flush()是SSE的灵魂。不调Flush,数据在Go的响应缓冲区里待着,客户端收不到。调了Flush,Go底层调用http.ResponseWriter的Flush方法,把缓冲区的数据推到TCP连接,客户端立即收到。
三、多客户端广播
单个客户端的推送只是演示。实际场景多个浏览器同时订阅告警,一个告警产生后要广播给所有在线客户端。用一个Hub管理所有连接。
packagesseimport("fmt""net/http""sync""time")// SSEHub 管理所有SSE客户端连接和消息广播typeSSEHubstruct{mu sync.RWMutex// 读写锁保护clientsclientsmap[chanstring]struct{}// 客户端通道集合}// NewSSEHub 创建广播HubfuncNewSSEHub()*SSEHub{return&SSEHub{clients:make(map[chanstring]struct{}),}}// Register 注册新客户端// 每个客户端分配一个带缓冲的通道// 缓冲大小决定客户端落后多少条消息不被丢弃func(h*SSEHub)Register()chanstring{ch:=make(chanstring,64)// 缓冲64条事件h.mu.Lock()h.clients[ch]=struct{}{}h.mu.Unlock()returnch}// Unregister 注销客户端func(h*SSEHub)Unregister(chchanstring){h.mu.Lock()delete(h.clients,ch)h.mu.Unlock()close(ch)// 关闭通道,通知客户端goroutine退出}// Broadcast 广播事件给所有客户端// 某个客户端通道满了就跳过,不影响其他客户端func(h*SSEHub)Broadcast(eventstring){h.mu.RLock()deferh.mu.RUnlock()forch:=rangeh.clients{// 非阻塞写入,通道满了就丢弃// 避免一个慢客户端阻塞整个广播select{casech<-event:default:// 通道满,该客户端处理不过来// 生产环境应记录日志并考虑踢掉}}}// HandleSSE 处理单个SSE连接// 注册到Hub,循环推送事件直到客户端断开func(h*SSEHub)HandleSSE(w http.ResponseWriter,r*http.Request){flusher,ok:=w.(http.Flusher)if!ok{http.Error(w,"不支持流式推送",http.StatusInternalServerError)return}w.Header().Set("Content-Type","text/event-stream")w.Header().Set("Cache-Control","no-cache")w.Header().Set("Connection","keep-alive")// 注册到Hub,获取事件通道ch:=h.Register()deferh.Unregister(ch)ctx:=r.Context()// 心跳定时器heartbeat:=time.NewTicker(15*time.Second)deferheartbeat.Stop()for{select{case<-ctx.Done():return// 客户端断开case<-heartbeat.C:fmt.Fprintf(w,": heartbeat\n\n")flusher.Flush()caseevent,ok:=<-ch:if!ok{return// Hub关闭了通道}// 推送事件给客户端fmt.Fprintf(w,"data: %s\n\n",event)flusher.Flush()}}}Hub模式和WebSocket的Hub思路一样,用channel解耦生产者和消费者。区别是SSE不需要WebSocket握手,直接用HTTP响应流推送。每个连接占一个goroutine,一个读r.Context()检测断开,一个读channel拿事件写响应。
四、踩坑经验:反向代理缓冲导致事件延迟
开篇说的生产环境事件延迟30秒到1分钟的坑,根因是Nginx的proxy_buffering默认开启。
Nginx作为反向代理时,默认会把后端响应缓冲到本地磁盘,攒够一批或上游响应结束后才转发给客户端。这个策略对普通HTTP请求是合理的,减少网络往返提高吞吐。但SSE要求实时推送,缓冲直接破坏了实时性。
修复方法有两个层面。
第一个层面是Nginx配置,对SSE路由关闭缓冲。
# nginx.conf - SSE路由关闭缓冲 location /events { proxy_pass http://backend; # 关闭缓冲,后端数据立即转发 proxy_buffering off; proxy_cache off; # 关闭TCP缓冲,禁用Nagle算法 # 让小数据包立即发送 tcp_nodelay on; # SSE是长连接,超时设长一点 proxy_read_timeout 3600s; # 传递真实Host和客户端IP proxy_set_header Host $host; proxy_set_header X-Real-IP $remote_addr; }第二个层面是Go服务端设置X-Accel-Buffering响应头,告诉Nginx不要缓冲这个响应。
// 在SSE handler中设置响应头func(h*SSEHub)HandleSSE(w http.ResponseWriter,r*http.Request){// 关键: 告诉Nginx不要缓冲此响应// 即使全局proxy_buffering on,这个头也能让Nginx关闭单条响应的缓冲w.Header().Set("X-Accel-Buffering","no")w.Header().Set("Content-Type","text/event-stream")w.Header().Set("Cache-Control","no-cache")w.Header().Set("Connection","keep-alive")// ... 后续处理逻辑同上}X-Accel-Buffering: no是Nginx的特殊响应头,Nginx看到这个头会对当前响应自动关闭缓冲。这是最省事的方案,不需要改Nginx配置,Go服务端自己控制。
五、对比分析
| 维度 | SSE | WebSocket | HTTP轮询 |
|---|---|---|---|
| 通信方向 | 服务端到客户端 | 双向 | 客户端发起 |
| 底层协议 | HTTP/1.1 | HTTP升级 | HTTP |
| 浏览器支持 | 原生EventSource | 原生WebSocket | 全兼容 |
| 自动重连 | 内置自动重连 | 需自己实现 | 无连接 |
| 断线续传 | Last-Event-ID支持 | 需自己实现 | 不适用 |
| 代理兼容 | HTTP友好 | 需特殊配置 | HTTP友好 |
| 连接数限制 | 6个/域名 | 无限制 | 每次新建 |
| 消息格式 | 文本 | 文本/二进制 | 任意 |
SSE最大的优势是简单。浏览器EventSource原生支持自动重连和断线续传,服务端代码量不到WebSocket的一半。HTTP协议天然友好,经过任何代理都不会被拦截。缺点是单向通信,客户端发不了消息,得靠额外的HTTP请求。浏览器对同域名SSE连接数限制6个,大量并发推送场景WebSocket更合适。
选型经验,单向推送选SSE,聊天协作选WebSocket。告警通知、日志流、行情推送这类服务端到客户端的场景,SSE是最轻量的方案。
总结
SSE是单向推送场景的最优解。协议简单,浏览器原生支持自动重连,服务端用http.Flusher实时刷新就行。多客户端用Hub模式管理,广播用channel解耦。部署时注意代理缓冲问题,设X-Accel-Buffering: no响应头让Nginx别缓冲。下一篇聊gRPC流式进阶,重点讲流控和背压。