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
相关产品推荐
相关产品推荐

