尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Go SSE服务器推送:EventSource实现

Go SSE服务器推送:EventSource实现 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/jsonfmtlognet/httptime)// AlertEvent 告警事件结构typeAlertEventstruct{Levelstringjson:level// 告警级别: info/warning/criticalMessagestringjson:message// 告警内容Timeint64json: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的响应缓冲区里待着客户端收不到。调了FlushGo底层调用http.ResponseWriter的Flush方法把缓冲区的数据推到TCP连接客户端立即收到。三、多客户端广播单个客户端的推送只是演示。实际场景多个浏览器同时订阅告警一个告警产生后要广播给所有在线客户端。用一个Hub管理所有连接。packagesseimport(fmtnet/httpsynctime)// SSEHub 管理所有SSE客户端连接和消息广播typeSSEHubstruct{mu sync.RWMutex// 读写锁保护clientsclientsmap[chanstring]struct{}// 客户端通道集合}// NewSSEHub 创建广播HubfuncNewSSEHub()*SSEHub{returnSSEHub{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服务端自己控制。五、对比分析维度SSEWebSocketHTTP轮询通信方向服务端到客户端双向客户端发起底层协议HTTP/1.1HTTP升级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流式进阶重点讲流控和背压。
返回列表