GDAX Websocket连接约2小时后自动断开,配置Ping仍失效
问题分析与修复方案
我帮你梳理下代码里的几个核心问题,这些应该就是导致程序运行2-2.5小时后连接断开且无法自动恢复的原因:
1. 未捕获WebSocket接收时的连接异常
你的代码中dane=json.loads(self.ws.recv())是在try-except块之外执行的。当WebSocket连接被远程关闭(比如GDAX服务器超时、网络波动),recv()会直接抛出连接关闭类的异常,而这个异常没有被捕获,导致程序直接崩溃。你需要把接收消息、处理消息、数据库操作的整个流程都包裹在异常捕获块中。
2. 心跳机制不完善且无重连逻辑
虽然你实现了ping发送,但:
- 没有处理ping失败的情况(比如连接已经断开时,
self.ws.ping()会抛出异常) - 核心的消息接收逻辑不在异常捕获范围内,一旦连接断开无法触发重连
- 缺少连接断开后的自动重连机制,程序断开后只能手动重启
3. 不安全且冗余的eval使用
代码里的eval("dane.get('" + b[i] + "', None)")完全没必要,直接用dane.get(b[i], None)就能实现相同功能,还能避免eval带来的安全风险和潜在语法错误。
4. 数据库操作未做异常处理
你定义的SQLite表中很多字段是INTEGER类型,但从GDAX WebSocket拿到的price、best_bid等数据是字符串格式,直接插入会导致类型不匹配错误,而且这个错误没有被捕获,会导致程序崩溃。
修改后的代码示例
下面是修复了上述问题的精简代码,重点加入了完整的异常捕获、自动重连、数据类型转换和安全的字典取值:
import time import json import sqlite3 from websocket import create_connection, WebSocketConnectionClosedException # 初始化持久化数据库(替换内存数据库,避免重启丢失数据) conn = sqlite3.connect("gdax_ticker.db") c = conn.cursor() c.execute("""CREATE TABLE IF NOT EXISTS `Ticker` ( `ID` INTEGER PRIMARY KEY AUTOINCREMENT, `Type` TEXT, `Sequence` INTEGER, `Product_id` TEXT, `Price` REAL, `Open_24h` REAL, `Volume_24h` REAL, `Low_24h` REAL, `High_24h` REAL, `Volume_30d` REAL, `Best_bid` REAL, `Best_ask` REAL, `side` TEXT, `time` TEXT, `trade_id` INTEGER, `last_size` REAL )""") conn.commit() class Websocket(): def __init__(self, wsurl="wss://ws-feed.gdax.com", agi1=2, produkty=['ETH-EUR', 'LTC-EUR', 'BTC-EUR']): self.wsurl = wsurl self.ws = None self.agi1 = agi1 self.ping_start = time.time() self.produkty = produkty self.newdict = {} print(f"连接地址: {self.wsurl}") print(f"订阅产品: {self.produkty}") def _parse_ticker_data(self, dane): """安全解析ticker数据并转换数值类型""" fields = [ 'type', 'sequence', 'product_id', 'price', 'open_24h', 'volume_24h', 'low_24h', 'high_24h', 'volume_30d', 'best_bid', 'best_ask', 'side', 'time', 'trade_id', 'last_size' ] for field in fields: value = dane.get(field, None) # 针对数值字段做类型转换 if field in ['price', 'open_24h', 'volume_24h', 'low_24h', 'high_24h', 'volume_30d', 'best_bid', 'best_ask', 'last_size']: if value is not None: try: value = float(value) except ValueError: value = None elif field in ['sequence', 'trade_id']: if value is not None: try: value = int(value) except ValueError: value = None self.newdict[field] = value def _insert_to_db(self): """插入数据到数据库,带异常处理""" try: c.execute("""INSERT INTO Ticker( type, sequence, product_id, price, open_24h, volume_24h, low_24h, high_24h, volume_30d, best_bid, best_ask, side, time, trade_id, last_size ) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", ( self.newdict['type'], self.newdict['sequence'], self.newdict['product_id'], self.newdict['price'], self.newdict['open_24h'], self.newdict['volume_24h'], self.newdict['low_24h'], self.newdict['high_24h'], self.newdict['volume_30d'], self.newdict['best_bid'], self.newdict['best_ask'], self.newdict['side'], self.newdict['time'], self.newdict['trade_id'], self.newdict['last_size'] )) conn.commit() except sqlite3.Error as e: print(f"{time.ctime()} 数据库插入错误: {e}") conn.rollback() def _connect_websocket(self): """建立WebSocket连接并发送订阅请求""" try: self.ws = create_connection(self.wsurl) # 确定订阅频道 if self.agi1 == 0: kanaly = None elif self.agi1 == 1: kanaly = "heartbeat" elif self.agi1 == 2: kanaly = ["ticker"] elif self.agi1 == 3: kanaly = "level2" else: print(f'ERROR: 无效参数! (arg1)= {self.agi1}') return False # 构造并发送订阅消息 subscribe_msg = {'type': 'subscribe', 'product_ids': self.produkty} if kanaly is not None: subscribe_msg['channels'] = [{"name": kanaly, 'product_ids': self.produkty}] self.ws.send(json.dumps(subscribe_msg)) self.ping_start = time.time() print(f"{time.ctime()} WebSocket连接成功") return True except Exception as e: print(f"{time.ctime()} 连接失败: {e}") return False def run(self): """主运行逻辑,包含自动重连""" while True: # 检测连接状态,断开则尝试重连 if self.ws is None or not self.ws.connected: if not self._connect_websocket(): print(f"{time.ctime()} 等待5秒后重试连接...") time.sleep(5) continue try: # 每20秒发送一次心跳ping if (time.time() - self.ping_start) >= 20: self.ws.ping("ping") print(f'{time.ctime()} 发送ping') self.ping_start = time.time() # 接收并处理消息 dane = json.loads(self.ws.recv()) typtranzakcji = dane.get('type', None) if typtranzakcji == 'error': print(f"{time.ctime()} 收到错误消息: {json.dumps(dane, indent=4)}") elif typtranzakcji == 'ticker' and self.agi1 == 2: self._parse_ticker_data(dane) self._insert_to_db() time.sleep(1) except WebSocketConnectionClosedException: print(f"{time.ctime()} WebSocket连接已关闭,准备重连") self.ws = None except json.JSONDecodeError: print(f"{time.ctime()} 消息解析失败,跳过该消息") except Exception as e: print(f"{time.ctime()} 未知错误: {e}") self.ws = None time.sleep(5) class Autotrader(): def __init__(self): self.produkty = ["BTC-EUR"] self.webs = Websocket(produkty=self.produkty) self.webs.run() if __name__ == "__main__": trader = Autotrader()
关键修复点说明
- 完整的异常捕获:把消息接收、处理、数据库操作都放到
try-except块中,专门捕获WebSocketConnectionClosedException来处理连接断开的情况。 - 自动重连逻辑:当检测到连接断开时,重置连接对象,循环尝试重新建立连接。
- 数据类型转换:将WebSocket返回的字符串数值转换为
REAL或INTEGER类型,匹配SQLite表的字段类型,避免插入错误。 - 移除
eval:用安全的字典get方法替代eval,提升代码安全性和可读性。 - 持久化数据库:把内存数据库改为文件数据库
gdax_ticker.db,避免程序重启后数据丢失。 - 模块化拆分:把连接、解析、插入操作拆分为单独的方法,代码更清晰易维护。
内容的提问来源于stack exchange,提问作者TRV
相关产品推荐
相关产品推荐

