如何用Pandas读取大型SAS文件并多进程导出为Parquet?代码排障
问题排查与修复
你的代码核心问题在于多线程任务分配逻辑错误,导致实际未正确处理数据块,因此没有生成Parquet文件:
错误原因分析
pd.read_sas(chunksize=...)返回的是迭代器,每次迭代产出单个chunk的DataFrameexecutor.map(sas_mult_process, file_reader)会把每个chunk(单个DataFrame)逐个传给sas_mult_process函数- 但你在函数里写了
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
相关产品推荐
相关产品推荐

