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

如何正确使用asyncio结合Pandas读取合并多CSV文件以提升效率?

为什么你的异步CSV读取没提速?正确实现+Concat优化技巧

嘿,我看了你的代码,问题出在你以为用了asyncio就能并行,但实际上pd.read_csv是同步阻塞的函数——它不会让出asyncio的事件循环,所以你的异步代码本质上还是串行执行每个文件的读取,和同步版本耗时自然没差。下面给你拆解正确的异步实现,再聊聊怎么优化concat的速度。


一、正确的异步CSV读取方式

要让文件读取真正并行,我们需要把同步的IO操作放到线程池里执行(因为文件IO是操作系统层面支持并行的),这样asyncio的事件循环可以同时调度多个读取任务,不会被单个read_csv卡住。

方法1:Python 3.9+ 用asyncio.to_thread(最简单)

这个方法是Python3.9新增的,可以直接把同步函数丢到线程池,不用自己管理线程:

import asyncio
import time
import pandas as pd
import glob2
import os

async def read_csv(filename):
    # 把同步的read_csv委托给线程池执行,让出事件循环
    return await asyncio.to_thread(pd.read_csv, filename, header=None)

async def main():
    path = r'C:\LRM_STGY_REPO\IB_IN'
    file_list = list(glob2.iglob(os.path.join(path, "*.txt")))
    # 创建所有读取任务,并发执行
    dfs = await asyncio.gather(*[read_csv(f) for f in file_list])
    # 合并结果
    df = pd.concat(dfs, ignore_index=True)
    return df

if __name__ == "__main__":
    t = time.time()
    df = asyncio.run(main())
    print(f"异步版本耗时: {time.time()-t:.4f}")

    # 同步版本对比
    t = time.time()
    df_sync = pd.concat(map(lambda f: pd.read_csv(f, header=None), glob2.iglob(os.path.join(path, "*.txt"))), ignore_index=True)
    print(f"同步版本耗时: {time.time()-t:.4f}")

方法2:Python3.8及以下 用ThreadPoolExecutor

如果你的Python版本较低,可以手动创建线程池配合asyncio:

from concurrent.futures import ThreadPoolExecutor
import asyncio
import time
import pandas as pd
import glob2
import os

async def read_csv(executor, filename):
    # 在指定的线程池中执行read_csv
    return await asyncio.get_event_loop().run_in_executor(executor, pd.read_csv, filename, header=None)

async def main():
    path = r'C:\LRM_STGY_REPO\IB_IN'
    file_list = list(glob2.iglob(os.path.join(path, "*.txt")))
    # 创建线程池,max_workers可以根据磁盘性能调整,一般8-16就行
    with ThreadPoolExecutor(max_workers=8) as executor:
        tasks = [read_csv(executor, f) for f in file_list]
        dfs = await asyncio.gather(*tasks)
        df = pd.concat(dfs, ignore_index=True)
    return df

if __name__ == "__main__":
    t = time.time()
    df = asyncio.run(main())
    print(f"异步版本耗时: {time.time()-t:.4f}")

二、优化pd.concat的耗时技巧

除了读取环节,合并DataFrame的速度也可以通过以下方式优化:

  • 提前指定数据类型:
    读取CSV时手动指定每列的dtype,避免pandas自动推断类型的开销,同时减少内存占用,合并时更快。比如:

    # 假设你的文件有3列,分别是浮点、整数、字符串
    dtype = {0: 'float64', 1: 'int32', 2: 'object'}
    return await asyncio.to_thread(pd.read_csv, filename, header=None, dtype=dtype)
    
  • 压缩内存占用:
    用downcast参数对数值型列进行内存压缩,比如:

    pd.read_csv(filename, header=None, downcast='float', downcast='integer')
    

    更小的内存占用意味着合并时更少的内存拷贝,速度自然更快。

  • 跳过不必要的操作:
    如果不需要重置索引,去掉ignore_index=True,节省索引重建的时间;如果所有文件的列顺序完全一致,pandas不需要对齐列,合并速度也会提升。

  • 超大数据量用Dask:
    如果你的文件总数据量很大(比如几十GB),pandas可能会内存不足,这时候可以用Dask来并行读取和合并,它会自动分块处理:

    import dask.dataframe as dd
    path = r'C:\LRM_STGY_REPO\IB_IN\*.txt'
    ddf = dd.read_csv(path, header=None)
    df = ddf.compute()  # 转换为pandas DataFrame
    

最后补充

异步读取的提速效果和你的磁盘性能有关:SSD的并行读取提升会很明显,机械硬盘可能提升有限。另外线程池的max_workers不要设置太大,否则会导致磁盘IO瓶颈,一般设置为CPU核心数的2-4倍就好。

内容的提问来源于stack exchange,提问作者Relax ZeroC

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:25:15