修复WebSocket客户端数据竞态问题技术求助
WebSocket客户端数据竞态问题排查求助
开发的WebSocket客户端在连接WS服务器时出现数据竞态问题,客户端工作流程为:连接服务器、监听服务器连接、订阅指定WS端点。基于SAC007 GoWebsocket封装了自定义Wrapper,未修改涉及竞态的逻辑,但无法定位竞态原因,请求帮助排查。
主程序代码
func main() { // 检查运行模式:生产/测试 var testing bool = true args := os.Args initClient(args, &testing, &stored_data.Base_currency) // 并发处理 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 // 设置K线周期和交易对 OHLCinterval := []int{1, 5} pairs := []string{"BTC/" + stored_data.Base_currency, "ETH/" + stored_data.Base_currency, "EOS/" + stored_data.Base_currency} // 创建WS连接 wg.Add(1) pubSocket, err := ws_client.ConnectToServer(&wg, "public", testing) if err != nil { fmt.Println(err) os.Exit(1) } // 创建WS数据通道 pubCh := make(chan interface{}) defer close(pubCh) // 监听WS连接 wg.Add(1) go pubSocket.PubListen(ctx, &wg, pubCh, testing) // 订阅K线数据流 for _, interval := range OHLCinterval { wg.Add(1) go pubSocket.SubscribeToOHLC(ctx, &wg, pairs, interval) } // 后续代码与数据竞态无关 }
WS封装自定义函数
func ConnectToServer(wg *sync.WaitGroup, server string, testing bool) (*Socket, error) { defer wg.Done() if server != "public" && server != "private" { return nil, errors.New("服务器类型只能是'public'或'private'") } var socket Socket if server == "public" { socket = New(stored_data.Pub_ws_url) } else { socket = New(stored_data.Priv_ws_url) } socket.OnConnected = func(socket Socket) { log.Println("已连接到服务器") } socket.OnTextMessage = func(message string, socket Socket) { pubJsonDecoder(message, testing) log.Println(message) } socket.Connect() return &socket, nil } 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) { res = pubJsonDecoder(message, testing) ch <- res } <-ctx.Done() log.Println("关闭公共Socket连接") return } func (socket *Socket) SubscribeToOHLC(ctx context.Context, wg *sync.WaitGroup, pairs []string, interval int) { defer wg.Done() sub, _ := json.Marshal(&types.Subscribe{ Event: "subscribe", Subscription: &types.OHLCSubscription{ Interval: interval, Name: "ohlc", }, Pair: pairs, }) socket.SendBinary(sub) <-ctx.Done() }
数据竞态提示信息1
================== WARNING: DATA RACE Read at 0x00c0001145d8 by goroutine 19: kraken_client/ws_client.(*Socket).Connect.func4() /Users/grantcanty/go/src/kraken_client/ws_client/websocket_client.go:156 +0x3b5 Previous write at 0x00c0001145d8 by goroutine 20: kraken_client/ws_client.(*Socket).PubListen() /Users/grantcanty/go/src/kraken_client/ws_client/websocket_wrapper.go:73 +0x234 main.main.func4() /Users/grantcanty/go/src/kraken_client/main.go:79 +0x8d Goroutine 19 (running) created at: kraken_client/ws_client.(*Socket).Connect() /Users/grantcanty/go/src/kraken_client/ws_client/websocket_client.go:136 +0xc75 kraken_client/ws_client.ConnectToServer() /Users/grantcanty/go/src/kraken_client/ws_client/websocket_wrapper.go:39 +0x6a4 main.main() /Users/grantcanty/go/src/kraken_client/main.go:59 +0x3d4 Goroutine 20 (running) created at: main.main() /Users/grantcanty/go/src/kraken_client/main.go:79 +0x6e4 ==================
数据竞态提示信息2
================== WARNING: DATA RACE Write at 0x00c000300000 by goroutine 19: kraken_client/ws_client.(*Socket).PubListen.func1() /Users/grantcanty/go/src/kraken_client/ws_client/websocket_wrapper.go:74 +0x84 kraken_client/ws_client.(*Socket).Connect.func4() /Users/grantcanty/go/src/kraken_client/ws_client/websocket_client.go:157 +0x4ab Previous write at 0x00c000300000 by goroutine 20: kraken_client/ws_client.(*Socket).PubListen() /Users/grantcanty/go/src/kraken_client/ws_client/websocket_wrapper.go:71 +0x111 main.main.func4() /Users/grantcanty/go/src/kraken_client/main.go:79 +0x8d Goroutine 19 (running) created at: kraken_client/ws_client.(*Socket).Connect() /Users/grantcanty/go/src/kraken_client/ws_client/websocket_client.go:136 +0xc75 kraken_client/ws_client.ConnectToServer() /Users/grantcanty/go/src/kraken_client/ws_client/websocket_wrapper.go:39 +0x6a4 main.main() /Users/grantcanty/go/src/kraken_client/main.go:59 +0x3d4 Goroutine 20 (running) created at: main.main() /Users/grantcanty/go/src/kraken_client/main.go:79 +0x6e4 ==================
问题分析
从竞态提示可以明确,问题出在OnTextMessage回调函数的并发读写上:
- Goroutine 19(由
Connect方法启动的协程)在读取OnTextMessage并执行它 - Goroutine 20(
PubListen的协程)在覆盖写入OnTextMessage - 同时,
OnTextMessage内部的逻辑还在对外部变量res进行写入,进一步加剧竞态
根本原因是Socket结构体的OnTextMessage字段被多个协程并发访问,且没有任何同步保护措施。
解决方案
方案1:给Socket结构体添加互斥锁保护回调字段
修改Socket结构体,加入sync.Mutex,在设置和调用回调时加锁:
import "sync" type Socket struct { // 原有字段... OnConnected func(socket Socket) OnTextMessage func(message string, socket Socket) mu sync.Mutex // 添加互斥锁 } // 设置OnTextMessage时加锁 func (s *Socket) SetOnTextMessage(f func(message string, socket Socket)) { s.mu.Lock() defer s.mu.Unlock() s.OnTextMessage = f } // 调用OnTextMessage时加锁 func (s *Socket) invokeOnTextMessage(message string) { s.mu.Lock() defer s.mu.Unlock() if s.OnTextMessage != nil { s.OnTextMessage(message, *s) } }
然后修改原有代码:
- 在
ConnectToServer中,用socket.SetOnTextMessage(...)替代直接赋值 - 在
PubListen中,同样用socket.SetOnTextMessage(...)替代直接赋值 - 原GoWebsocket库中调用
OnTextMessage的地方,替换为调用invokeOnTextMessage
方案2:避免在并发场景下修改回调
调整流程,在Connect完成后再设置最终的OnTextMessage回调,避免在协程中覆盖:
比如在main中,连接完成后直接设置监听回调,不再启动PubListen协程修改回调:
// 创建WS连接后直接设置回调 pubSocket.SetOnTextMessage(func(message string, socket ws_client.Socket) { res := pubJsonDecoder(message, testing) pubCh <- res })
同时移除PubListen相关的协程逻辑,减少并发修改的可能。
额外优化:避免回调访问外部未同步变量
在PubListen中,回调直接访问外部的res变量,这也是潜在竞态点。可以将res移到回调内部,或者用通道传递时直接处理:
socket.SetOnTextMessage(func(message string, socket Socket) { res := pubJsonDecoder(message, testing) ch <- res })
内容的提问来源于stack exchange,提问作者Grant
相关产品推荐
相关产品推荐

