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

PySpark批量转换多JSON至Parquet及循环代码语法错误排查

批量将不同Schema的JSON文件转换为Parquet解决方案

一、遍历目录代码的语法错误排查

原遍历代码存在3处问题:

  1. 多余转义符:Directory定义行末尾的\\无意义,导致语法中断,直接删除即可。
  2. 方法名错误:os.join应为os.path.join,且无需path=f''的冗余写法。
  3. 逻辑误区:os.listdir和open仅支持本地文件系统,无法直接处理ABFS(Azure Data Lake Storage)路径,需用PySpark或Azure专属工具访问。

修正后的本地文件遍历代码(仅适用于本地路径):

import os 
directory = 'path/to/local/directory'
for filename in os.listdir(directory):
    if filename.endswith('.json'):
        file_path = os.path.join(directory, filename)
        print(file_path)

二、批量转换不同Schema JSON到Parquet的完整实现

由于每个JSON文件Schema不同,必须逐个读取、转换、写入。以下是基于PySpark的适配ABFS路径的方案:

1. 初始化PySpark会话

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("Batch JSON to Parquet Conversion") \
    .getOrCreate()

2. 列出ABFS路径下的所有JSON文件

使用Spark内置的Hadoop文件系统API遍历ABFS目录:

# 初始化Hadoop文件系统客户端
fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(spark._jsc.hadoopConfiguration())
target_path = spark._jvm.org.apache.hadoop.fs.Path('abfss://folder_j@{AZ_NM}.dfs.core.windows.net/files_j/all_files')

# 过滤出所有JSON文件
file_status_list = fs.listStatus(target_path)
json_files = [
    status.getPath().toString() 
    for status in file_status_list 
    if status.getPath().getName().endswith('.json')
]

3. 逐个转换并写入Parquet

遍历每个JSON文件,自动推断Schema后写入同名Parquet文件:

for json_path in json_files:
    # 生成对应的Parquet路径
    parquet_path = json_path.replace('.json', '.parquet')
    
    # 读取JSON(自动推断每个文件的Schema)
    df = spark.read.json(json_path)
    
    # 写入Parquet,覆盖已有文件(可根据需求改为append/ignore)
    df.write.mode('overwrite').parquet(parquet_path)
    
    print(f"转换完成:{json_path} -> {parquet_path}")

# 停止Spark会话
spark.stop()

关键注意事项

  • ABFS权限配置:确保PySpark环境已配置Azure Data Lake Storage的访问凭证(如Service Principal、SAS Token),否则无法读取ABFS路径文件。
  • Schema推断:spark.read.json会自动识别每个文件的Schema,适配多Schema场景。
  • 写入模式:mode('overwrite')用于覆盖已存在的Parquet文件,可根据业务需求调整。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 19:17:45