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.
解决方案
问题根源
- 实例混淆:代码中创建了局部
client = WebsocketClient(),但监听connectError和调用open时用的是C.client,两个实例完全独立,导致connect事件监听绑定到了未使用的局部实例上,而connectError的监听又无法触发当前循环的ended。 - 事件监听累积:每次循环都会给
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
相关产品推荐
相关产品推荐

