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

如何高效将大体积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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 11:23:20