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

Python Asyncio Polygon WebSocket流数据运行数分钟后报错求助

Polygon WebSocket期权流数据报错排查与修复

核心问题分析

从代码和运行数分钟后报错的现象来看,主要问题集中在文件资源泄漏和WebSocket连接稳定性两方面:

1. 文件句柄泄漏

每次收到消息都重复打开文件但未关闭,系统可用文件句柄耗尽后会抛出IO异常:

  • 原代码中open("PolygonQuote1.json", "w")和open("PolygonTrade1.json", "w")未调用close(),也未用with语句自动管理资源。
  • 频繁打开写入会积累大量未释放的文件句柄,运行一段时间后触发系统级资源限制报错。

2. 未处理WebSocket连接异常

Polygon WebSocket会定期发送心跳包,若客户端未正确响应或网络波动,连接可能断开,原代码没有重连机制和异常捕获逻辑。

3. 消息处理逻辑缺陷

只取msgs列表的最后一条消息,会丢失批量推送的历史消息,虽不是报错直接原因,但会导致数据不完整。


修复后的代码

from polygon import WebSocketClient
from polygon.websocket.models import WebSocketMessage
from typing import List
import json
import asyncio
import nest_asyncio

nest_asyncio.apply()

# 用模块级变量替代全局变量,更规范
quote_list = []
trade_list = []

def handle_msg(msgs: List[WebSocketMessage]):
    # 遍历所有消息,避免丢数据
    for m in msgs:
        print(m)
        if m.event_type == 'Q':
            quote_list.append({
                'TimeStamp': m.timestamp, 
                'BidPrice': m.bid_price, 
                'BidSize': m.bid_size, 
                'AskPrice': m.ask_price, 
                'AskSize': m.ask_size, 
                'Symbol': m.symbol
            })
        elif m.event_type == 'T':
            trade_list.append({
                'TimeStamp': m.timestamp,
                'Price': m.price, 
                'Size': m.size, 
                'Symbol': m.symbol
            })
    
    # 使用with语句自动管理文件资源,避免句柄泄漏
    with open("PolygonQuote1.json", "w") as out_file1:
        json.dump(quote_list, out_file1, indent=4)
    with open("PolygonTrade1.json", "w") as out_file2:
        json.dump(trade_list, out_file2, indent=4)

async def run_websocket():
    ws = WebSocketClient(
        api_key="YOUR_API_KEY", 
        feed='socket.polygon.io', 
        market='options',
        subscriptions=[
            "T.O:SPY230210C00408000", "Q.O:SPY230210C00408000", 
            "T.O:SPY230210P00407000", "Q.O:SPY230210P00407000"
        ]
    )
    try:
        await ws.run(handle_msg)
    except Exception as e:
        print(f"连接异常: {e}")
        # 捕获异常后延迟重连
        await asyncio.sleep(5)
        await run_websocket()

if __name__ == "__main__":
    asyncio.run(run_websocket())

额外优化建议

  • 减少文件写入频率:每次消息都写入会导致IO性能瓶颈,可设置定时批量写入(比如每30秒)或累计一定数量消息后再写入。
  • 使用线程安全结构:如果后续扩展多线程处理,建议用queue.Queue替代普通列表,避免数据竞争。
  • 替换打印为日志:用logging模块替代print,方便排查长期运行的问题。
  • 验证API权限:确保你的Polygon API密钥拥有期权数据流的访问权限,权限不足也可能导致连接中断。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 09:21:35