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

Go WebSocket客户端channel传值仅返回nil问题求助

问题描述

WebSocket客户端开发中需要将接收的响应传入channel供主协程做后续业务处理,运行时出现异常:channel仅返回一次nil值后就无任何输出。排查发现OnTextMessage回调内部可正常打印解码后的响应值,但回调外部的传值位置无日志输出,主协程仅能读取到nil值。


原始问题代码

main函数实现

package main

import (
    "context"
    "fmt"
    "kraken_client/stored_data"
    "kraken_client/ws_client"
    "os"
    "os/signal"
    "strings"
    "sync"
    "syscall"
)

func main() {
    // check if in production or testing mode & find base curency
    var testing bool = true
    args := os.Args
    isTesting(args, &testing, &stored_data.Base_currency)

    // go routine handler
    comms := make(chan os.Signal, 1)
    signal.Notify(comms, os.Interrupt, syscall.SIGTERM)
    ctx := context.Background()
    ctx, cancel := context.WithCancel(ctx)
    var wg sync.WaitGroup

    // set ohlc interval and pairs
    OHLCinterval := 5
    pairs := []string{"BTC/" + stored_data.Base_currency, "EOS/" + stored_data.Base_currency}

    // create ws connections
    pubSocket, err := ws_client.ConnectToServer("public", testing)
    if err != nil {
        fmt.Println(err)
        os.Exit(1)
    }

    // listen to websocket connections
    ch := make(chan interface{})
    wg.Add(1)
    go pubSocket.PubListen(ctx, &wg, ch, testing)

    // subscribe to a stream
    pubSocket.SubscribeToOHLC(pairs, OHLCinterval)

    go func() {
        for c := range ch {
            fmt.Println(c)
        }
    }()

    <-comms
    cancel()
    wg.Wait()
    defer close(ch)
}

PubListen函数实现

func (socket *Socket) PubListen(ctx context.Context, wg *sync.WaitGroup, ch chan interface{}, testing bool) {
    defer wg.Done()
    defer socket.Close()

    var res interface{}

    socket.OnTextMessage = func(message string, socket Socket) {
        //log.Println(message)
        res = pubJsonDecoder(message, testing) // this function decodes the message and returns an interface
        log.Println(res) // this is printing the correctly decoded value.

    }

    ch <- res
    log.Println(res) // does not print a value
    log.Println(ch) // does not print a value

    <-ctx.Done()
    log.Println("closing public socket")
    return
}

根因分析
  • 执行时序错误:PubListen中给socket.OnTextMessage赋值回调后,立刻执行ch <- res,此时回调还未被任何WebSocket消息触发,res是接口类型的零值nil,因此第一次发送到channel的值就是nil。
  • 传值位置错误:向channel发送数据的逻辑写在了回调外部,OnTextMessage是WebSocket收到消息时才会异步触发的回调,外部代码不会等待回调执行就会继续向下运行,后续回调中对res的更新根本不会触发新的channel发送,自然不会有后续内容输出。
  • 并发安全问题:用外部的res变量中转消息,多个WebSocket消息同时到达时会并发读写res,存在数据竞争风险。
  • 关闭逻辑错误:defer close(ch)写在主协程的<-comms阻塞逻辑之后,且channel应由发送方负责关闭,放在接收侧的主协程关闭不符合channel使用惯例。

修正后代码

修正后的PubListen函数

func (socket *Socket) PubListen(ctx context.Context, wg *sync.WaitGroup, ch chan interface{}, testing bool) {
    defer wg.Done()
    defer socket.Close()
    // 发送方在函数退出前关闭channel,通知接收方数据传输结束
    defer close(ch)

    socket.OnTextMessage = func(message string, socket Socket) {
        // 回调内直接解码消息,不使用外部变量中转,避免并发问题
        res := pubJsonDecoder(message, testing)
        log.Println(res)
        // 加select处理退出信号,避免程序退出时发送操作阻塞
        select {
        case ch <- res:
            // 消息发送成功
        case <-ctx.Done():
            // 收到退出信号,终止发送
            return
        }
    }

    <-ctx.Done()
    log.Println("closing public socket")
    return
}

main函数修正点

  • 移除主协程中错误位置的defer close(ch)逻辑
  • 将读取channel的goroutine启动逻辑移到PubListen启动之前,避免极端时序下的发送阻塞问题,调整后相关片段如下:
// listen to websocket connections
ch := make(chan interface{})
// 先启动channel接收协程
go func() {
    for c := range ch {
        fmt.Println(c)
    }
}()

wg.Add(1)
go pubSocket.PubListen(ctx, &wg, ch, testing)

// subscribe to a stream
pubSocket.SubscribeToOHLC(pairs, OHLCinterval)

<-comms
cancel()
wg.Wait()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 01:54:35