如何正确使用asyncio结合Pandas读取合并多CSV文件以提升效率?
嘿,我看了你的代码,问题出在你以为用了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

