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
相关产品推荐
相关产品推荐

