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

如何合并多文件的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()

关键注意事项

  1. 内存控制:两种方案都确保每个进程仅加载单个chunk到内存,处理完成后释放,避免一次性加载全量数据导致内存溢出。
  2. 任务设计:pool.map要求输入可迭代的独立任务列表,任务需无状态——要么是单个文件路径,要么是单个chunk的定位参数。
  3. TextFileReader正确用法:不要尝试合并多个读取器,而是迭代每个读取器获取实际chunk数据,将chunk作为并行处理的最小单元。

内容的提问来源于stack exchange,提问作者Lukasz Hoszowski

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 06:45:32