Binance websocket行情跨脚本生成pandas DataFrame报错如何解决?
错误根因
- 进程内存隔离:第一个脚本运行时存储在numpy数组中的数据属于该进程的私有内存,独立启动的第二个脚本无法直接访问该内存空间
- DataFrame构造参数不合法:直接传入存储了JSON字符串/嵌套字典的numpy数组,不符合pandas.DataFrame的入参要求,触发ValueError
实现方案(新手友好版)
选择JSON Lines格式的本地文件作为跨进程数据中转介质,实现成本最低,无需额外依赖。
第一个脚本(Websocket接收端,后台运行)
原有对接Binance 24hrMiniTicker推送的逻辑无需大幅修改,收到消息后追加写入本地文件即可,无需暂存到numpy数组:
import json import websocket # Binance Websocket 24hrMiniTicker地址 BINANCE_WS_URL = "wss://stream.binance.com:9443/ws/!miniTicker@arr" def on_message(ws, message): ticker_list = json.loads(message) for ticker in ticker_list: # 过滤ADAUSDT行情 if ticker.get("s") == "ADAUSDT": # 追加写入JSONL文件,每行对应一条行情数据 with open("ada_usdt_24h_ticker.jsonl", "a", encoding="utf-8") as f: f.write(json.dumps(ticker) + "\n") def on_error(ws, error): print(error) def on_close(ws, close_status_code, close_msg): print("Websocket连接关闭") def on_open(ws): print("Websocket连接已建立,开始接收行情数据") if __name__ == "__main__": ws = websocket.WebSocketApp(BINANCE_WS_URL, on_open=on_open, on_message=on_message, on_error=on_error, on_close=on_close) ws.run_forever()
后台运行该脚本即可持续接收数据并写入本地文件。
第二个脚本(数据处理端)
按行读取JSONL文件转换为字典列表后,再传入DataFrame构造函数即可避免报错:
import pandas as pd import json data_list = [] with open("ada_usdt_24h_ticker.jsonl", "r", encoding="utf-8") as f: for line in f: line = line.strip() if not line: continue data_list.append(json.loads(line)) # 传入字典列表构造DataFrame df = pd.DataFrame(data_list) # 后续处理逻辑可直接基于df实现
可选优化方案
- 数据量较大时可改用Sqlite作为存储介质,第一个脚本将数据写入Sqlite表,第二个脚本直接调用
pd.read_sql()读取数据生成DataFrame,读写效率更高 - 要求低延迟共享数据时可使用multiprocessing模块的共享内存/队列实现跨进程数据传输,无需落地磁盘,但实现复杂度相对更高
内容的提问来源于stack exchange,提问作者Renaat Vandewiele
相关产品推荐
相关产品推荐

