如何合并多文件的Pandas TextFileReader块以实现多进程处理
问题描述
我长期浏览Stack Overflow,此前均能自行找到答案或解决问题,这是我首次提问。
环境配置
- Python 3.10.4
- Pandas 1.4.2
需求规划
我需对3个各约90万行的大型CSV文件执行ETL操作,现有脚本运行较慢。受IO限制,选择使用multiprocessing利用全部CPU资源。测试时可并行执行read_csv,但脚本内存占用过高,很快耗尽RAM。因此计划对文件进行chunk处理,将所有文件的块收集到容器中交给Pool,让worker高效处理,但遇到问题。
核心问题
我不知道如何合并多个read_csv生成的TextFileReader对象,能否创建类似队列的多文件块集合?
当前临时方案
将多个大型CSV合并为一个2.5GB的文件,只需对单个文件分块即可获得完整块集合,但该方案繁琐且耗时。
期望方案
多CSV文件 → 分块 → 合并为统一的块集合
尝试的代码
def csv2chunks(file_list): chunks = [] for file in file_list: chunks.append(pd.read_csv(filename, dtype='string[pyarrow]', usecols=['Year', 'Month', 'Term'], chunksize=10000, iterator=True)) return chunks def main(): files = os.path.join("./", "volume PN FY*.csv") files = glob.glob(files) file_list = files with Pool(processes=3) as pool: # or whatever your hardware can support df_list = pool.map(csv2chunks), file_list) combined_df = pd.concat(df_list, ignore_index=True) return combined_df if __name__ == '__main__': df = main()
我感觉遗漏了关键内容,希望得到关于TextFileReader对象特性及操作方法的建议与帮助。
解决方案与说明
关于TextFileReader的特性
TextFileReader是pandas返回的迭代器对象,并非内存中的数据集集合——它不会一次性加载所有数据,而是每次迭代时才读取一个chunk到内存。因此你不能直接把多个TextFileReader合并成一个容器,因为它们本身只是"读取器",而非实际数据块的集合。
实现多文件分块并行处理的正确思路
不需要提前合并读取器,而是将单个文件/单个chunk的处理逻辑拆分为独立任务,交给进程池执行:
方案1:每个进程处理单个完整文件的分块
适合ETL逻辑需要对单个文件的所有chunk做关联处理的场景:
import os import glob import pandas as pd from multiprocessing import Pool def process_single_file(file_path): processed_chunks = [] # 分块读取单个文件并执行ETL操作 for chunk in pd.read_csv(file_path, dtype='string[pyarrow]', usecols=['Year', 'Month', 'Term'], chunksize=10000): # 替换为你的实际ETL逻辑,比如清洗、转换 processed_chunk = chunk.drop_duplicates() processed_chunks.append(processed_chunk) # 合并当前文件的所有处理后chunk return pd.concat(processed_chunks, ignore_index=True) def main(): files = glob.glob(os.path.join("./", "volume PN FY*.csv")) # 用进程池并行处理每个文件 with Pool(processes=3) as pool: processed_dfs = pool.map(process_single_file, files) # 合并所有文件的最终结果 combined_df = pd.concat(processed_dfs, ignore_index=True) return combined_df if __name__ == '__main__': df = main()
方案2:将所有文件的chunk拆分为独立任务(细粒度并行)
适合ETL逻辑仅针对单chunk即可完成的场景,最大化CPU利用率:
import os import glob import pandas as pd from multiprocessing import Pool def process_single_chunk(args): file_path, chunk_idx, chunksize = args # 根据索引定位并读取指定chunk chunk = pd.read_csv( file_path, dtype='string[pyarrow]', usecols=['Year', 'Month', 'Term'], chunksize=chunksize, skiprows=1 + chunk_idx * chunksize, # 跳过表头和之前的chunk header=None, names=['Year', 'Month', 'Term'] ) chunk = next(chunk) # 替换为你的实际ETL逻辑 processed_chunk = chunk.drop_duplicates() return processed_chunk def generate_chunk_tasks(files, chunksize=10000): tasks = [] for file in files: # 计算当前文件的总chunk数(跳过表头行) with open(file, 'r') as f: total_rows = sum(1 for _ in f) - 1 total_chunks = (total_rows + chunksize - 1) // chunksize # 生成每个chunk的任务参数 for chunk_idx in range(total_chunks): tasks.append((file, chunk_idx, chunksize)) return tasks def main(): files = glob.glob(os.path.join("./", "volume PN FY*.csv")) chunk_tasks = generate_chunk_tasks(files, chunksize=10000) # 用全部CPU核心处理所有chunk任务 with Pool(processes=os.cpu_count()) as pool: processed_chunks = pool.map(process_single_chunk, chunk_tasks) combined_df = pd.concat(processed_chunks, ignore_index=True) return combined_df if __name__ == '__main__': df = main()
关键注意事项
- 内存控制:两种方案都确保每个进程仅加载单个chunk到内存,处理完成后释放,避免一次性加载全量数据导致内存溢出。
- 任务设计:
pool.map要求输入可迭代的独立任务列表,任务需无状态——要么是单个文件路径,要么是单个chunk的定位参数。 TextFileReader正确用法:不要尝试合并多个读取器,而是迭代每个读取器获取实际chunk数据,将chunk作为并行处理的最小单元。
内容的提问来源于stack exchange,提问作者Lukasz Hoszowski
相关产品推荐
相关产品推荐

