在Go的Fiber框架中结合Channel实现SSE遇问题求助
问题分析
你的代码出现panic和逻辑异常的核心原因有两个:
SetBodyStreamWriter的异步特性:Fiber(底层基于fasthttp)的SetBodyStreamWriter是异步执行的,调用后sseHandler会立刻继续执行直到返回,导致defer函数提前触发,直接关闭了sseChan并置为nil,同时c.Context()在Handler返回后可能被回收,后续StreamWriter里访问c.Context().Done()就会触发空指针panic。- 全局Channel的设计缺陷:全局的
sseChan只能支持单个客户端连接,多客户端同时连接时会覆盖Channel,导致之前的客户端无法收到消息。
修正后的代码
package main import ( "bufio" "fmt" "sync" "time" "github.com/gofiber/fiber/v2" "github.com/valyala/fasthttp" ) // 用map管理所有活跃客户端的专属Channel,支持多客户端连接 var clientChannels = make(map[chan string]struct{}) var mu sync.Mutex func sseHandler(c *fiber.Ctx) error { // 设置SSE标准响应头 c.Set("Content-Type", "text/event-stream") c.Set("Cache-Control", "no-cache") c.Set("Connection", "keep-alive") c.Set("Transfer-Encoding", "chunked") // 为当前客户端创建带缓冲的专属Channel,避免发送阻塞 clientChan := make(chan string, 10) mu.Lock() clientChannels[clientChan] = struct{}{} mu.Unlock() // 客户端断开连接时自动清理资源 defer func() { mu.Lock() delete(clientChannels, clientChan) mu.Unlock() close(clientChan) fmt.Println("客户端连接关闭,资源已清理") }() fmt.Println("新客户端已建立SSE连接") ctx := c.Context() return ctx.SetBodyStreamWriter(fasthttp.StreamWriter(func(w *bufio.Writer) { for { select { case message := <-clientChan: // 写入符合SSE格式的消息 if _, err := fmt.Fprintf(w, "data: %s\n\n", message); err != nil { fmt.Printf("写入消息失败: %v\n", err) return } // 强制刷新缓冲区,确保消息即时发送到客户端 if err := w.Flush(); err != nil { fmt.Printf("刷新缓冲区失败: %v\n", err) return } case <-ctx.Done(): fmt.Println("客户端主动断开连接") return } } })) } func fireEvent(c *fiber.Ctx) error { // 向所有活跃客户端广播消息 mu.Lock() defer mu.Unlock() msg := time.Now().Format("15:04:05") fmt.Printf("向所有客户端广播消息: %s\n", msg) for ch := range clientChannels { // 非阻塞发送,避免单个客户端异常阻塞整个广播流程 select { case ch <- msg: default: fmt.Println("客户端Channel缓冲已满,消息丢弃") } } return c.SendString("事件已触发,消息已广播") } func main() { app := fiber.New() app.Get("/sse", sseHandler) app.Post("/fire", fireEvent) fmt.Println("服务启动在 :3000") if err := app.Listen(":3000"); err != nil { fmt.Printf("服务启动失败: %v\n", err) } }
关键修改说明
- 多客户端支持:用
clientChannels配合互斥锁管理每个客户端的专属Channel,解决了全局Channel只能支持单客户端的问题。 - 资源生命周期管理:每个客户端连接创建独立Channel,断开时通过defer自动清理,避免资源泄漏;直接复用Fiber的Context Done通道,确保客户端断开时及时退出循环。
- 异步执行问题修复:不再依赖WaitGroup阻塞Handler返回,让
SetBodyStreamWriter的异步回调自行处理消息循环,Handler返回后Fiber会自动保持连接直到回调退出。 - 消息发送可靠性:使用带缓冲的Channel,发送时用非阻塞select避免流程卡住,写入后强制刷新缓冲区确保消息即时送达客户端。
内容的提问来源于stack exchange,提问作者Mustaghees
相关产品推荐
相关产品推荐

