PySpark新增列时遇到mismatched input expecting EOF及字段不存在报错如何解决
问题根源
- 混淆了Driver端执行的本地代码和Executor端执行的Spark分布式列表达式的执行时机:
'statusBit:' in jsonDF.schema.simpleString()是Driver端运行的Python代码,执行时机远早于Spark作业提交,直接把这段代码写进expr()的SQL字符串、或者和when的Column表达式混写,Spark SQL解析器无法识别,会触发语法错、类型不匹配错误。- 最后尝试的代码无论字段是否存在都引用了
col("statusBit"),Spark在Driver端生成执行计划阶段就会校验Schema是否存在该字段,不存在直接报错,根本轮不到运行when的条件判断逻辑。
最优解决方案(Schema全局固定场景)
绝大多数场景下读取JSON生成的DataFrame Schema是全局统一的,直接在Driver端先做Schema判断,再生成对应的列即可,无额外性能损耗:
from pyspark.sql.functions import col, lit from pyspark.sql.types import StringType # 可替换为实际匹配的字段类型,如IntegerType # Driver端先判断Schema中是否存在statusBit字段 if "statusBit" in [field.name for field in jsonDF.schema.fields]: res_df = jsonDF.withColumn("STATUS_BIT", col("statusBit")) else: res_df = jsonDF.withColumn("STATUS_BIT", lit(None).cast(StringType()))
特殊场景解决方案(每行JSON结构不一致)
如果数据源是每行存储独立JSON字符串、每行结构不统一,读取时未解析为结构化Schema,可直接用get_json_object函数取值,不存在的字段会自动返回null:
from pyspark.sql.functions import get_json_object res_df = jsonDF.withColumn("STATUS_BIT", get_json_object(col("JSON字符串列名"), "$.statusBit"))
内容的提问来源于stack exchange,提问作者Moritz
相关产品推荐
相关产品推荐

