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

PySpark读取JSON写入Delta Lake全为Null的问题解决与日志解析

PySpark读取JSON全为Null及Delta Lake临时目录错误的解决方案

一、解决读取JSON全为Null的问题

读取JSON返回全Null的核心原因是Schema与JSON字段不匹配或JSON格式不符合Spark默认读取规则,以下是具体修复步骤和示例代码:

常见问题排查

  1. Schema字段不匹配
    • Spark默认区分大小写,若JSON字段是user_id,Schema定义成userId会直接匹配失败返回Null;
    • 嵌套JSON结构必须对应StructType嵌套StructField,否则解析失效。
  2. JSON格式不符合Spark规则
    • Spark默认要求JSON文件是每行一个独立的JSON对象(行分隔式JSON),如果你的JSON是数组格式(如[{"a":1},{"b":2}]),必须添加multiLine=True参数才能正确读取。

修改后的可运行代码

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数据输出。

二、临时目录删除失败的警告与错误解析

原因

  1. 文件被占用:本地开发环境中,Spark临时目录(默认是/tmp/spark-*或Windows的C:\Temp\spark-*)被文件资源管理器、终端或其他进程锁定,导致Spark无法删除;
  2. 权限不足:Spark运行用户没有临时目录的写入/删除权限,比如Linux下/tmp权限被限制,或Windows下无管理员权限;
  3. 作业异常中断: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 15:52:43