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

拼接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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 10:42:57