如何借助asyncio与aiohttp从API取数后高效写入DataFrame到CSV
优化异步获取后DataFrame的CSV写入效率
核心问题分析
你当前迭代写入的低效主要来自重复的文件IO操作(打开、写入、关闭文件的开销),以及如果直接在异步函数里调用同步的to_csv,会阻塞事件循环,完全浪费异步特性。
优化方案
1. 合并DataFrame后批量写入(优先推荐,若业务允许)
如果不需要为每个id生成单独的CSV文件,直接合并所有DataFrame后一次性写入,能大幅减少IO次数:
import pandas as pd # 合并所有时序DataFrame(ignore_index重置索引) combined_df = pd.concat(df_list, ignore_index=True) # 一次性写入CSV combined_df.to_csv("all_timedata.csv", index=False)
2. 异步写入的正确实现(需单独文件时)
pandas的to_csv是同步阻塞操作,直接在async函数中调用会卡住事件循环。正确的做法是用线程池包装同步IO操作,让写入任务在后台线程执行,不占用异步事件循环:
import asyncio import pandas as pd from concurrent.futures import ThreadPoolExecutor async def write_single_csv(df, filename): loop = asyncio.get_running_loop() # 用线程池执行同步的to_csv,避免阻塞事件循环 await loop.run_in_executor( None, # 默认使用全局线程池 df.to_csv, filename, False # 不写入索引 ) async def batch_write_csv(df_list, filenames): # 创建所有写入任务,并发执行 tasks = [ write_single_csv(df, fname) for df, fname in zip(df_list, filenames) ] await asyncio.gather(*tasks) # 调用示例:假设df_list是你的DataFrame列表,文件名对应每个id # asyncio.run(batch_write_csv(df_list, [f"data_{id}.csv" for id in id_list]))
- 这里用
run_in_executor把同步IO丢到线程池,异步事件循环可以继续处理其他任务,真正利用了异步的并发特性。 - 可以自定义线程池大小(比如
ThreadPoolExecutor(max_workers=4)),避免线程过多导致调度开销。
3. 额外优化细节
- 调整
to_csv的buffering参数,增大缓冲区大小(比如buffering=1024*1024),减少磁盘写入次数。 - 如果是大量小文件,可考虑先写入一个带id列的总文件,后续用pandas按id拆分(适合后续需要批量处理的场景)。
- 避免在写入时做额外的数据转换,尽量在API响应转DataFrame时完成预处理。
内容的提问来源于stack exchange,提问作者brokkoo
相关产品推荐
相关产品推荐

