拼接DataFrame遇MemoryError求助(超1万JSON文件场景)
解决大量嵌套JSON转DataFrame的MemoryError问题
针对你处理5830个JSON文件时出现的内存不足错误(需41.3GiB内存),以及后续要处理超10000个Azure数据湖文件的需求,给出以下实用方案:
1. 分块处理+增量写入
不要一次性加载所有文件到内存,而是分批处理并直接写入输出介质,避免保留全量DataFrame:
- 批量遍历Azure数据湖的文件,比如每500个文件为一组;
- 处理每组文件(提取嵌套JSON的目标字段)后,将数据追加写入到Parquet/CSV文件,或者直接写入数据库;
- 处理完每组后,手动清空当前DataFrame并调用
gc.collect()释放内存。
示例代码片段:
import gc import json import pandas as pd from azure.storage.filedatalake import DataLakeServiceClient # 初始化ADLS客户端 service_client = DataLakeServiceClient.from_connection_string("your_conn_str") file_system_client = service_client.get_file_system_client(file_system="your-container") # 分批处理 batch_size = 500 file_list = [file.name for file in file_system_client.get_paths() if file.is_file] for i in range(0, len(file_list), batch_size): batch_files = file_list[i:i+batch_size] batch_df = pd.DataFrame() for file_name in batch_files: file_client = file_system_client.get_file_client(file_name) json_data = file_client.download_file().readall() # 提取嵌套JSON的目标字段,生成单条数据的DataFrame single_df = pd.json_normalize(json.loads(json_data), record_path=['target_path'], meta=['meta_col1']) batch_df = pd.concat([batch_df, single_df], ignore_index=True) # 增量写入Parquet(推荐,比CSV更省空间且支持分块) batch_df.to_parquet("output.parquet", mode='append', engine='pyarrow') # 释放内存 del batch_df gc.collect()
2. 裁剪数据字段+优化数据类型
报错显示DataFrame有5112列,多数场景下不需要这么多字段,同时优化数据类型可大幅降低内存占用:
- 处理嵌套JSON时,只提取业务需要的字段,避免加载冗余数据;
- 将
float64类型转为float32(精度损失可接受时),整数类型按需缩小(比如int64转int32/int16); - 读取JSON时直接指定数据类型:
# 示例:指定部分字段的类型 dtype_spec = {'numeric_col1': 'float32', 'numeric_col2': 'int32'} single_df = pd.json_normalize(json.loads(json_data), record_path=['target_path'], meta=['meta_col1'], dtype=dtype_spec)
- 处理完数据后,用
df.memory_usage(deep=True)查看内存占用,针对性优化。
3. 用Dask替代Pandas处理超大规模数据
Dask支持并行分块处理,无需加载全量数据到内存,天然适配大数据场景:
- 使用
dask.dataframe.read_json批量读取ADLS上的JSON文件,Dask会自动分块; - 直接将Dask DataFrame写入Parquet或数据库,全程无需加载全量数据:
import dask.dataframe as dd # 读取ADLS上的所有JSON文件(路径格式参考azure://your-container/*.json) ddf = dd.read_json("azure://your-container/*.json", blocksize='64MB') # 提取嵌套字段(类似pandas的json_normalize) ddf_normalized = ddf.map_partitions(lambda df: pd.json_normalize(df.to_dict('records'), record_path=['target_path'])) # 写入Parquet ddf_normalized.to_parquet("azure://your-output-container/output.parquet", engine='pyarrow')
4. 直接写入数据库跳过中间DataFrame
如果最终数据要存入数据库,可跳过拼接大DataFrame的步骤,每批数据直接插入:
- 用SQLAlchemy或数据库SDK实现批量插入;
- 示例(SQLAlchemy批量插入):
from sqlalchemy import create_engine, Table, MetaData engine = create_engine("your-db-connection-string") metadata = MetaData() table = Table('your_table', metadata, autoload_with=engine) # 处理每批文件后,将数据转为字典列表 data_mappings = batch_df.to_dict('records') # 批量插入 with engine.connect() as conn: conn.execute(table.insert(), data_mappings) conn.commit()
内容的提问来源于stack exchange,提问作者rclee
相关产品推荐
相关产品推荐

