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

修复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回调函数的并发读写上:

  1. Goroutine 19(由Connect方法启动的协程)在读取OnTextMessage并执行它
  2. Goroutine 20(PubListen的协程)在覆盖写入OnTextMessage
  3. 同时,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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 21:12:29