从非空JSON文件选取字段后Spark DataFrame显示为空的问题
PySpark读取JSON返回空DataFrame的排查与解决
问题背景
使用PySpark 3.2.0执行以下操作后返回空DataFrame:
- 初始化SparkSession(已开启大小写敏感配置)
- 手动定义包含
Orderstate字段的Schema - 读取
myfile.json并选取Orderstate字段,执行df.show()无数据返回,手动指定Schema后问题依旧。
相关代码片段:
SparkSession初始化
import sys from pyspark.sql import SparkSession from pyspark.sql.functions import col,explode from pyspark.sql.functions import row_number,max,month, unix_timestamp,floor from pyspark.sql.functions import col,monotonically_increasing_id import logging from pyspark.sql.window import Window from all_items_data_loropiana import items_data_all_orders spark1 = SparkSession.builder\ .appName("book_and_pack")\ .config("spark.mongodb.input.uri", "mongodb://localhost:27017/")\ .config("spark.sql.repl.eagerEval.enabled", True)\ .config('spark.sql.caseSensitive', True)\ .config("inferSchema",True)\ .config("spark.mongodb.input.sampleSize", 50000)\ .config('spark.jars.packages','org.mongodb.spark:mongo-spark-connector_2.12:2.4.2')\ .config("spark.driver.extraJavaOptions", "-Dlog4j.configuration=file:log4j.properties")\ .enableHiveSupport()\ .getOrCreate() logger = logging.getLogger("org.mongodb.spark.MongsoInferSchema") logger.setLevel(logging.ERROR)
手动定义的Schema
schema = StructType([ StructField("orderNumber", StringType(), True), StructField("Orderstate", ArrayType( StructType([ StructField("id", StringType(), True), StructField("state", StringType(), True), StructField("modifiedAt", StructType([StructField("$date", StringType(), True)]), True) ]) ), True) ])
读取JSON的代码
ordersips_df=spark.read.json("myfile.json") df=ordersips_df.select("Orderstate")
可能原因
- JSON文件格式不兼容:PySpark默认要求JSON是每行一个独立对象(JSON Lines格式),如果文件是多行嵌套的完整JSON结构,会导致读取失败或无数据。
- 字段名大小写不匹配:虽然开启了
spark.sql.caseSensitive=True,但如果JSON文件中实际字段名是orderState(首字母小写)或其他大小写变体,会导致选取不到数据。 - 手动Schema与实际数据结构不符:比如
Orderstate的数组元素结构、modifiedAt.$date的类型与实际数据不一致,导致Spark无法解析字段,返回空值。 - 文件路径错误或权限问题:Spark无法找到目标文件,或者没有读取权限,导致返回空DataFrame。
- 原始数据本身无有效内容:所有记录的
Orderstate字段为null或空数组。
解决方案
1. 适配JSON文件格式
如果JSON文件是多行嵌套结构,读取时开启multiLine选项:
# 读取时指定multiLine和手动Schema ordersips_df = spark.read.option("multiLine", True).json("myfile.json", schema=schema) df = ordersips_df.select("Orderstate") df.show()
2. 验证字段名与Schema匹配
先让Spark自动推断Schema,对比手动定义的结构:
# 自动推断Schema并打印 temp_df = spark.read.option("multiLine", True).json("myfile.json") temp_df.printSchema()
根据打印结果调整手动Schema的字段名、类型,确保和实际数据完全一致。
3. 确认文件路径与权限
- 本地环境下,检查文件路径是否为绝对路径,或者相对路径是否正确;
- 集群环境下,确认文件存储路径(如HDFS、S3)的访问权限,使用对应工具验证文件存在性(如
hdfs dfs -ls /path/to/myfile.json)。
4. 检查原始数据内容
验证原始数据中Orderstate字段是否有有效值:
# 统计非空Orderstate的记录数 non_null_count = ordersips_df.filter(col("Orderstate").isNotNull()).count() print(f"非空Orderstate记录数:{non_null_count}") # 查看原始数据的前几行 ordersips_df.show(5, truncate=False)
内容的提问来源于stack exchange,提问作者zied salhi
相关产品推荐
相关产品推荐

