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

如何批量将多实体JSON文件转换为DataFrame?(十亿级规模)

解决方案:流式解析+列式存储处理超大规模JSON数据

针对你遇到的ValueError: Value is too big!问题,本质是pd.read_json一次性加载整个超大JSON文件到内存导致的内存溢出。要处理十亿级规模的文件,必须采用流式解析+分块写入列式存储的方案,核心思路是避免一次性加载全部数据,逐节点处理并持久化。

核心依赖库

  • ijson: 流式JSON解析库,仅加载当前处理的节点,不占用全量内存
  • pandas: 小批量数据转换为DataFrame
  • pyarrow: 用于将数据写入Parquet列式存储(压缩率高、查询效率远优于JSON)

先安装依赖:

pip install ijson pandas pyarrow

单个JSON文件处理:提取指定实体并写入Parquet

以下函数会流式读取JSON中的指定顶级键(比如transactions),分批次将数据写入Parquet文件,避免内存过载:

import ijson
import pandas as pd
import pyarrow as pa
from pyarrow import parquet as pq
import os

def extract_entity_to_parquet(json_file_path, entity_key, output_parquet_dir, chunk_size=10000):
    """
    从超大JSON文件中流式提取指定实体,分块写入Parquet
    :param json_file_path: 输入JSON文件路径
    :param entity_key: 要提取的顶级键(如'transactions')
    :param output_parquet_dir: 输出Parquet数据集的目录
    :param chunk_size: 每次写入的批次大小,根据内存调整
    """
    os.makedirs(output_parquet_dir, exist_ok=True)
    
    with open(json_file_path, 'rb') as f:
        # 定位到顶级键对应的数组节点,逐个读取元素
        entities = ijson.items(f, f'{entity_key}.item')
        
        chunk = []
        for idx, item in enumerate(entities, 1):
            chunk.append(item)
            
            # 达到批次大小则写入
            if idx % chunk_size == 0:
                df_chunk = pd.DataFrame(chunk)
                table = pa.Table.from_pandas(df_chunk)
                # 首次写入创建数据集,后续追加
                if idx == chunk_size:
                    pq.write_table(table, os.path.join(output_parquet_dir, 'data.parquet'))
                else:
                    pq.write_to_dataset(table, root_path=output_parquet_dir, append=True)
                chunk = []
        
        # 处理剩余的不足批次大小的数据
        if chunk:
            df_chunk = pd.DataFrame(chunk)
            table = pa.Table.from_pandas(df_chunk)
            pq.write_to_dataset(table, root_path=output_parquet_dir, append=True)

批量处理多个JSON文件

循环遍历所有JSON文件,将同类型实体合并到同一个Parquet数据集:

def batch_process_json_files(json_dir, output_root_dir, entities=['block', 'transactions', 'logs', 'token_transfers', 'traces']):
    """
    批量处理目录下所有JSON文件,提取指定实体到对应Parquet目录
    :param json_dir: JSON文件所在目录
    :param output_root_dir: 所有实体Parquet文件的根目录
    :param entities: 需要提取的顶级键列表
    """
    # 初始化每个实体的输出目录
    for entity in entities:
        os.makedirs(os.path.join(output_root_dir, entity), exist_ok=True)
    
    # 遍历所有JSON文件
    for filename in os.listdir(json_dir):
        if not filename.endswith('.json'):
            continue
        json_path = os.path.join(json_dir, filename)
        print(f'Processing: {json_path}')
        
        for entity in entities:
            output_dir = os.path.join(output_root_dir, entity)
            extract_entity_to_parquet(json_path, entity, output_dir)

调用示例:

# 替换为你的实际路径
json_directory = '/your/json/files/dir'
output_root = '/your/parquet/output/dir'

batch_process_json_files(json_directory, output_root)

关键说明

  1. 流式解析:ijson仅读取当前处理的JSON节点,内存占用稳定在批次数据的大小,完全支持十亿级文件处理
  2. 列式存储:Parquet比JSON压缩率高5-10倍,后续用pd.read_parquet或Dask读取时,可只加载需要的列,大幅提升查询效率
  3. 分块写入:chunk_size可根据机器内存调整,比如内存充足可设为100000,内存紧张则设为1000

特殊场景处理

  • 如果某个顶级键是单个对象(比如block)而非数组,修改ijson的读取路径为entity_key而非{entity_key}.item,直接将单个对象转为DataFrame的一行后写入
  • 若需要分布式处理超大规模数据,可替换pandas为Dask,配合Dask的Parquet读写实现并行处理

内容的提问来源于stack exchange,提问作者Mert Deveci

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 05:30:46