Python多进程多线程处理BTC K线流数据报错求助
问题描述
我尝试从Web API流式获取BTC K线数据,每约2秒需对数据流进行一次数据计算,但计算耗时超过2秒,因此需要实现多进程来启动计算工作进程。作为编程新手,这超出了我的现有知识范围,恳请提供帮助。
当前代码
import websocket import json import pandas as pd import threading import multiprocessing import ccxt import time import os class BinanceKlineStream(): # kline subscription def __init__(self, binance): self.kline_sub = 'wss://stream.binance.com:9443/ws/btcusdt@kline_1m' self.binance = binance def thread_stream(self): worker_thread = threading.Thread(target=self.stream) worker_thread.start() def stream(self): self.ws = websocket.WebSocketApp(self.kline_sub, on_message=self.on_message, on_error=self.on_error) self.ws.run_forever(ping_interval=60) def on_message(self, ws, message): json_message = json.loads(message) self.kline_current = pd.DataFrame([[json_message['E'], json_message['k']['o'], json_message['k']['h'], json_message['k']['l'], json_message['k']['c'], json_message['k']['v'], json_message['k']['x'], json_message['k']['t'], json_message['k']['T']]], columns=['id', 'open', 'high', 'low', 'close', 'volume', 'kline_closed', 'start_time', 'close_time']) print('process1: ', os.getpid()) # spawn a new process when on_message is called worker_process = multiprocessing.Process(target=type(self).get_kline_combined) worker_process.start() worker_process.join() def on_error(self, ws, error): print(error) @staticmethod def get_kline_combined(): print('process2: ', os.getpid()) time.sleep(5) if __name__ == "__main__": # start streaming kline information binance_spot = ccxt.binance() stream = BinanceKlineStream(binance_spot) stream.thread_stream()
报错信息
Traceback (most recent call last): File "<string>", line 1, in <module> File "C:\Users\Ted Teng\AppData\Local\Programs\Python\Python310\lib\multiprocessing\spawn.py", line 116, in spawn_main exitcode = _main(fd, parent_sentinel) File "C:\Users\Ted Teng\AppData\Local\Programs\Python\Python310\lib\multiprocessing\spawn.py", line 126, in _main self = reduction.pickle.load(from_parent) AttributeError: Can't get attribute 'BinanceKlineStream.get_kline_combined' on <module '__main__' (built-in)>
解决方案
这个错误源于Windows系统多进程的spawn启动机制:子进程会重新导入主模块,但类内的静态方法在子进程中无法被正确序列化识别。另外你的代码还有两个问题:join()会阻塞主线程,导致下一条K线消息处理延迟;没有把最新K线数据传递给计算进程。
以下是修复后的代码:
import websocket import json import pandas as pd import threading import multiprocessing import ccxt import time import os # 将计算逻辑独立成函数,避免类方法序列化问题 def calculate_kline_data(kline_df): print(f'计算进程ID: {os.getpid()}') # 替换为你的实际计算逻辑,比如技术指标计算 time.sleep(5) # 模拟耗时计算 print(f'计算完成,数据摘要: {kline_df[["open", "close", "volume"]].iloc[0]}') class BinanceKlineStream(): def __init__(self, binance): self.kline_sub = 'wss://stream.binance.com:9443/ws/btcusdt@kline_1m' self.binance = binance def thread_stream(self): worker_thread = threading.Thread(target=self.stream) worker_thread.start() def stream(self): self.ws = websocket.WebSocketApp(self.kline_sub, on_message=self.on_message, on_error=self.on_error) self.ws.run_forever(ping_interval=60) def on_message(self, ws, message): json_message = json.loads(message) kline_current = pd.DataFrame([[json_message['E'], json_message['k']['o'], json_message['k']['h'], json_message['k']['l'], json_message['k']['c'], json_message['k']['v'], json_message['k']['x'], json_message['k']['t'], json_message['k']['T']]], columns=['id', 'open', 'high', 'low', 'close', 'volume', 'kline_closed', 'start_time', 'close_time']) print(f'主线程进程ID: {os.getpid()}') # 启动子进程并传递K线数据,不调用join()避免阻塞主线程 worker_process = multiprocessing.Process(target=calculate_kline_data, args=(kline_current,)) worker_process.start() # 设置守护进程,主进程退出时自动终止子进程 worker_process.daemon = True def on_error(self, ws, error): print(error) if __name__ == "__main__": # 所有启动逻辑必须放在该判断内,防止子进程重复执行初始化代码 binance_spot = ccxt.binance() stream = BinanceKlineStream(binance_spot) stream.thread_stream()
关键修复点:
- 把计算逻辑移到独立函数,解决类方法在子进程中的序列化问题
- 移除
worker_process.join(),让主线程可以持续处理新的K线消息 - 通过
args参数将最新K线数据传递给计算进程 - 设置
daemon=True,避免主进程退出后残留计算进程 - 确保所有启动代码都在
if __name__ == "__main__":块内,符合Windows多进程要求
内容的提问来源于stack exchange,提问作者Ted Teng
相关产品推荐
相关产品推荐

