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

如何让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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 03:15:14