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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 13:06:17