如何使用多线程运行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
相关产品推荐
相关产品推荐

