如何批量将多实体JSON文件转换为DataFrame?(十亿级规模)
解决方案:流式解析+列式存储处理超大规模JSON数据
针对你遇到的ValueError: Value is too big!问题,本质是pd.read_json一次性加载整个超大JSON文件到内存导致的内存溢出。要处理十亿级规模的文件,必须采用流式解析+分块写入列式存储的方案,核心思路是避免一次性加载全部数据,逐节点处理并持久化。
核心依赖库
ijson: 流式JSON解析库,仅加载当前处理的节点,不占用全量内存pandas: 小批量数据转换为DataFramepyarrow: 用于将数据写入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)
关键说明
- 流式解析:
ijson仅读取当前处理的JSON节点,内存占用稳定在批次数据的大小,完全支持十亿级文件处理 - 列式存储:Parquet比JSON压缩率高5-10倍,后续用
pd.read_parquet或Dask读取时,可只加载需要的列,大幅提升查询效率 - 分块写入:
chunk_size可根据机器内存调整,比如内存充足可设为100000,内存紧张则设为1000
特殊场景处理
- 如果某个顶级键是单个对象(比如
block)而非数组,修改ijson的读取路径为entity_key而非{entity_key}.item,直接将单个对象转为DataFrame的一行后写入 - 若需要分布式处理超大规模数据,可替换
pandas为Dask,配合Dask的Parquet读写实现并行处理
内容的提问来源于stack exchange,提问作者Mert Deveci
相关产品推荐
相关产品推荐

