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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 11:18:15