如何解决GDAX WebSocket流数据存入Pandas DataFrame失败问题?
问题分析与解决方案
首先得点出你代码里的几个核心问题:
wsClient.start()并非可迭代对象:这个方法只是启动WebSocket客户端的后台线程,不会返回任何可迭代的数据,所以for i in wsClient.start()这个循环根本不会执行,自然拿不到数据。df.append()不会原地修改DataFrame:Pandas的append()方法是返回一个新的DataFrame实例,你如果不把结果重新赋值给df1,原DataFrame不会有任何变化。- WebSocket是异步运行的:客户端在独立线程推送数据,你没法用普通的for循环去同步“抓取”这些实时数据。
正确的实现方式
你需要重写GDAX WebsocketClient的on_message方法,每当有新数据推送过来时,就把它添加到DataFrame中。另外要注意线程安全,因为WebSocket的回调是在单独线程执行的,我们可以用锁来避免并发操作导致的问题。
import pandas as pd from gdax import WebsocketClient import threading # 初始化DataFrame和线程锁 df1 = pd.DataFrame(columns=['type', 'price', 'size', 'product_id']) data_lock = threading.Lock() class CustomWebsocketClient(WebsocketClient): def on_message(self, msg): # 只处理我们需要的ticker数据(可根据GDAX消息类型调整) if msg['type'] == 'ticker': global df1, data_lock # 加锁避免多线程操作DataFrame时的冲突 with data_lock: # 创建单行数据的DataFrame,用concat更新原DataFrame new_row = pd.DataFrame([{ 'type': msg['type'], 'price': float(msg['price']), 'size': float(msg['size']), 'product_id': msg['product_id'] }]) df1 = pd.concat([df1, new_row], ignore_index=True) # 启动自定义的WebSocket客户端 wsClient = CustomWebsocketClient(url="wss://ws-feed.gdax.com", products="LTC-USD") wsClient.start() # 可选:在主线程定时打印最新数据,验证是否成功写入 import time while True: with data_lock: print("最新数据:") print(df1.tail(3)) time.sleep(3)
关键说明
- 我们自定义了
CustomWebsocketClient继承自官方的WebsocketClient,重写on_message方法来处理每条推送的消息。 - 使用
threading.Lock()保证DataFrame的读写安全,避免WebSocket线程和主线程同时操作DataFrame引发异常。 - 用
pd.concat()替代已被Pandas 2.0+弃用的append(),并且必须重新赋值给df1,才能完成DataFrame的更新。 - 最后的循环定时打印数据,方便你验证是否成功将流式数据存入DataFrame。
内容的提问来源于stack exchange,提问作者RustyShackleford
相关产品推荐
相关产品推荐

