如何加速Python中5GB级.json.gz大文件的迭代处理?
处理超大JSON.gz文件的加速方案
问题背景
有一个超大型goodreads_reviews_dedup.json.gz文件,需要提取user_id、book_id、rating三个字段保存为Parquet格式的DataFrame,用于推荐系统开发。当前逐行写入DataFrame的代码处理速度极慢:1500万行数据预计耗时20小时,30分钟仅处理35万行。已尝试ultrajson和Google Colab,效果不佳。
原代码:
n_review=0 df = pd.DataFrame(columns = ['user_id', 'book_id', 'rating']) with gzip.open('goodreads_reviews_dedup.json.gz') as fin: for l in fin: d = ujson.loads(l) df.loc[n_review] = [d['user_id'],d['book_id'],d['rating']] #counter to check progress if n_review % 1000 == 0: print(n_review, end=',') n_review += 1
核心瓶颈
原代码的最大问题是逐行使用df.loc写入DataFrame。DataFrame是列导向的数据结构,每次loc操作都会触发内存重排,时间复杂度为O(n²),数据量越大,性能下降越明显。
加速方案
1. 先收集数据到列表,再一次性转DataFrame
这是最直接的优化,将逐行写入改为列表批量收集,最后一次性生成DataFrame,时间复杂度降至O(n)。
import gzip import ujson import pandas as pd data = [] n_review = 0 with gzip.open('goodreads_reviews_dedup.json.gz') as fin: for l in fin: d = ujson.loads(l) # 仅收集需要的字段到列表 data.append((d['user_id'], d['book_id'], d['rating'])) if n_review % 10000 == 0: # 调整打印间隔减少IO开销 print(f"Processed {n_review} rows", end='\r') n_review += 1 # 批量生成DataFrame df = pd.DataFrame(data, columns=['user_id', 'book_id', 'rating']) # 保存为Parquet(默认snappy压缩,空间和速度平衡) df.to_parquet('goodreads_reviews_slim.parquet', index=False)
2. 用Pandas原生read_json直接处理
Pandas的read_json针对JSON行文件做了底层优化,支持直接读取压缩文件,还能指定仅加载需要的列,效率远高于手动逐行解析。
import pandas as pd # 直接读取压缩JSON行文件,仅加载目标字段 df = pd.read_json( 'goodreads_reviews_dedup.json.gz', compression='gzip', lines=True, usecols=['user_id', 'book_id', 'rating'] ) df.to_parquet('goodreads_reviews_slim.parquet', index=False)
3. 分块处理(内存不足时用)
如果文件过大导致内存不够,可通过chunksize分块读取,逐块写入Parquet文件。
import pandas as pd chunk_size = 100000 # 每块处理10万行 first_chunk = True for chunk in pd.read_json( 'goodreads_reviews_dedup.json.gz', compression='gzip', lines=True, usecols=['user_id', 'book_id', 'rating'], chunksize=chunk_size ): if first_chunk: chunk.to_parquet('goodreads_reviews_slim.parquet', index=False) first_chunk = False else: # 追加模式写入 chunk.to_parquet('goodreads_reviews_slim.parquet', mode='append', index=False) print(f"Processed {chunk_size} rows...")
4. 多核并行处理(超大数据集)
如果本地CPU有多个核心,可使用Dask库自动实现并行处理,充分利用硬件资源。
import dask.dataframe as dd # Dask会自动拆分文件并多核并行处理 ddf = dd.read_json( 'goodreads_reviews_dedup.json.gz', compression='gzip', lines=True, usecols=['user_id', 'book_id', 'rating'] ) # 保存为Parquet(支持分块存储) ddf.to_parquet('goodreads_reviews_slim_dask.parquet', write_index=False)
额外优化建议
- 确保使用64位Python,且Pandas、Numpy通过Conda安装(预编译优化版本,性能优于pip安装);
- 减少打印频率(比如每1万行打印一次),避免频繁IO操作拖慢速度;
- 本地处理优先于Colab,避免网络IO和资源限制导致的降速。
内容的提问来源于stack exchange,提问作者krill_445
相关产品推荐
相关产品推荐

