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

PySpark读取Blob嵌套JSON异常:路径过滤字段为空,指定文件路径正常

问题分析与解决方案

你的问题核心是Spark自动推断Schema时,根路径下存在Schema不一致的文件,导致合并后部分字段被置为Null。当指定单个IngestionDate目录时,Spark仅基于该目录内的文件推断Schema,结构一致所以能正常读取。

解决思路

直接读取最近7天对应的IngestionDate分区目录,而非读取根目录后再过滤,既避免Schema冲突,又提升读取效率。

具体实现方法

方法1:生成最近7天的分区路径(适合已知日期格式的场景)

from pyspark.sql import functions as F
from datetime import datetime, timedelta

# 根路径
root_path = "/mnt/.../FileFormat=json/SchemaId=v0/"

# 生成最近7天的日期字符串(需和目录中的日期格式一致,比如yyyy-MM-dd)
date_list = [(datetime.today() - timedelta(days=i)).strftime("%Y-%m-%d") for i in range(7)]

# 构建目标分区路径
target_paths = [f"{root_path}/IngestionDate={date}" for date in date_list]

# 读取指定路径的文件
df_raw = spark.read.json(target_paths)

方法2:扫描目录并过滤最近7天的分区(适合动态获取存在的日期目录)

如果部分日期可能没有数据,可先扫描所有IngestionDate分区,再筛选最近7天的存在路径:

from pyspark.sql import functions as F
from datetime import datetime, timedelta
from py4j.java_gateway import java_import

# 获取Hadoop文件系统实例
java_import(spark._jvm, "org.apache.hadoop.fs.Path")
fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(spark._jsc.hadoopConfiguration())

root_path = "/mnt/.../FileFormat=json/SchemaId=v0/"
root_path_obj = spark._jvm.Path(root_path)

# 扫描所有IngestionDate分区目录
date_dirs = []
for status in fs.listStatus(root_path_obj):
    if status.isDirectory():
        dir_name = status.getPath().getName()
        if dir_name.startswith("IngestionDate="):
            date_str = dir_name.split("=")[1]
            date_dirs.append(date_str)

# 筛选最近7天的日期
today = datetime.today().date()
recent_dates = [
    d for d in date_dirs 
    if (today - datetime.strptime(d, "%Y-%m-%d").date()).days <= 7
]

# 构建目标路径并读取
target_paths = [f"{root_path}/IngestionDate={d}" for d in recent_dates]
df_raw = spark.read.json(target_paths)

额外优化:指定固定Schema

如果你的JSON结构是固定的,可提前定义Schema,彻底避免自动推断的问题:

from pyspark.sql.types import StructType, StructField, StringType, IntegerType, TimestampType

# 根据实际JSON结构定义Schema
custom_schema = StructType([
    StructField("field1", StringType(), nullable=True),
    StructField("field2", IntegerType(), nullable=True),
    StructField("IngestionDate", StringType(), nullable=True),
    # 其他字段...
])

# 读取时指定Schema
df_raw = spark.read.schema(custom_schema).json(target_paths)

原理说明

  • 原方法先读取根目录下所有文件,Spark会合并所有文件的Schema,若存在字段缺失/类型不一致的文件,会导致部分字段被填充为Null。
  • 直接读取目标分区目录,Spark仅基于这些目录内的文件推断Schema,保证结构一致性;同时减少了不必要的文件扫描,提升读取性能。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 17:02:59