如何基于websocket-client实现Binance WebSocket断线重连及告警?
问题描述
我正在编写从Binance采集加密货币数据的代码,Binance会在24小时后自动断开连接。我想实现断线后自动重连,本以为run_forever能实现,但程序遇到错误就崩溃了。该程序需要在服务器7*24运行,同时还要添加Telegram/Discord机器人告警功能,请问如何实现重连,告警代码应写在何处?
报错信息
Traceback (most recent call last): File "exchanges/binance/binance_ticker.py", line 97, in <module> start() File "exchanges/binance/binance_ticker.py", line 94, in start rel.dispatch() File "/home/pyjobs/.local/lib/python3.8/site-packages/rel/rel.py", line 205, in dispatch registrar.dispatch() File "/home/pyjobs/.local/lib/python3.8/site-packages/rel/registrar.py", line 72, in dispatch if not self.loop(): File "/home/pyjobs/.local/lib/python3.8/site-packages/rel/registrar.py", line 81, in loop e = self.check_events() File "/home/pyjobs/.local/lib/python3.8/site-packages/rel/registrar.py", line 232, in check_events self.callback('read', fd) File "/home/pyjobs/.local/lib/python3.8/site-packages/rel/registrar.py", line 125, in callback self.events[etype][fd].callback() File "/home/pyjobs/.local/lib/python3.8/site-packages/rel/listener.py", line 108, in callback if not self.cb(*self.args) and not self.persist and self.active: File "/home/pyjobs/.local/lib/python3.8/site-packages/websocket/_app.py", line 349, in read op_code, frame = self.sock.recv_data_frame(True) File "/home/pyjobs/.local/lib/python3.8/site-packages/websocket/_core.py", line 401, in recv_data_frame frame = self.recv_frame() File "/home/pyjobs/.local/lib/python3.8/site-packages/websocket/_core.py", line 440, in recv_frame return self.frame_buffer.recv_frame() File "/home/pyjobs/.local/lib/python3.8/site-packages/websocket/_abnf.py", line 352, in recv_frame payload = self.recv_strict(length) File "/home/pyjobs/.local/lib/python3.8/site-packages/websocket/_abnf.py", line 373, in recv_strict bytes_ = self.recv(min(16384, shortage)) File "/home/pyjobs/.local/lib/python3.8/site-packages/websocket/_core.py", line 524, in _recv return recv(self.sock, bufsize) File "/home/pyjobs/.local/lib/python3.8/site-packages/websocket/_socket.py", line 122, in recv raise WebSocketConnectionClosedException( websocket._exceptions.WebSocketConnectionClosedException: Connection to remote host was lost.
现有代码
import websocket import rel uri = "wss://stream.binance.com:9443/ws/!ticker@arr" def on_message(ws, message): print(message) def on_error(ws, error): print(error) write_logs(error) def on_close(ws, close_status_code, close_msg): print("### closed ###") write_logs(str(close_status_code) + str(close_msg)) start( def on_open(ws): print("Opened connection") start() websocket.enableTrace(True) ws = websocket.WebSocketApp(uri, on_open = on_open, on_message=on_message, on_error = on_error, on_close (on_close) ws.run_forever(dispatcher=rel) #Set the dispatcher to automatic reconnection. rel.signal(2, rel.abort) # Keyboard Interrupt rel.dispatch() start()
解决方案
一、实现自动重连
你的现有代码存在语法错误(如on_close (on_close)缺少等号、start(未闭合),且run_forever配合rel的自动重连逻辑未正确配置。以下是修正后的可靠重连实现:
- 把WebSocket初始化逻辑封装到
start函数,确保断开后能重新创建连接实例 - 利用
run_forever的reconnect参数,让库自动处理重连间隔,避免手动延迟的繁琐 - 在错误和关闭回调中触发重连,同时保证异常不会中断程序
修正后的核心代码:
import websocket import rel import time uri = "wss://stream.binance.com:9443/ws/!ticker@arr" def write_logs(content): # 补全你的日志写入逻辑 with open("binance_ticker_logs.txt", "a", encoding="utf-8") as f: log_line = f"{time.strftime('%Y-%m-%d %H:%M:%S')} - {str(content)}\n" f.write(log_line) def on_message(ws, message): print(message) # 这里添加你的数据处理/存储逻辑 def on_error(ws, error): error_info = f"连接错误: {str(error)}" print(error_info) write_logs(error_info) send_alert(error_info) # 触发告警 def on_close(ws, close_status_code, close_msg): close_info = f"连接关闭,状态码: {close_status_code}, 消息: {close_msg}" print(f"### closed ### {close_info}") write_logs(close_info) send_alert(close_info) # 触发告警 def on_open(ws): success_info = "Binance数据采集连接已成功建立" print(success_info) write_logs(success_info) send_alert(success_info) # 可选:连接恢复通知 def start(): websocket.enableTrace(True) ws = websocket.WebSocketApp( uri, on_open=on_open, on_message=on_message, on_error=on_error, on_close=on_close ) # 使用rel调度器,reconnect参数指定自动重连间隔(秒) ws.run_forever(dispatcher=rel, reconnect=5) # 注册Ctrl+C退出信号 rel.signal(2, rel.abort) # 启动服务 start() rel.dispatch()
关键说明:
ws.run_forever(reconnect=5):由websocket库自动处理重连逻辑,间隔5秒,比手动在回调中重启更稳定- 封装
start函数为可重复调用的逻辑,每次重连都会创建全新的WebSocket实例,避免旧实例的状态污染 - 保留日志写入,方便后续排查连接问题
二、Telegram/Discord告警代码位置
告警逻辑需封装为独立的send_alert函数,在三个关键节点调用:
on_error:连接出现异常时发送告警on_close:连接被主动断开时发送告警on_open:连接成功重建时发送恢复通知(可选)
Telegram告警示例
需先创建Telegram机器人,获取Bot Token和个人Chat ID:
import requests TELEGRAM_BOT_TOKEN = "你的机器人Token" TELEGRAM_CHAT_ID = "你的Chat ID" def send_alert(message): try: url = f"https://api.telegram.org/bot{TELEGRAM_BOT_TOKEN}/sendMessage" payload = {"chat_id": TELEGRAM_CHAT_ID, "text": message} response = requests.post(url, data=payload) response.raise_for_status() except Exception as e: error_msg = f"Telegram告警发送失败: {str(e)}" print(error_msg) write_logs(error_msg)
Discord告警示例
需创建Discord机器人并生成Webhook URL:
import requests DISCORD_WEBHOOK_URL = "你的Discord Webhook URL" def send_alert(message): try: payload = {"content": message} response = requests.post(DISCORD_WEBHOOK_URL, json=payload) response.raise_for_status() except Exception as e: error_msg = f"Discord告警发送失败: {str(e)}" print(error_msg) write_logs(error_msg)
将send_alert函数放在代码顶部,然后在on_error、on_close、on_open中按需调用即可。
额外建议
- 进程守护:服务器上用
systemd或supervisord托管Python进程,即使程序意外崩溃也能自动重启 - 重连次数限制:可添加重连计数逻辑,多次重连失败后停止尝试并发送紧急告警
- 日志轮转:用
logging模块替代简单文件写入,实现日志自动轮转,避免日志文件过大
内容的提问来源于stack exchange,提问作者Nez
相关产品推荐
相关产品推荐

