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

如何基于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的自动重连逻辑未正确配置。以下是修正后的可靠重连实现:

  1. 把WebSocket初始化逻辑封装到start函数,确保断开后能重新创建连接实例
  2. 利用run_forever的reconnect参数,让库自动处理重连间隔,避免手动延迟的繁琐
  3. 在错误和关闭回调中触发重连,同时保证异常不会中断程序

修正后的核心代码:

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函数,在三个关键节点调用:

  1. on_error:连接出现异常时发送告警
  2. on_close:连接被主动断开时发送告警
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 06:15:40