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

如何借助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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 19:31:23