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

排查Bitfinex Websocket交易数据间歇性延迟与缺失问题

Bitfinex WebSocket交易数据采集脚本间歇性断连+数据滞后问题排查与自动重启方案

问题背景

你提到用Python脚本通过WebSocket采集Bitfinex的BTC/ETH交易数据,前几天运行正常,但之后出现间歇性数小时数据缺失、消息时间戳滞后的问题,重启后暂时恢复但会复发,其他交易所脚本无异常,部署在AWS t2.medium的tmux环境中,同时运行10个其他脚本。下面是针对这个问题的排查建议和自动重启的实现方案:

一、问题排查方向

1. WebSocket连接稳定性问题

Bitfinex的WebSocket服务器可能会因为长时间无交互主动断开连接,而你的脚本目前没有处理重连和心跳逻辑:

  • 添加心跳机制:Bitfinex要求客户端定期发送{"event": "ping"}心跳包(建议每25-30秒一次),服务器会返回pong响应。如果超过1-2分钟没收到pong或任何消息,说明连接已失效,需要主动断开重连。
  • 细化日志记录:用logging模块代替print,记录每次连接建立、断开、消息接收、心跳交互的时间戳和细节,这样能精准定位断连发生的时间点和频率,方便排查是网络波动还是服务器资源问题。
  • 重连逻辑优化:当前脚本在连接断开后会自动重启,但没有限制重连间隔,短时间内频繁重连可能会触发交易所的限流机制,建议添加重连间隔(比如30秒)和失败次数限制。

2. 数据处理逻辑的性能瓶颈

你的脚本维护了全局的trades列表,每次收到消息都重新生成DataFrame并只保留最后一行,这会导致内存占用随着运行时间增加而上升,最终可能引发脚本卡顿:

  • 直接处理单条交易数据:不要维护大列表,收到te消息后直接将单条交易数据生成单行DataFrame,然后追加到CSV文件,减少内存消耗。
  • 优化CSV写入逻辑:当前判断文件是否存在的方式没问题,但可以提前定义列名,避免每次都依赖DataFrame自动生成,减少不必要的计算。

3. 服务器资源限制

AWS t2.medium是突发性能实例,当CPU积分耗尽时会降速,可能导致脚本无法及时处理WebSocket消息:

  • 监控服务器资源:用htop或top查看CPU、内存使用率,尤其是脚本运行时的资源占用情况,确认是否因为资源不足导致消息处理滞后。
  • 调整实例类型:如果CPU积分经常耗尽,可以考虑升级到t3.medium(无CPU积分限制),或者将Bitfinex的脚本单独分配更多资源。
  • 检查其他脚本影响:同时运行的10个其他脚本可能占用了大量资源,导致Bitfinex脚本无法及时响应WebSocket消息,可以尝试暂停部分脚本测试是否恢复正常。

4. Bitfinex API特定限制

  • 限流检查:虽然WebSocket的限流比较宽松,但如果短时间内发送过多订阅请求可能会被限制,不过你的脚本只订阅了一个交易对,这种概率较低,但可以查看Bitfinex官方文档确认限流规则。
  • 消息格式验证:打印完整的WebSocket消息,确认是否有异常消息被忽略,比如交易所可能发送的系统通知或格式变化。

二、自动重启实现方案

1. 脚本内集成自动重连与重启逻辑

修改脚本,添加心跳检测、超时重连和崩溃自动重启功能,以下是优化后的代码:

import websocket
import pandas as pd
import json
import time
import datetime
import os
import logging

# 配置日志,记录到文件
logging.basicConfig(
    filename='bitfinex_btcusd_trades.log',
    level=logging.INFO,
    format='%(asctime)s - %(levelname)s - %(message)s'
)

# 配置参数
CSV_FOLDER = r'/realtimedata/trades/bitfinex/btcusd/bitfinex_btcusd_trades_'
TRADE_COLUMNS = ['id', 'time', 'amount', 'price']
HEARTBEAT_INTERVAL = 25  # 每25秒发送一次心跳
RECONNECT_WAIT = 30       # 断连后等待30秒重连
NO_MESSAGE_TIMEOUT = 120  # 120秒没收到消息则重连

def send_heartbeat(ws):
    """向Bitfinex服务器发送心跳包"""
    try:
        ws.send(json.dumps({"event": "ping"}))
        logging.info("Sent heartbeat ping")
    except Exception as e:
        logging.error(f"Failed to send heartbeat: {str(e)}")

def on_message(ws, message):
    global last_msg_time
    last_msg_time = time.time()
    msg = json.loads(message)
    
    # 处理服务器返回的pong响应
    if isinstance(msg, dict) and msg.get('event') == 'pong':
        logging.info("Received pong response")
        return
    
    # 处理交易执行消息(te)
    if isinstance(msg, list) and msg[1] == 'te':
        trade_data = msg[2]
        logging.info(f"Received trade: ID={trade_data[0]}, Price={trade_data[3]}")
        
        # 生成单行DataFrame并写入CSV
        trade_df = pd.DataFrame([trade_data], columns=TRADE_COLUMNS)
        csv_path = f"{CSV_FOLDER}{datetime.datetime.today().strftime('%Y_%m_%d')}.csv"
        
        if not os.path.exists(csv_path):
            trade_df.to_csv(csv_path, header=TRADE_COLUMNS, index=False)
        else:
            trade_df.to_csv(csv_path, mode='a', header=False, index=False)

def on_error(ws, error):
    logging.error(f"WebSocket error occurred: {str(error)}")

def on_close(ws):
    logging.info("WebSocket connection closed")

def on_open(ws):
    # 订阅BTCUSD交易频道
    ws.send(json.dumps({"event": "subscribe", "channel": "trades", "pair": "BTCUSD"}))
    logging.info("Successfully subscribed to BTCUSD trades channel")

def run_websocket_client():
    global last_msg_time
    last_msg_time = time.time()
    
    ws = websocket.WebSocketApp(
        "wss://api.bitfinex.com/ws/2",
        on_message=on_message,
        on_error=on_error,
        on_close=on_close
    )
    ws.on_open = on_open
    
    # 运行客户端,同时处理心跳和超时检查
    while True:
        # 启动WebSocket客户端,设置心跳间隔
        ws.run_forever(ping_interval=HEARTBEAT_INTERVAL)
        
        # 连接断开后,检查是否超时无消息
        if time.time() - last_msg_time > NO_MESSAGE_TIMEOUT:
            logging.warning("No messages received for over 2 minutes, triggering reconnection")
        
        logging.info(f"Reconnecting in {RECONNECT_WAIT} seconds...")
        time.sleep(RECONNECT_WAIT)

if __name__ == "__main__":
    # 脚本崩溃后自动重启
    while True:
        try:
            run_websocket_client()
        except Exception as e:
            logging.critical(f"Script crashed with error: {str(e)}, restarting in 60 seconds...")
            time.sleep(60)

2. 系统级定时重启(Cron)

如果脚本内的重连机制不够稳定,可以用Linux的Cron定时每天重启脚本:

  1. 创建启动脚本start_bitfinex_trades.sh:
#!/bin/bash
# 终止旧进程
pkill -f "bitfinex_btcusd_trades.py"
# 在tmux会话中启动新脚本
tmux new-session -d -s bitfinex_btcusd "python3 /path/to/your/bitfinex_btcusd_trades.py"
  1. 赋予执行权限:chmod +x start_bitfinex_trades.sh
  2. 编辑Crontab:crontab -e,添加每天凌晨2点重启的任务:
0 2 * * * /path/to/start_bitfinex_trades.sh >> /path/to/cron_bitfinex.log 2>&1

3. 用Supervisor监控进程(推荐)

Supervisor是一个专业的进程管理工具,能自动监控脚本状态,崩溃时自动重启:

  1. 安装Supervisor(Ubuntu/Debian):
sudo apt-get update && sudo apt-get install supervisor
  1. 创建配置文件/etc/supervisor/conf.d/bitfinex_trades.conf:
[program:bitfinex_btcusd_trades]
command=python3 /path/to/your/bitfinex_btcusd_trades.py
directory=/path/to/your/script/directory
user=ubuntu  # 替换为你的服务器用户名
autostart=true
autorestart=true  # 崩溃时自动重启
stderr_logfile=/var/log/bitfinex_trades.err.log
stdout_logfile=/var/log/bitfinex_trades.out.log
  1. 重新加载配置并启动服务:
sudo supervisorctl reread
sudo supervisorctl update
sudo supervisorctl start bitfinex_btcusd_trades
  1. 查看进程状态:sudo supervisorctl status bitfinex_btcusd_trades

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 04:24:14