PySpark读取解析目录下多JSON格式TXT文件的解决方案求助
解决PySpark批量读取嵌套JSON格式TXT文件的问题
核心问题原因
传统Python文件读取方法(如open)不支持解析*通配符,必须依赖Shell展开;而PySpark的数据源API原生支持通配符匹配与目录批量读取,直接使用通配符路径即可解决文件找不到的错误。
正确实现步骤
1. 批量读取文件
PySpark的spark.read.json()可直接识别通配符路径,且会自动合并所有结构一致文件的Schema:
from pyspark.sql import SparkSession # 初始化SparkSession spark = SparkSession.builder \ .appName("BatchNestedJSONProcessor") \ .getOrCreate() # 读取目录下所有txt格式的JSON文件 df = spark.read.json("/home/input/*.txt") # 验证数据结构与内容 df.printSchema() df.show(5, truncate=False)
2. 处理嵌套结构
针对嵌套JSON,可通过pyspark.sql.functions展开或提取指定字段:
from pyspark.sql import functions as F # 示例:展开名为`user_profile`的嵌套列 flattened_df = df.select( "order_id", "create_time", F.col("user_profile.name").alias("user_name"), F.col("user_profile.contact.phone").alias("user_phone"), F.col("user_profile.address.city").alias("user_city") ) flattened_df.show()
3. 替代方案:读取整个目录
若无需精确匹配后缀,直接指定目录路径也会加载目录下所有文件(可通过recursiveFileLookup控制是否读取子目录):
# 读取目录及子目录下所有JSON格式文件 df = spark.read.json("/home/input/", recursiveFileLookup=True)
常见异常处理
- 若文件是单个文件包含多个JSON对象(非每行一个JSON),需添加
multiLine=True参数:df = spark.read.json("/home/input/*.txt", multiLine=True) - 确保Spark进程拥有
/home/input/目录的读取权限
内容的提问来源于stack exchange,提问作者pravek
相关产品推荐
相关产品推荐

