如何用Python并行分块读取处理20GB级序列化JSON大文件?
针对你这种20GB、2000万行的JSON行文件处理需求,结合Python3.6的特性,我整理了几个经过实践验证的高效方案,从易上手到进阶优化都有,你可以根据自己的场景选择:
1. 多进程分块并行处理(最易落地的基础方案)
Python的GIL限制了多线程在CPU密集型任务的效率,而处理JSON解析属于CPU+IO混合场景,用多进程是最直接的选择。核心思路是把文件分成几个连续的块,每个进程负责读取并处理一个块,注意要保证块的结尾是完整的行(避免截断JSON)。
实现步骤:
- 先获取文件总大小,按内存情况划分块(比如每个块100MB-1GB,不要太小导致进程开销过高)
- 每个进程定位到对应块的起始位置,然后向后找到第一个换行符,从这里开始读取到块的结束位置(同样找最后一个换行符)
- 用
concurrent.futures.ProcessPoolExecutor来管理进程池,提交分块处理任务
示例代码:
import os import ujson # 替换标准json库,解析速度快3-5倍,Python3.6支持 from concurrent.futures import ProcessPoolExecutor def process_chunk(file_path, start_pos, end_pos): results = [] with open(file_path, 'rb') as f: f.seek(start_pos) # 跳过可能不完整的首行 if start_pos != 0: f.readline() # 读取到块结束位置 while f.tell() < end_pos: line = f.readline().decode('utf-8').strip() if not line: continue try: data = ujson.loads(line) # 这里写你的处理逻辑,比如提取字段、计算等 processed = {"id": data.get("id"), "value": data.get("value") * 2} results.append(processed) except Exception as e: print(f"处理行失败: {e}") continue return results def split_file_into_chunks(file_path, chunk_size=100*1024*1024): # 100MB每块 file_size = os.path.getsize(file_path) chunks = [] start = 0 while start < file_size: end = min(start + chunk_size, file_size) chunks.append((start, end)) start = end return chunks if __name__ == "__main__": file_path = "your_large_file.jsonl" chunks = split_file_into_chunks(file_path) # 进程数建议等于CPU核心数,避免过度调度 with ProcessPoolExecutor(max_workers=os.cpu_count()) as executor: futures = [executor.submit(process_chunk, file_path, start, end) for start, end in chunks] # 收集所有结果 all_results = [] for future in futures: all_results.extend(future.result()) # 这里可以把结果写入文件或数据库 print(f"处理完成,共处理{len(all_results)}条数据")
2. Dask并行数据流处理(前沿大数据工具,适合超大规模数据)
如果你需要处理更大规模的文件(甚至TB级),或者想更优雅地处理并行任务,Dask是非常合适的选择。它是一个并行计算库,可以模拟Pandas/NumPy的接口,但能处理超出内存的数据集,底层自动分块并行。
示例代码(Python3.6兼容):
import dask.bag as db import ujson # 读取JSON行文件,自动分块 b = db.read_text("your_large_file.jsonl").map(ujson.loads) # 定义处理函数 def process_data(data): return {"id": data.get("id"), "value": data.get("value") * 2} # 并行处理并收集结果 processed_bag = b.map(process_data) all_results = processed_bag.compute() # 触发计算 print(f"处理完成,共处理{len(all_results)}条数据")
注意:Python3.6需要安装兼容版本的Dask,比如pip install dask==2021.12.0(更高版本可能不再支持3.6)。
3. PyArrow高效分块读取(性能优先,适合序列化数据)
PyArrow是专门为大数据处理设计的库,它的IO性能远超Python标准库,能快速读取大文件并分批次处理,结合多进程可以进一步提升效率。
示例代码:
import os import pyarrow.json as pajson from concurrent.futures import ProcessPoolExecutor import ujson def process_batch(batch): # 将Arrow批次转为Python字典列表 data_list = batch.to_pylist() results = [] for data in data_list: try: processed = {"id": data.get("id"), "value": data.get("value") * 2} results.append(processed) except Exception as e: print(f"处理数据失败: {e}") continue return results if __name__ == "__main__": # 读取JSON行文件为Arrow批次,设置批次大小 reader = pajson.open_json("your_large_file.jsonl", read_options=pajson.ReadOptions(block_size=100*1024*1024)) batches = [] try: while True: batch = reader.read_next_batch() batches.append(batch) except StopIteration: pass # 用进程池处理每个批次 with ProcessPoolExecutor(max_workers=os.cpu_count()) as executor: futures = [executor.submit(process_batch, batch) for batch in batches] all_results = [] for future in futures: all_results.extend(future.result()) print(f"处理完成,共处理{len(all_results)}条数据")
关键最佳实践
- 用ujson替代标准json库:ujson的解析速度是标准库的3-5倍,能大幅降低CPU耗时,Python3.6完全支持。
- 合理设置分块大小:分块太小会导致进程创建/销毁开销过高,太大则会占用过多内存,建议根据你的内存情况设置100MB-1GB每块。
- 避免进程间传递大对象:让每个进程直接读取对应块的内容,不要在主进程读取后再传递给子进程,减少IPC开销。
- 错误处理不可少:大文件难免有格式错误的行,一定要捕获异常,避免单个错误导致整个任务崩溃。
- 如果是IO瓶颈:可以考虑用SSD存储文件,或者开启文件预读(比如
open时用buffering=8*1024*1024增大缓冲区)。
内容的提问来源于stack exchange,提问作者Nodirbek Shamsiev
相关产品推荐
相关产品推荐

