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

Julia SimpleWebsockets重连逻辑问题:connectError的notify失效

问题

使用Julia的SimpleWebsockets库建立WebSocket连接,需要实现断开后自动重连,但标注的notify调用失效,导致连接失败后程序停滞在等待状态。

代码如下:

using SimpleWebsockets

function connect_loop(C::Coinbase, msg_channel::Channel)
    while true
        client = WebsocketClient()

        ended = Condition()

        listen(client, :connect) do ws
            listen(ws, :message) do message
                put!(msg_channel, message)
            end
            listen(ws, :error) do reason
                @warn "Websocket connection error" reason...
                notify(ended)
            end
            listen(ws, :close) do reason
                @warn "Websocket connection closed" reason...
                notify(ended)
            end
            # for msg in C.command_channel
            #     send(ws, JSON.json(msg))
            # end
        end

        listen(C.client, :connectError) do err
            println("Websocket Connect Error.")
            notify(ended)  # 这里失效了!!!!!!!!!
        end

        println("connecting")
        @async open(C.client, "wss://ws-feed.pro.coinbase.com")
        println("wait for ended")
        wait(ended)
        println("ended")
        sleep(2)
    end
end

测试流程:先通过WiFi建立连接,成功后关闭WiFi,预期程序进入重连循环,但实际停滞在等待notify的状态。程序输出:

connecting
wait for ended
┌ Warning: Websocket connection closed
│   code = 1006
│   description = "could not ping the server."
└ @ Main ~/src/tibra/Tibra.jl/src/MarketDataIO/Coinbase.jl:60
ended
connecting
wait for ended
Websocket Connect Error.
解决方案

问题根源

  1. 实例混淆:代码中创建了局部client = WebsocketClient(),但监听connectError和调用open时用的是C.client,两个实例完全独立,导致connect事件监听绑定到了未使用的局部实例上,而connectError的监听又无法触发当前循环的ended。
  2. 事件监听累积:每次循环都会给C.client新增connect和connectError监听,旧监听绑定的是上一轮循环的ended变量,新通知会触发旧的Condition,而非当前循环等待的实例。

修复代码

using SimpleWebsockets

function connect_loop(C::Coinbase, msg_channel::Channel)
    while true
        ended = Condition()
        # 清除上一轮循环绑定的旧监听,避免干扰
        off(C.client, :connect)
        off(C.client, :connectError)

        listen(C.client, :connect) do ws
            # 清除WebSocket连接上的旧事件监听
            off(ws, :message)
            off(ws, :error)
            off(ws, :close)
            
            listen(ws, :message) do message
                put!(msg_channel, message)
            end
            listen(ws, :error) do reason
                @warn "Websocket connection error" reason...
                notify(ended)
            end
            listen(ws, :close) do reason
                @warn "Websocket connection closed" reason...
                notify(ended)
            end
            # for msg in C.command_channel
            #     send(ws, JSON.json(msg))
            # end
        end

        listen(C.client, :connectError) do err
            println("Websocket Connect Error.")
            notify(ended; all=true)  # 确保所有等待该Condition的任务被唤醒
        end

        println("connecting")
        @async open(C.client, "wss://ws-feed.pro.coinbase.com")
        println("wait for ended")
        wait(ended)
        println("ended")
        sleep(2)
    end
end

额外优化点

  • 添加超时机制:可以在wait(ended)前启动一个异步任务,超时后主动notify(ended),避免极端情况下无限等待。
  • 限制重连次数:记录重连失败次数,达到阈值时停止循环,防止无意义的重复尝试。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 19:45:41