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

如何使用多线程运行for循环加速加密货币历史数据批量获取

加密货币历史数据拉取多线程改造方案

改造核心思路

你的场景属于典型的IO密集型任务,大部分耗时是等待接口返回数据,用多线程并行发起请求可以极大压缩整体耗时。我们用Python标准库的concurrent.futures.ThreadPoolExecutor实现线程池,既可以控制并发数避免触发交易所API限流,也不用手动管理线程生命周期。

具体改造实现

第一步:导入依赖

在代码头部新增导入:

import concurrent.futures
import pandas as pd
# 你原来的client导入逻辑保留

第二步:改造crypto类

把单个交易对+单个时间粒度的拉取处理逻辑抽成独立方法,再用线程池并行执行:

class crypto:
    symbols = []
    with open('allsymbols', 'r+') as f:
        for line in f:
            symbol = line.strip('\n') + 'BTC'
            symbols.append(symbol)

    intervals = ['1m','5m']
    # 并发数,可根据API限流规则调整,建议初始10-20
    MAX_WORKERS = 15

    @staticmethod
    def _fetch_single_task(symbol, interval):
        """单个拉取任务:独立处理单个交易对+单个时间粒度的数据拉取和预处理"""
        historical_data = client.get_historical_klines(symbol, interval, '11/19/2021', limit=1000)[-22:-1]
        historical_df = pd.DataFrame(historical_data)
        historical_df.columns = ['time', '1', 'high', 'low', 'close', '5', '6', '7', '8', '9', '10', '11']
        historical_df.drop(columns=['1', '5', '6', '7', '8', '9', '10', '11'], axis=1, inplace=True)
        historical_df['interval'] = interval
        historical_df['symbol'] = symbol
        historical_df[['high', 'low', 'close']] = historical_df[['high', 'low', 'close']].apply(pd.to_numeric, axis=1)
        historical_df['time'] = pd.to_datetime(historical_df['time'] / 1000, unit='s')
        return historical_df.to_dict()

    @staticmethod
    def historical_data():
        historical_list = []
        # 生成所有需要执行的任务参数
        task_params = [(symbol, interval) for symbol in crypto.symbols for interval in crypto.intervals]
        
        # 线程池并发执行任务
        with concurrent.futures.ThreadPoolExecutor(max_workers=crypto.MAX_WORKERS) as executor:
            # 提交所有任务,按顺序收集返回结果
            results = executor.map(lambda x: crypto._fetch_single_task(*x), task_params)
        
        # 把结果转成列表,格式和原来完全一致
        for res in results:
            historical_list.append(res)
        return historical_list

    @staticmethod
    def refactor_list(historical_list):
        # 原有逻辑完全保留,不用修改
        historical_list_refactored = []
        for i in range(len(historical_list)):
            single_key_data = historical_list[i]
            single_key_data['high'] = list(single_key_data['high'].values())
            single_key_data['low'] = list(single_key_data['low'].values())
            single_key_data['close'] = list(single_key_data['close'].values())
            single_key_data['interval'] = list(single_key_data['interval'].values())
            single_key_data['symbol'] = list(single_key_data['symbol'].values())
            single_key_data['time'] = list(single_key_data['time'].values())
            historical_list_refactored.append(single_key_data)
        return historical_list_refactored

注意事项

  • 调整MAX_WORKERS参数时不要过大,否则会触发交易所API限流,导致请求失败,建议初始值设为10~20,可根据实际请求返回情况调整
  • 可在_fetch_single_task方法中添加异常捕获和重试逻辑,避免个别接口请求失败导致整体任务中断
  • 如果你使用的接口client不是线程安全的,需要在_fetch_single_task方法内部初始化client实例,不要复用全局的client

内容的提问来源于stack exchange,提问作者Jellyfish

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 02:15:01