You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于UserID的Go专属用户Server-Sent Events实现求助

Go SSE 针对特定用户推送消息实现方案

核心设计思路

  • 维护一个用户连接管理器,按UserID分组存储客户端SSE连接
  • 每个客户端连接启动独立goroutine,监听专属消息通道,同时监听连接关闭事件以触发清理
  • 基于gorilla-sessions完成身份校验,仅允许已认证用户建立SSE连接

完整可运行代码

package main

import (
	"context"
	"fmt"
	"net/http"
	"sync"
	"time"

	"github.com/gorilla/sessions"
)

// 全局存储session,实际生产建议用Redis等持久化存储
var store = sessions.NewCookieStore([]byte("your-secret-key-here"))

// Client 代表一个SSE客户端连接
type Client struct {
	msgChan chan string
	userID  string
}

// ConnectionManager 管理所有用户的SSE连接
type ConnectionManager struct {
	mu      sync.RWMutex
	clients map[string][]*Client // key: UserID, value: 该用户的所有连接
}

func NewConnectionManager() *ConnectionManager {
	return &ConnectionManager{
		clients: make(map[string][]*Client),
	}
}

// AddClient 添加新的客户端连接到管理器
func (cm *ConnectionManager) AddClient(userID string, client *Client) {
	cm.mu.Lock()
	defer cm.mu.Unlock()
	cm.clients[userID] = append(cm.clients[userID], client)
}

// RemoveClient 从管理器移除指定客户端连接
func (cm *ConnectionManager) RemoveClient(userID string, client *Client) {
	cm.mu.Lock()
	defer cm.mu.Unlock()

	clients, ok := cm.clients[userID]
	if !ok {
		return
	}

	// 找到并删除目标客户端
	for i, c := range clients {
		if c == client {
			cm.clients[userID] = append(clients[:i], clients[i+1:]...)
			// 如果该用户没有剩余连接,删除key
			if len(cm.clients[userID]) == 0 {
				delete(cm.clients, userID)
			}
			break
		}
	}
	close(client.msgChan)
}

// SendToUser 向指定UserID的所有连接发送消息
func (cm *ConnectionManager) SendToUser(userID string, message string) {
	cm.mu.RLock()
	defer cm.mu.RUnlock()

	clients, ok := cm.clients[userID]
	if !ok {
		return
	}

	// 遍历所有连接发送消息,避免阻塞
	for _, client := range clients {
		select {
		case client.msgChan <- message:
		case <-time.After(1 * time.Second):
			// 发送超时,标记该连接可能已失效,后续清理
			go cm.RemoveClient(userID, client)
		}
	}
}

// SSEHandler 处理SSE连接请求
func (cm *ConnectionManager) SSEHandler(w http.ResponseWriter, r *http.Request) {
	// 从session获取UserID,完成身份校验
	session, err := store.Get(r, "auth-session")
	if err != nil {
		http.Error(w, "Unauthorized", http.StatusUnauthorized)
		return
	}

	userID, ok := session.Values["user_id"].(string)
	if !ok || userID == "" {
		http.Error(w, "Unauthorized", http.StatusUnauthorized)
		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", "*") // 根据实际跨域需求调整

	// 创建客户端实例
	client := &Client{
		msgChan: make(chan string, 10),
		userID:  userID,
	}
	cm.AddClient(userID, client)
	defer cm.RemoveClient(userID, client)

	// 监听连接关闭事件(通过请求上下文的Done通道)
	ctx := r.Context()
	for {
		select {
		case <-ctx.Done():
			// 连接已关闭,退出循环,触发defer清理
			return
		case msg := <-client.msgChan:
			// 按SSE格式发送消息
			fmt.Fprintf(w, "data: %s\n\n", msg)
			// 强制刷新响应缓冲区,确保消息立即发送
			if flusher, ok := w.(http.Flusher); ok {
				flusher.Flush()
			}
		}
	}
}

// 模拟业务逻辑:向指定用户发送通知
func sendNotification(cm *ConnectionManager, userID string, content string) {
	cm.SendToUser(userID, content)
}

func main() {
	cm := NewConnectionManager()

	// 模拟登录接口(实际需替换为真实认证逻辑)
	http.HandleFunc("/login", func(w http.ResponseWriter, r *http.Request) {
		session, _ := store.Get(r, "auth-session")
		// 这里直接模拟设置UserID,实际应从登录请求中校验用户后获取
		session.Values["user_id"] = "user_123"
		session.Save(r, w)
		fmt.Fprintln(w, "Logged in successfully")
	})

	// SSE连接接口
	http.HandleFunc("/sse", cm.SSEHandler)

	// 模拟发送消息的接口(供业务调用)
	http.HandleFunc("/send", func(w http.ResponseWriter, r *http.Request) {
		userID := r.URL.Query().Get("user_id")
		content := r.URL.Query().Get("content")
		if userID == "" || content == "" {
			http.Error(w, "user_id and content are required", http.StatusBadRequest)
			return
		}
		sendNotification(cm, userID, content)
		fmt.Fprintln(w, "Message sent")
	})

	fmt.Println("Server starting on :8080...")
	http.ListenAndServe(":8080", nil)
}

关键细节说明

  1. 连接生命周期管理:

    • 通过请求上下文ctx.Done()监听连接关闭事件,一旦触发就执行清理逻辑,确保goroutine正常退出,避免泄漏
    • 每个客户端的msgChan在移除时会被关闭,对应的监听goroutine会自然终止
  2. 消息发送可靠性:

    • 发送消息时使用select加超时机制,避免因客户端连接异常导致发送阻塞
    • 发送后强制刷新响应缓冲区,确保消息即时推送给客户端
  3. 身份校验集成:

    • 直接复用gorilla-sessions的session数据获取UserID,和你现有认证逻辑无缝对接

轮询 vs SSE 选择建议

  • 如果你的业务对实时性要求较低(比如允许5秒延迟),Ajax轮询实现更简单,但会产生大量重复请求,增加服务器负载
  • SSE适合需要低延迟实时推送的场景,长连接模式能减少请求数,资源利用率更高。只要处理好连接管理和goroutine泄漏(本示例已解决),比轮询更优

解决你提到的现有方案问题

  • 无法向单个用户发送:本示例通过ConnectionManager按UserID分组管理连接,调用SendToUser即可精准推送
  • 连接关闭后goroutine不停止:通过监听ctx.Done()触发清理,关闭消息通道,goroutine会正常退出
  • 隐私窗口关闭重开失效:每次新连接都会重新注册到管理器,只要用户session有效,就能正常接收消息

内容的提问来源于stack exchange,提问作者Digital_ghost

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.12 06:45:33