使用Tornado Websocket对接Binance流时出现异常错误求助
问题描述
程序使用Tornado连接Binance的WebSocket流,初期可正常接收数据,但运行不到1分钟就抛出错误:
{"error":{"code":3,"msg":"Invalid JSON: expected value at line 1 column 1"}}
我的代码如下:
#!/usr/bin/env python # -*- coding: utf-8 -*- import json from tornado.ioloop import IOLoop, PeriodicCallback from tornado import gen from tornado.websocket import websocket_connect api_data = { "method": "SUBSCRIBE", "params": [ "btcusdt@trade" ], "id": "BTCUSDT", } class Client(object): def __init__(self, url, timeout): self.url = url self.timeout = timeout self.ioloop = IOLoop.instance() self.ws = None self.connect() PeriodicCallback(self.keep_alive, 20000).start() self.ioloop.start() @gen.coroutine def connect(self): print("trying to connect") try: self.ws = yield websocket_connect(self.url) except Exception as e: print("connection error") else: yield self.ws.write_message(json.dumps(api_data)) print("connected") self.run() @gen.coroutine def run(self): while True: msg = yield self.ws.read_message() if msg is None: print("connection closed") self.ws = None break else: print(msg) exit() def keep_alive(self): if self.ws is None: self.connect() else: self.ws.write_message("keep alive") if __name__ == "__main__": client = Client("wss://stream.binance.com:9443/ws", 5)
解决方案
错误根源是Binance WebSocket API要求心跳必须是JSON格式的PING消息,你的代码发送的是普通字符串"keep alive",不符合API规范,因此被服务器判定为无效JSON并断开连接。
修改步骤
- 调整心跳函数,发送符合Binance要求的JSON格式PING消息
- 优化协程内的退出逻辑,避免使用
exit()导致IOLoop异常
修改后的代码
#!/usr/bin/env python # -*- coding: utf-8 -*- import json from tornado.ioloop import IOLoop, PeriodicCallback from tornado import gen from tornado.websocket import websocket_connect api_data = { "method": "SUBSCRIBE", "params": [ "btcusdt@trade" ], "id": "BTCUSDT", } # Binance要求的PING消息格式 ping_msg = json.dumps({"method": "PING", "id": 1}) class Client(object): def __init__(self, url, timeout): self.url = url self.timeout = timeout self.ioloop = IOLoop.instance() self.ws = None self.connect() # Binance建议每30秒发送一次心跳,调整为25秒更稳妥 PeriodicCallback(self.keep_alive, 25000).start() self.ioloop.start() @gen.coroutine def connect(self): print("trying to connect") try: self.ws = yield websocket_connect(self.url) except Exception as e: print(f"connection error: {e}") # 连接失败后延迟3秒重试 self.ioloop.call_later(3, self.connect) else: yield self.ws.write_message(json.dumps(api_data)) print("connected") self.run() @gen.coroutine def run(self): while True: if self.ws is None: break msg = yield self.ws.read_message() if msg is None: print("connection closed, reconnecting...") self.ws = None self.connect() break else: # 处理PONG响应,验证心跳有效性 msg_data = json.loads(msg) if msg_data.get("method") == "PONG": print(f"received pong, id: {msg_data.get('id')}") else: print(msg) @gen.coroutine def keep_alive(self): if self.ws is None: self.connect() else: try: yield self.ws.write_message(ping_msg) print("sent ping") except Exception as e: print(f"send ping error: {e}") self.ws = None self.connect() if __name__ == "__main__": client = Client("wss://stream.binance.com:9443/ws", 5)
关键修改点
- 将心跳消息改为JSON格式的
{"method": "PING", "id": 1},完全符合Binance API规范 - 把
keep_alive改为协程函数,处理发送心跳时的异常情况 - 优化连接失败后的重试逻辑,避免频繁重试导致服务器限制
- 增加PONG消息的处理,确认心跳交互正常
- 替换
exit()为重新连接逻辑,保证程序持续运行
内容的提问来源于stack exchange,提问作者chiwal
相关产品推荐
相关产品推荐

