使用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配置,存在编码异常、格式异常的文件会被直接跳过,不会抛出报错
- Spark默认
修复方案
步骤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
相关产品推荐
相关产品推荐

