Go结构体使用多通道实现SSE时仅首条消息发送成功后挂起
问题根因
核心是三个逻辑错误导致通道阻塞:
Format()方法没有循环消费逻辑:当前实现仅从data通道读取1次数据、处理1次就直接退出协程,第一条消息发送完成后就没有消费者再读取data通道的后续数据,所有往data通道写入的操作都会永久挂起。- 协程启动逻辑错误:每次SSE请求进来都会重复启动
Connect()协程重复连接WhatsApp服务,同时重复启动Format()协程,全局唯一的无缓冲通道被多个协程抢占读写,直接导致链路混乱。 - 通道读写没有退出保护:无论是
Format()读通道、Connect()写通道,还是请求循环读logs通道,都没有绑定请求/程序的生命周期做退出判断,客户端断开后相关协程全部泄漏,后续写入操作会因为没有接收方永久阻塞。
修复步骤
1. 改造Format方法为持续消费模式
给方法加生命周期上下文,循环读取data通道数据,处理后写入logs通道,任意环节收到退出信号立刻停止:
func (p *DataPasser) Format(ctx context.Context) { for { select { case data := <-p.data: var msg string if len(data.event) > 0 { msg = fmt.Sprintf("event: %v\ndata: %v\n\n", data.event, data.message) } else { msg = fmt.Sprintf("data: %v\n\n", data.message) } // 写logs通道也加退出判断,避免阻塞 select { case p.logs <- msg: case <-ctx.Done(): return } case <-ctx.Done(): return } } }
2. 调整HandleSignal逻辑,绑定请求生命周期
- 移除每次请求启动
Connect()的逻辑,WhatsApp连接全局只需要启动一次 - 连接配额的获取和释放用
defer统一处理,避免异常路径漏释放 - 把请求上下文传给
Format()协程,客户端断开时自动终止协程 - 写SSE响应时增加错误判断,客户端断开时直接退出
func (p *DataPasser) HandleSignal(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "text/event-stream; charset=utf-8") w.Header().Set("Cache-Control", "no-cache") w.Header().Set("Connection", "keep-alive") setupCORS(&w, r) fmt.Println("Client connected from IP:", r.RemoteAddr) // 校验最大连接数,满了直接返回 select { case p.connection <- struct{}{}: default: http.Error(w, "Reach max client limit", 429) return } defer func() { <-p.connection }() // 函数退出自动释放连接名额 flusher, ok := w.(http.Flusher) if !ok { http.Error(w, "Internal error", 500) return } fmt.Fprint(w, "event: notification\ndata: Connection to WhatsApp server ...\n\n") flusher.Flush() // 启动和当前请求绑定的格式化协程 ctx := r.Context() go p.Format(ctx) for { select { case c := <-p.logs: // 写响应失败说明客户端已断开 if _, err := fmt.Fprint(w, c); err != nil { fmt.Println("Write to client failed, connection closed") return } flusher.Flush() case <-ctx.Done(): fmt.Println("Connection closed") return } } }
3. 调整main函数全局初始化逻辑
把Connect()移到main里全局启动一次,避免重复连接WhatsApp服务:
var passer *DataPasser const maxClients = 1 func main() { // 初始化通道 passer = &DataPasser{ data: make(chan sseData), logs: make(chan string), connection: make(chan struct{}, maxClients), } // 全局只启动一次WhatsApp连接逻辑 go Connect() http.HandleFunc("/sse", passer.HandleSignal) go func() { if err := http.ListenAndServe(":1234", nil); err != nil { panic(err) } }() // 监听系统退出信号 c := make(chan os.Signal, 1) signal.Notify(c, os.Interrupt, syscall.SIGTERM) <-c if client.IsConnected() { client.Disconnect() } }
4. 给Connect函数的通道写操作加保护
所有往passer.data写入数据的逻辑,都增加全局退出信号监听,避免程序退出或无消费者时永久阻塞。示例改造:
// 原来的直接写: // passer.data <- sseData{...} // 改成带退出判断的写法: select { case passer.data <- sseData{ event: "notification", message: "Reconnecting to WhatsApp server ...", }: case <-ctx.Done(): // 这里传入全局程序退出的context即可 return }
额外说明
当前实现仅支持单客户端连接(和设置的maxClients=1匹配),如果后续需要支持多客户端SSE推送,不能用全局单data/logs通道,需要实现订阅广播机制:每个客户端连接时注册自己的消息通道,消息产生时批量推送给所有在线的连接通道,断开时注销通道避免泄漏。
内容的提问来源于stack exchange,提问作者Hasan A Yousef
相关产品推荐
相关产品推荐

