如何高效将大体积JSONL文件转换为Parquet格式?
大规模JSONL转Parquet的内存友好方案
针对数十亿行的超大规模JSONL文件,无需全量加载到内存的转换方案主要有以下几种:
1. 用Dask实现单机/轻量集群处理
Dask是基于Python的并行计算框架,支持分块读取和处理数据,内存占用可控:
import dask.dataframe as dd # 分块读取JSONL,lines=True表示每行是一个JSON对象 ddf = dd.read_json('mydata.jsonl', lines=True) # 处理编码问题,修复乱码 def fix_encoding(col): return col.str.encode('utf8', 'replace').str.decode('utf8') for col in ddf.columns: if ddf[col].dtype == object: ddf[col] = fix_encoding(ddf[col]) # 写入Parquet,默认按分块存储,可通过blocksize参数调整分区大小 ddf.to_parquet('mydata_dask.parquet', engine='fastparquet')
Dask会自动将数据拆分为多个分区,每次仅处理单个分区,内存占用由分区大小决定,适合单机处理百GB级别的数据。
2. 用PySpark处理超大规模集群数据
如果数据量达到TB/PB级别,单机性能不足,PySpark的分布式计算可以完全规避全量加载问题:
from pyspark.sql import SparkSession from pyspark.sql.functions import udf, col from pyspark.sql.types import StringType # 初始化Spark会话 spark = SparkSession.builder.appName("JSONLtoParquet").getOrCreate() # 读取JSONL文件,multiLine=False表示每行一个JSON df = spark.read.json('mydata.jsonl', multiLine=False) # 定义UDF修复编码问题 fix_encoding_udf = udf(lambda x: str(x).encode('utf8', 'replace').decode('utf8') if x else x, StringType()) for c in df.columns: if df.schema[c].dataType.typeName() == 'string': df = df.withColumn(c, fix_encoding_udf(col(c))) # 写入Parquet,支持按字段分区存储,优化后续查询 df.write.parquet('mydata_spark.parquet')
Spark会将数据分布式存储在集群节点上,每个节点仅处理部分数据,完全不需要将全量数据加载到单节点内存,是处理数十亿行数据的最优选择之一。
3. 手动分批处理(纯单机轻量方案)
如果不想依赖框架,可以手动控制批次大小,逐批读取、处理并写入Parquet:
import json import pandas as pd import fastparquet # 设置每批处理的行数,根据机器内存调整 BATCH_SIZE = 100000 parquet_writer = None with open('mydata.jsonl', 'r') as f: batch = [] for line in f: try: # 解析单条JSON记录 batch.append(json.loads(line)) except json.JSONDecodeError: # 跳过无效行 continue # 达到批次大小则处理 if len(batch) >= BATCH_SIZE: df_batch = pd.json_normalize(batch) # 修复编码 for col in df_batch.columns: if df_batch[col].dtype == object: df_batch[col] = df_batch[col].apply(lambda x: str(x).encode('utf8', 'replace').decode('utf8') if pd.notna(x) else x) # 写入Parquet,首次创建writer,后续追加 if not parquet_writer: parquet_writer = fastparquet.ParquetWriter('mydata_batch.parquet', df_batch.columns, index=False) parquet_writer.write_table(fastparquet.api.Table.from_pandas(df_batch)) batch = [] # 处理剩余的最后一批数据 if batch: df_batch = pd.json_normalize(batch) for col in df_batch.columns: if df_batch[col].dtype == object: df_batch[col] = df_batch[col].apply(lambda x: str(x).encode('utf8', 'replace').decode('utf8') if pd.notna(x) else x) if parquet_writer: parquet_writer.write_table(fastparquet.api.Table.from_pandas(df_batch)) else: fastparquet.write('mydata_batch.parquet', df_batch, index=False) if parquet_writer: parquet_writer.close()
这种方式完全通过手动控制内存占用,适合没有集群资源但需要处理大文件的场景,调整BATCH_SIZE即可适配不同机器的内存容量。
内容的提问来源于stack exchange,提问作者alvas
相关产品推荐
相关产品推荐

