PySpark读取JSON写入Delta Lake全为Null的问题解决与日志解析
PySpark读取JSON全为Null及Delta Lake临时目录错误的解决方案
一、解决读取JSON全为Null的问题
读取JSON返回全Null的核心原因是Schema与JSON字段不匹配或JSON格式不符合Spark默认读取规则,以下是具体修复步骤和示例代码:
常见问题排查
- Schema字段不匹配
- Spark默认区分大小写,若JSON字段是
user_id,Schema定义成userId会直接匹配失败返回Null; - 嵌套JSON结构必须对应
StructType嵌套StructField,否则解析失效。
- Spark默认区分大小写,若JSON字段是
- JSON格式不符合Spark规则
- Spark默认要求JSON文件是每行一个独立的JSON对象(行分隔式JSON),如果你的JSON是数组格式(如
[{"a":1},{"b":2}]),必须添加multiLine=True参数才能正确读取。
- Spark默认要求JSON文件是每行一个独立的JSON对象(行分隔式JSON),如果你的JSON是数组格式(如
修改后的可运行代码
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, IntegerType import json # 初始化带Delta Lake支持的SparkSession spark = SparkSession.builder \ .appName("JSONtoDelta") \ .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \ .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \ .getOrCreate() # 生成符合Spark要求的JSON数据(每行一个对象) sample_data = [ {"user_id": 1, "username": "alice", "age": 30}, {"user_id": 2, "username": "bob", "age": 25} ] # 写入文件时确保每行一个JSON对象,而非数组 with open("sample_data.json", "w") as f: for item in sample_data: json.dump(item, f) f.write("\n") # 定义与JSON完全匹配的Schema user_schema = StructType([ StructField("user_id", IntegerType(), nullable=False), StructField("username", StringType(), nullable=False), StructField("age", IntegerType(), nullable=True) ]) # 读取JSON文件(若JSON是数组格式,将multiLine改为true) df = spark.read \ .schema(user_schema) \ .option("multiLine", "false") \ .json("sample_data.json") # 验证数据是否正常读取 df.show() df.printSchema() # 写入Delta Lake表 df.write \ .format("delta") \ .mode("overwrite") \ .save("./delta_user_table")
验证步骤
- 运行
df.printSchema()确认Schema与JSON字段完全对应; - 运行
df.show()查看是否有非Null数据输出。
二、临时目录删除失败的警告与错误解析
原因
- 文件被占用:本地开发环境中,Spark临时目录(默认是
/tmp/spark-*或Windows的C:\Temp\spark-*)被文件资源管理器、终端或其他进程锁定,导致Spark无法删除; - 权限不足:Spark运行用户没有临时目录的写入/删除权限,比如Linux下
/tmp权限被限制,或Windows下无管理员权限; - 作业异常中断:Spark作业中途崩溃,临时文件未被正常清理。
解决办法
- 手动清理临时目录:通过
spark.conf.get("spark.local.dir")查看当前临时目录路径,手动删除其中的spark-*文件夹; - 自定义临时目录:在SparkSession初始化时配置有权限的自定义临时目录:
spark = SparkSession.builder \ .appName("JSONtoDelta") \ .config("spark.local.dir", "/path/to/your/custom/temp/dir") # 替换为有权限的本地路径 .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \ .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \ .getOrCreate() - 释放文件锁:关闭打开临时目录的文件管理器、终端窗口,解除文件占用;
- Linux权限调整:测试环境下可给临时目录添加读写权限:
chmod -R 777 /path/to/temp/dir(生产环境需谨慎操作)
内容的提问来源于stack exchange,提问作者NO2 SIIZEXL
相关产品推荐
相关产品推荐

