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

如何用Pandas读取大型SAS文件并多进程导出为Parquet?代码排障

问题排查与修复

你的代码核心问题在于多线程任务分配逻辑错误,导致实际未正确处理数据块,因此没有生成Parquet文件:

错误原因分析

  1. pd.read_sas(chunksize=...)返回的是迭代器,每次迭代产出单个chunk的DataFrame
  2. executor.map(sas_mult_process, file_reader)会把每个chunk(单个DataFrame)逐个传给sas_mult_process函数
  3. 但你在函数里写了for i, df in enumerate(data)——这里的data已经是单个DataFrame,enumerate会遍历它的列名而非数据块,自然不会生成有效文件

另外补充:ThreadPoolExecutor适合IO密集型任务(如文件读写),如果是CPU密集型的格式转换,ProcessPoolExecutor能规避GIL限制,性能更优。

修复后的代码

版本1:修正线程池逻辑(保留ThreadPoolExecutor)

import pandas as pd
from concurrent.futures import ThreadPoolExecutor

arq = "path_to_my_file"

def process_chunk(args):
    chunk, idx = args
    chunk.to_parquet(f"hist_dif_base_pt{idx}.parquet")

# 生成带索引的chunk迭代器,确保每个块有唯一文件名
file_reader = pd.read_sas(arq, chunksize=100000, encoding='ISO-8859-1', format='sas7bdat')
chunk_with_idx = ((chunk, i) for i, chunk in enumerate(file_reader))

with ThreadPoolExecutor(max_workers=10) as executor:
    executor.map(process_chunk, chunk_with_idx)

版本2:改用进程池(适合CPU密集场景)

import pandas as pd
from concurrent.futures import ProcessPoolExecutor

arq = "path_to_my_file"

def process_chunk(args):
    chunk, idx = args
    chunk.to_parquet(f"hist_dif_base_pt{idx}.parquet")

file_reader = pd.read_sas(arq, chunksize=100000, encoding='ISO-8859-1', format='sas7bdat')
chunk_with_idx = ((chunk, i) for i, chunk in enumerate(file_reader))

# Windows环境下使用ProcessPoolExecutor必须将执行代码放入该判断块内
if __name__ == '__main__':
    with ProcessPoolExecutor(max_workers=4) as executor:
        executor.map(process_chunk, chunk_with_idx)

额外优化建议

  • 可以指定Parquet压缩格式减少文件体积:chunk.to_parquet(..., compression='snappy')
  • 若后续需合并小文件,可使用pyarrow.parquet.read_table批量读取后合并保存,提升效率

内容的提问来源于stack exchange,提问作者Gabriel Otavio Canieto Costa

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 13:55:23