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

在Go的Fiber框架中结合Channel实现SSE遇问题求助

问题分析

你的代码出现panic和逻辑异常的核心原因有两个:

  1. SetBodyStreamWriter的异步特性:Fiber(底层基于fasthttp)的SetBodyStreamWriter是异步执行的,调用后sseHandler会立刻继续执行直到返回,导致defer函数提前触发,直接关闭了sseChan并置为nil,同时c.Context()在Handler返回后可能被回收,后续StreamWriter里访问c.Context().Done()就会触发空指针panic。
  2. 全局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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 18:34:55