如何让Pandas apply()异步调用API 加速14000行数据处理
优化Bloomberg API请求效率的方案
最优方案:批量请求(无需并行)
Bloomberg的bdp接口本身支持批量传入多个ticker,这比单条请求效率高几个数量级,完全没必要逐个调用。直接把所有ISIN转换成/isin/{ISIN}格式后一次性请求:
# 处理所有ISIN为Bloomberg格式 tickers = df['ISIN'].apply(lambda x: f'/isin/{x}').tolist() # 批量请求数据 result = blp.bdp(tickers=tickers, flds=['market_sector_des']) # 将结果合并回原DataFrame df['new_column'] = result['market_sector_des'].values
这种方式只需要1次API请求,耗时会从8小时压缩到几秒到几十秒(取决于Bloomberg服务器响应),是最推荐的方案。
并行请求方案(IO密集型场景)
如果因为特殊限制无法批量请求,可使用Python的concurrent.futures.ThreadPoolExecutor做异步并行请求(API调用是IO密集型,线程池比进程池更高效):
步骤1:改造原函数为线程安全版本
def market_sector_des(isin): isin = f'/isin/{isin}' # xbbg的blp对象是线程安全的,无需额外处理 return blp.bdp(tickers=isin, flds=['market_sector_des']).iloc[0, 0]
步骤2:用线程池并行执行
from concurrent.futures import ThreadPoolExecutor # 设置线程数(根据Bloomberg API限制调整,建议10-20,避免触发限流) max_workers = 15 with ThreadPoolExecutor(max_workers=max_workers) as executor: # 并行处理所有ISIN df['new_column'] = list(executor.map(market_sector_des, df['ISIN'].tolist()))
这种方式能把耗时从8小时降到约30分钟(14000/15≈933次并行循环,每次2秒的话就是933*2≈1866秒≈31分钟),注意不要设置过大的线程数,避免被Bloomberg限流。
Dask问题排查与修复
你提到用Dask遇到问题,大概率是因为blp对象无法在Dask分布式进程中序列化,或者任务调度与API上下文冲突。如果一定要用Dask,需要在每个任务中重新初始化blp连接:
import dask.dataframe as dd def dask_market_sector_des(isin): # 每个任务重新初始化blp连接 from xbbg import blp isin = f'/isin/{isin}' return blp.bdp(tickers=isin, flds=['market_sector_des']).iloc[0, 0] # 将DataFrame转为Dask DataFrame ddf = dd.from_pandas(df, npartitions=10) # 并行计算 ddf['new_column'] = ddf['ISIN'].apply(dask_market_sector_des, meta=('new_column', 'object')) # 转回Pandas DataFrame df = ddf.compute()
但这种方式效率不如批量请求或线程池,仅作为备选。
内容的提问来源于stack exchange,提问作者user21584139
相关产品推荐
相关产品推荐

