PySpark批量转换多JSON至Parquet及循环代码语法错误排查
批量将不同Schema的JSON文件转换为Parquet解决方案
一、遍历目录代码的语法错误排查
原遍历代码存在3处问题:
- 多余转义符:
Directory定义行末尾的\\无意义,导致语法中断,直接删除即可。 - 方法名错误:
os.join应为os.path.join,且无需path=f''的冗余写法。 - 逻辑误区:
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
相关产品推荐
相关产品推荐

