如何优化Pandas DataFrame循环中的API调用性能?
问题描述
我有一个名为df_flights的DataFrame,每行代表一次API调用的输入数据,数据集最大可达1000行。需要提取每行各列的值发起API调用,将返回结果整合为新的DataFrame。
当前实现的简化版本如下:
# 初始化空DataFrame df_bereik = pd.DataFrame() # 遍历df_flights的每一行 for index, row in df_flights.iterrows(): # 使用当前行的值发起API调用 df = get_new_df(index, row, variable, token) # 将结果合并到df_bereik df_bereik = pd.concat([df_bereik, df], ignore_index=True) def get_new_df(index, row, target_group, token): advertiser = row["advertiser"] product = row["product"] # 其他字段提取... df = call_api(variable, advertiser, product) # 清洗API返回的数据 df = clean_up(df, advertiser, product, variable, ...) return df def clean_up(df): # 数据清洗逻辑 return df
目前因使用iterrows导致运行速度极慢,求更快的替代方案。
优化方案
1. 避免循环内频繁执行concat
每次concat都会创建新的DataFrame,带来额外内存开销与时间消耗。改为先将所有API返回的DataFrame存入列表,最后一次性合并:
df_list = [] for index, row in df_flights.iterrows(): df = get_new_df(index, row, variable, token) df_list.append(df) # 最后一次性合并所有结果 df_bereik = pd.concat(df_list, ignore_index=True)
2. 用itertuples替代iterrows
iterrows返回Series对象,字段访问开销较高;itertuples返回命名元组,访问速度更快。修改遍历逻辑:
df_list = [] for row in df_flights.itertuples(index=True): # 通过属性名访问字段,如row.advertiser df = get_new_df(row.Index, row, variable, token) df_list.append(df) df_bereik = pd.concat(df_list, ignore_index=True) # 对应调整get_new_df的字段提取逻辑 def get_new_df(index, row, target_group, token): advertiser = row.advertiser product = row.product # 其他字段提取... # 后续API调用与清洗逻辑不变
3. 并行调用API(核心优化)
API调用属于IO密集型任务,串行等待会浪费大量时间。使用多线程或异步IO可同时发起多个请求,大幅缩短总耗时。
多线程实现(基于concurrent.futures.ThreadPoolExecutor)
from concurrent.futures import ThreadPoolExecutor def process_row(row_tuple): index = row_tuple.Index row = row_tuple return get_new_df(index, row, variable, token) # 线程数可根据API并发限制调整,建议设为10-20 with ThreadPoolExecutor(max_workers=15) as executor: # 批量提交所有行的处理任务 df_list = list(executor.map(process_row, df_flights.itertuples(index=True))) df_bereik = pd.concat(df_list, ignore_index=True)
异步IO实现(适合高并发场景)
若API支持异步请求,用aiohttp配合asyncio效率更高:
import asyncio import aiohttp # 将同步API调用改为异步版本 async def async_call_api(variable, advertiser, product): async with aiohttp.ClientSession() as session: # 替换为实际API请求逻辑 async with session.get(f"your_api_endpoint?var={variable}&adv={advertiser}&prod={product}") as resp: data = await resp.json() return pd.DataFrame(data) # 同步的get_new_df改为异步版本 async def async_get_new_df(index, row, target_group, token): advertiser = row.advertiser product = row.product df = await async_call_api(variable, advertiser, product) df = clean_up(df, advertiser, product, variable, ...) return df async def main(): tasks = [async_get_new_df(row.Index, row, variable, token) for row in df_flights.itertuples(index=True)] df_list = await asyncio.gather(*tasks) return pd.concat(df_list, ignore_index=True) df_bereik = asyncio.run(main())
4. 其他细节优化
- 检查
clean_up函数是否可通过向量化操作优化,避免在每行清洗逻辑中使用循环; - 若API提供批量请求接口,优先使用批量调用减少请求次数(例如一次传入10行参数,返回10组结果);
- 为API请求添加超时重试机制,避免个别请求失败导致整个流程中断。
内容的提问来源于stack exchange,提问作者Jordy
相关产品推荐
相关产品推荐

