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

使用PySpark读取Databricks DBFS中JSON文件仅能读取部分文件问题咨询

问题根因定位

首先分两种场景排查:PySpark 未扫描到全量文件、PySpark 扫描到文件但解析阶段丢弃数据

  • 场景1:未扫描到全量文件
    • 默认路径匹配规则未覆盖所有子目录下的.json后缀文件,未明确指定过滤规则时,Spark可能跳过隐藏文件、无后缀文件,或因通配符规则未匹配到深层目录文件
    • 若路径存在key=value格式的分区目录结构,分区推断逻辑与recursiveFileLookup配置存在冲突,导致部分目录被排除
  • 场景2:扫描到文件但解析阶段丢弃
    • Spark默认multiLine=False,仅按行解析JSON,若文件为单文件单个JSON对象/JSON数组的跨行格式,会判定为无效内容直接丢弃
    • 默认采样推断的Schema与全量文件Schema不匹配,字段类型、字段存在性差异会导致不符合Schema的行被静默丢弃
    • 若开启了spark.sql.files.ignoreCorruptFiles=true配置,存在编码异常、格式异常的文件会被直接跳过,不会抛出报错
修复方案

步骤1:确认扫描范围是否正常

执行以下代码查看Spark实际识别到的文件数量,判断是扫描问题还是解析问题:

raw_df = spark.read.json('/mnt/testdatabricks/metrics-raw/', recursiveFileLookup=True)
print("Spark识别到的文件数量:", len(raw_df.inputFiles()))

步骤2:修复扫描范围异常

直接传入已验证的全量文件路径列表,规避路径匹配规则问题:

from glob import glob
# 读取全量.json文件路径(与你Pandas验证的路径逻辑一致)
local_json_paths = glob('/dbfs/mnt/testdatabricks/metrics-raw/**/*.json', recursive=True)
# 转换为Spark支持的DBFS路径格式
dbfs_json_paths = [path.replace('/dbfs', '') for path in local_json_paths]

步骤3:修复解析丢弃问题

调整JSON读取配置,保留异常记录便于排查,适配多场景JSON格式:

raw_df = spark.read.json(
    dbfs_json_paths,
    multiLine=True, # 支持跨行的单文件JSON格式
    mode="PERMISSIVE", # 解析失败的内容存入异常字段,不直接丢弃
    columnNameOfCorruptRecord="corrupt_record",
    encoding="UTF-8" # 可根据实际文件编码调整,比如GBK
)
# 查看是否存在解析异常的记录
print("解析异常记录数:", raw_df.filter("corrupt_record is not null").count())

步骤4:Schema兼容优化

避免采样Schema与全量数据不匹配导致的数据丢失:

# 先用全量文件推断统一Schema
schema = spark.read.json(dbfs_json_paths, samplingRatio=1.0, multiLine=True).schema
# 用统一Schema读取全量数据
raw_df = spark.read.json(dbfs_json_paths, schema=schema, multiLine=True)
Pandas全量读取报错修复

Pandas运行在Driver节点,全量读取会把所有数据加载到Driver内存导致OOM,改为分布式并行读取即可:

from pyspark.sql.types import StructType
import pandas as pd

# 采样获取文件Schema
sample_pd = pd.read_json(local_json_paths[0])
spark_schema = StructType.fromJson(sample_pd._schema.json())

# 并行读取所有文件
raw_df = spark.createDataFrame(
    spark.sparkContext.parallelize(local_json_paths)
    .map(lambda p: pd.read_json(p).to_dict("records"))
    .flatMap(lambda x: x),
    schema=spark_schema
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 19:36:03