Go语言SSE服务端推送事件与前端EventSource实时通知
导语
在实时Web应用中,服务器向客户端推送数据是一项基本需求。传统的实现方案包括:短轮询(Short Polling)、长轮询(Long Polling)、WebSocket和SSE(Server-Sent Events)。其中,SSE是一种基于HTTP的轻量级推送技术,具有实现简单、自动重连、支持自定义事件类型等优点。与WebSocket相比,SSE是单向通信(服务端→客户端),非常适合股票行情、新闻推送、日志流式输出等场景。本文将深入探讨Go语言中实现SSE的完整方案,并结合前端EventSource API实现实时通知功能。
核心技术知识点讲解
1. SSE协议基础
SSE是HTML5规范的一部分,允许服务端向客户端流式推送文本数据。
HTTP响应头要求:
- Content-Type: text/event-stream
- Cache-Control: no-cache
- Connection: keep-alive
消息格式:
event: 事件类型(可选)
data: 消息内容
id: 事件ID(可选)
retry: 重连延时(可选,单位毫秒)
(空行表示消息结束)
示例:
event: update
data: {"time": "2026-06-08 12:00:00", "value": 42}
event: ping
data: heartbeat
2. 前端EventSource API
浏览器原生支持SSE,通过EventSource接口实现:
const eventSource = new EventSource('/events');
// 监听默认消息(没有event字段的消息)
eventSource.onmessage = function(event) {
console.log('收到消息:', event.data);
};
// 监听自定义事件
eventSource.addEventListener('update', function(event) {
console.log('收到update事件:', event.data);
});
// 错误处理和自动重连
eventSource.onerror = function(error) {
console.log('连接错误,将自动重连');
};
优点:
- 自动重连机制(默认3秒)
- 支持CORS
- 轻量级,无需额外JS库
3. Go语言实现SSE的关键点
- 设置正确的响应头
- 禁用输出缓冲(使用Flusher接口)
- 保持连接不关闭
- 使用channel向客户端推送消息
- 检测客户端断开连接
4. SSE vs WebSocket vs 长轮询
| 通信方向 | 单向(服务端→客户端) | 双向 | 单向 |
| 协议 | HTTP | 独立协议(ws://) | HTTP |
| 浏览器支持 | 好(除IE) | 好 | 好 |
| 自动重连 | 是 | 需手动实现 | 需手动实现 |
| 复杂度 | 低 | 中 | 高 |
实战代码演示/项目案例总结
完整的SSE服务端实现(Go)
package main
import (
"encoding/json"
"fmt"
"log"
"net/http"
"os"
"os/signal"
"strings"
"sync"
"syscall"
"time"
)
// 消息结构
type SSEMessage struct {
Event string `json:"event,omitempty"`
Data interface{} `json:"data"`
ID string `json:"id,omitempty"`
Retry int `json:"retry,omitempty"`
}
// SSE客户端连接
type SSEClient struct {
ID string
Channel chan SSEMessage
Done chan struct{}
}
// SSE事件总线
type SSEBroker struct {
clients map[string]*SSEClient
subscribe chan *SSEClient
unsubscribe chan *SSEClient
broadcast chan SSEMessage
mutex sync.RWMutex
}
// 创建新的SSE代理
func NewSSEBroker() *SSEBroker {
return &SSEBroker{
clients: make(map[string]*SSEClient),
subscribe: make(chan *SSEClient),
unsubscribe: make(chan *SSEClient),
broadcast: make(chan SSEMessage),
}
}
// 启动代理主循环
func (b *SSEBroker) Start() {
for {
select {
case client := <-b.subscribe:
b.mutex.Lock()
b.clients[client.ID] = client
b.mutex.Unlock()
log.Printf("✓ 客户端 %s 已连接 (在线: %d)\\n", client.ID, len(b.clients))
// 发送欢迎消息
welcome := SSEMessage{
Event: "welcome",
Data: fmt.Sprintf("欢迎 %s 连接SSE服务", client.ID),
ID: fmt.Sprintf("%d", time.Now().UnixNano()),
}
select {
case client.Channel <- welcome:
case <-time.After(5 * time.Second):
log.Printf("发送欢迎消息超时: %s", client.ID)
}
case client := <-b.unsubscribe:
b.mutex.Lock()
if _, ok := b.clients[client.ID]; ok {
delete(b.clients, client.ID)
close(client.Channel)
close(client.Done)
log.Printf("✗ 客户端 %s 已断开 (在线: %d)\\n", client.ID, len(b.clients))
}
b.mutex.Unlock()
case msg := <-b.broadcast:
b.mutex.RLock()
for _, client := range b.clients {
select {
case client.Channel <- msg:
case <-time.After(5 * time.Second):
log.Printf("向客户端 %s 发送消息超时", client.ID)
}
}
b.mutex.RUnlock()
}
}
}
// 订阅SSE
func (b *SSEBroker) Subscribe(client *SSEClient) {
b.subscribe <- client
}
// 取消订阅
func (b *SSEBroker) Unsubscribe(client *SSEClient) {
b.unsubscribe <- client
}
// 广播消息
func (b *SSEBroker) Broadcast(msg SSEMessage) {
b.broadcast <- msg
}
// 获取在线客户端数量
func (b *SSEBroker) ClientCount() int {
b.mutex.RLock()
defer b.mutex.RUnlock()
return len(b.clients)
}
// SSE处理器(核心实现)
func sseHandler(broker *SSEBroker) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
// 设置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", "*") // CORS支持
// 检查是否支持Flusher(必须支持,否则无法流式输出)
flusher, ok := w.(http.Flusher)
if !ok {
http.Error(w, "Streaming unsupported", http.StatusInternalServerError)
return
}
// 创建客户端
clientID := fmt.Sprintf("client_%d", time.Now().UnixNano())
client := &SSEClient{
ID: clientID,
Channel: make(chan SSEMessage, 10),
Done: make(chan struct{}),
}
// 订阅
broker.Subscribe(client)
defer broker.Unsubscribe(client)
// 监听客户端断开
notify := r.Context().Done()
// 发送初始消息(连接建立确认)
fmt.Fprintf(w, ": SSE连接已建立\\n\\n")
flusher.Flush()
// 消息推送循环
for {
select {
case msg := <-client.Channel:
// 构造SSE消息格式
var builder strings.Builder
if msg.ID != "" {
fmt.Fprintf(&builder, "id: %s\\n", msg.ID)
}
if msg.Event != "" {
fmt.Fprintf(&builder, "event: %s\\n", msg.Event)
}
// 将数据序列化为JSON
dataJSON, err := json.Marshal(msg.Data)
if err != nil {
log.Printf("序列化消息失败: %v", err)
continue
}
fmt.Fprintf(&builder, "data: %s\\n\\n", dataJSON)
// 写入响应
fmt.Fprint(w, builder.String())
flusher.Flush()
log.Printf("→ 推送消息到 %s: %s", clientID, msg.Event)
case <-notify:
// 客户端断开连接
log.Printf("客户端 %s 断开连接", clientID)
return
case <-client.Done:
// 服务端主动关闭
return
}
}
}
}
// 发送消息的HTTP接口
func sendMessageHandler(broker *SSEBroker) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
http.Error(w, "Method not allowed", http.StatusMethodNotAllowed)
return
}
var msg SSEMessage
if err := json.NewDecoder(r.Body).Decode(&msg); err != nil {
http.Error(w, "Invalid request body", http.StatusBadRequest)
return
}
// 设置消息ID
msg.ID = fmt.Sprintf("%d", time.Now().UnixNano())
// 广播消息
broker.Broadcast(msg)
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(map[string]interface{}{
"success": true,
"clients": broker.ClientCount(),
"message": "消息已广播",
})
}
}
// 首页(包含EventSource演示)
func indexHandler(w http.ResponseWriter, r *http.Request) {
html := `<!DOCTYPE html>
<html lang="zh-CN">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<title>SSE实时通知演示</title>
<style>
body { font-family: 'Microsoft YaHei', sans-serif; max-width: 800px; margin: 0 auto; padding: 20px; background: #f5f5f5; }
h1 { color: #333; border-bottom: 2px solid #4CAF50; padding-bottom: 10px; }
#messages { background: white; border: 1px solid #ddd; border-radius: 5px; padding: 15px; height: 400px; overflow-y: auto; margin: 20px 0; }
.message { padding: 8px; margin: 5px 0; border-left: 3px solid #4CAF50; background: #f9f9f9; }
.message.error { border-left-color: #f44336; }
.message.info { border-left-color: #2196F3; }
button { background: #4CAF50; color: white; border: none; padding: 10px 20px; border-radius: 5px; cursor: pointer; margin: 5px; }
button:hover { background: #45a049; }
#status { padding: 10px; border-radius: 5px; margin: 10px 0; }
.connected { background: #d4edda; color: #155724; }
.disconnected { background: #f8d7da; color: #721c24; }
</style>
</head>
<body>
<h1>🔔 SSE实时通知演示</h1>
<div id="status" class="disconnected">❌ 未连接</div>
<div>
<button onclick="connectSSE()">连接SSE</button>
<button onclick="disconnectSSE()">断开连接</button>
<button onclick="sendTestMessage()">发送测试消息</button>
</div>
<h3>实时消息:</h3>
<div id="messages"></div>
<script>
let eventSource = null;
function connectSSE() {
if (eventSource) {
eventSource.close();
}
eventSource = new EventSource('/events');
const statusDiv = document.getElementById('status');
statusDiv.className = 'connected';
statusDiv.innerHTML = '✅ 已连接SSE服务';
// 监听默认消息
eventSource.onmessage = function(event) {
addMessage('默认消息', event.data, 'info');
};
// 监听自定义事件
eventSource.addEventListener('welcome', function(event) {
addMessage('欢迎', event.data, 'info');
});
eventSource.addEventListener('notification', function(event) {
addMessage('通知', event.data, '');
});
eventSource.addEventListener('update', function(event) {
addMessage('更新', event.data, 'info');
});
eventSource.onerror = function(error) {
statusDiv.className = 'disconnected';
statusDiv.innerHTML = '❌ 连接断开,正在自动重连…';
addMessage('错误', 'SSE连接错误', 'error');
};
}
function disconnectSSE() {
if (eventSource) {
eventSource.close();
eventSource = null;
const statusDiv = document.getElementById('status');
statusDiv.className = 'disconnected';
statusDiv.innerHTML = '❌ 已手动断开连接';
}
}
function sendTestMessage() {
fetch('/api/send', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
event: 'notification',
data: { message: '这是一条测试通知', time: new Date().toLocaleString() }
})
}).then(res => res.json()).then(data => {
console.log('发送结果:', data);
});
}
function addMessage(type, data, className) {
const messagesDiv = document.getElementById('messages');
const msgDiv = document.createElement('div');
msgDiv.className = 'message ' + className;
try {
const dataObj = JSON.parse(data);
msgDiv.innerHTML = '<strong>[' + new Date().toLocaleTimeString() + '] ' + type + ':</strong> ' + JSON.stringify(dataObj, null, 2);
} catch(e) {
msgDiv.innerHTML = '<strong>[' + new Date().toLocaleTimeString() + '] ' + type + ':</strong> ' + data;
}
messagesDiv.appendChild(msgDiv);
messagesDiv.scrollTop = messagesDiv.scrollHeight;
}
// 自动连接
connectSSE();
</script>
</body>
</html>`
w.Header().Set("Content-Type", "text/html; charset=utf-8")
w.Write([]byte(html))
}
func main() {
// 创建SSE代理
broker := NewSSEBroker()
go broker.Start()
// 启动定时推送(模拟实时数据)
go func() {
ticker := time.NewTicker(10 * time.Second)
defer ticker.Stop()
for range ticker.C {
msg := SSEMessage{
Event: "update",
Data: map[string]interface{}{
"time": time.Now().Format("2006-01-02 15:04:05"),
"value": rand.Intn(100),
"clients": broker.ClientCount(),
},
ID: fmt.Sprintf("%d", time.Now().UnixNano()),
}
broker.Broadcast(msg)
log.Printf("定时推送: update 事件 (在线: %d)", broker.ClientCount())
}
}()
// 注册路由
http.HandleFunc("/", indexHandler)
http.HandleFunc("/events", sseHandler(broker))
http.HandleFunc("/api/send", sendMessageHandler(broker))
// 健康检查
http.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) {
w.Write([]byte(fmt.Sprintf("OK (SSE客户端: %d)", broker.ClientCount())))
})
fmt.Println("===========================================")
fmt.Println("SSE服务端推送事件演示启动")
fmt.Println("===========================================")
fmt.Println("访问地址: http://localhost:8080")
fmt.Println("SSE端点: http://localhost:8080/events")
fmt.Println("发送消息: POST http://localhost:8080/api/send")
fmt.Println("健康检查: http://localhost:8080/health")
fmt.Println("===========================================")
// 优雅关闭
go func() {
sigChan := make(chan os.Signal, 1)
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM)
<-sigChan
log.Println("收到关闭信号,退出…")
os.Exit(0)
}()
log.Fatal(http.ListenAndServe(":8080", nil))
}
开发痛点与报错避坑指南
痛点1:SSE连接立即断开
问题描述:前端EventSource连接后立即触发onerror,然后不断重连。
原因分析:
- 服务端没有正确设置响应头
- 服务端在没有发送任何数据的情况下关闭了连接
- 代理服务器(如Nginx)缓冲了响应
解决方案:
// 1. 确保设置正确的响应头
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("Cache-Control", "no-cache")
w.Header().Set("Connection", "keep-alive")
// 2. 发送初始注释保持连接
fmt.Fprintf(w, ": SSE连接已建立\\n\\n")
flusher.Flush()
// 3. Nginx配置(如果使用了反向代理)
// proxy_buffering off;
// proxy_cache off;
痛点2:消息格式错误导致前端无法解析
问题描述:前端收不到消息,或者消息格式不正确。
原因分析:SSE消息格式要求严格,必须空行结束。
正确格式:
event: update
data: {"key": "value"}
(注意:必须有两个换行符)
解决方案:使用封装函数确保格式正确:
func sendSSEMessage(w http.ResponseWriter, event, data string) {
fmt.Fprintf(w, "event: %s\\ndata: %s\\n\\n", event, data)
w.(http.Flusher).Flush()
}
痛点3:浏览器自动重连导致重复消息
问题描述:客户端重连后,收到重复的消息。
原因分析:没有使用Last-Event-ID头追踪已接收的消息。
解决方案:
// 前端重连时会发送 Last-Event-ID 头
lastEventID := r.Header.Get("Last-Event-ID")
// 服务端根据lastEventID返回缺失的消息
if lastEventID != "" {
// 从存储中获取该ID之后的所有消息
missedMessages := getMessagesAfterID(lastEventID)
for _, msg := range missedMessages {
sendSSEMessage(w, msg.Event, msg.Data)
}
}
痛点4:大量并发连接导致内存溢出
问题描述:1000+并发SSE连接时,服务内存占用过高。
原因分析:
- 每个连接都占用一个goroutine
- channel缓冲区设置过大
优化方案:
全文总结+技术进阶展望
总结
本文详细介绍了Go语言中实现SSE(Server-Sent Events)的完整方案:
网硕互联帮助中心




评论前必须登录!
注册