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

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()

关键修复点说明

  1. 完整的异常捕获:把消息接收、处理、数据库操作都放到try-except块中,专门捕获WebSocketConnectionClosedException来处理连接断开的情况。
  2. 自动重连逻辑:当检测到连接断开时,重置连接对象,循环尝试重新建立连接。
  3. 数据类型转换:将WebSocket返回的字符串数值转换为REAL或INTEGER类型,匹配SQLite表的字段类型,避免插入错误。
  4. 移除eval:用安全的字典get方法替代eval,提升代码安全性和可读性。
  5. 持久化数据库:把内存数据库改为文件数据库gdax_ticker.db,避免程序重启后数据丢失。
  6. 模块化拆分:把连接、解析、插入操作拆分为单独的方法,代码更清晰易维护。

内容的提问来源于stack exchange,提问作者TRV

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 09:35:55