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

Spark UDF执行报错:SyntaxError: unexpected EOF while parsing

问题解决:PySpark读取DBFS JSON文件时的解析异常

问题描述

使用PySpark读取DBFS上的JSON文件时,通过自定义UDF转换字符串后调用from_json解析指定Schema,执行df.show()或写入DataFrame时抛出如下异常:

PythonException: An exception was thrown from a UDF: 'SyntaxError: unexpected EOF while parsing', from , line 6. Full traceback below:
Traceback (most recent call last):
File "", line 6, in
File "/usr/lib/python3.9/ast.py", line 62, in literal_eval
node_or_string = parse(node_or_string, mode='eval')
File "/usr/lib/python3.9/ast.py", line 50, in parse
return compile(source, filename, mode, flags,
File "", line 0

SyntaxError: unexpected EOF while parsing

涉及代码如下:

import ast
import json
from pyspark.sql import functions as F
from pyspark.sql.types import *

schema = StructType([StructField('in_network', ArrayType(StructType([StructField('billing_code', StringType(), True), StructField('billing_code_type', StringType(), True), StructField('billing_code_type_version', StringType(), True), StructField('description', StringType(), True), StructField('name', StringType(), True), StructField('negotiated_rates', ArrayType(StructType([StructField('negotiated_prices', ArrayType(StructType([StructField('additional_information', StringType(), True), StructField('billing_class', StringType(), True), StructField('billing_code_modifier', ArrayType(StringType(), True), True), StructField('expiration_date', StringType(), True), StructField('negotiated_rate', DoubleType(), True), StructField('negotiated_type', StringType(), True), StructField('service_code', ArrayType(StringType(), True), True)]), True), True), StructField('provider_references', ArrayType(LongType(), True), True)]), True), True), StructField('negotiation_arrangement', StringType(), True)]), True), True), StructField('last_updated_on', StringType(), True), StructField('provider_references', ArrayType(StructType([StructField('provider_group_id', LongType(), True), StructField('provider_groups', ArrayType(StructType([StructField('npi', ArrayType(LongType(), True), True), StructField('tin', StructType([StructField('type', StringType(), True), StructField('value', StringType(), True)]), True)]), True), True)]), True), True), StructField('reporting_entity_name', StringType(), True), StructField('reporting_entity_type', StringType(), True), StructField('version', StringType(), True)])

dict_to_json = F.udf(lambda x: json.dumps(ast.literal_eval(x)))
df = spark.read.text("dbfs:/mnt/transparency/mkjl/in_frt_qSW.json").withColumn("value", F.from_json(dict_to_json("value"), schema)).select("value.*")
df.show()

原因分析

  1. 数据格式不合法:ast.literal_eval()抛出EOF异常,说明部分行的字符串存在格式问题,可能是空行、不完整的字典/JSON结构,或者引号不匹配、括号未闭合等语法错误。
  2. 冗余转换操作:先通过ast.literal_eval()转成字典,再用json.dumps()转回JSON字符串属于多余操作,除非原始数据是Python字典格式(而非标准JSON),否则直接用from_json解析原始文本即可。

解决方案

方案1:直接读取JSON文件(推荐)

如果文件是标准JSON格式(每行一个JSON对象,或整个文件是JSON数组),直接用spark.read.json()读取,无需文本模式+UDF转换:

from pyspark.sql.types import *

# 复用原定义的Schema
schema = StructType([StructField('in_network', ArrayType(StructType([StructField('billing_code', StringType(), True), StructField('billing_code_type', StringType(), True), StructField('billing_code_type_version', StringType(), True), StructField('description', StringType(), True), StructField('name', StringType(), True), StructField('negotiated_rates', ArrayType(StructType([StructField('negotiated_prices', ArrayType(StructType([StructField('additional_information', StringType(), True), StructField('billing_class', StringType(), True), StructField('billing_code_modifier', ArrayType(StringType(), True), True), StructField('expiration_date', StringType(), True), StructField('negotiated_rate', DoubleType(), True), StructField('negotiated_type', StringType(), True), StructField('service_code', ArrayType(StringType(), True), True)]), True), True), StructField('provider_references', ArrayType(LongType(), True), True)]), True), True), StructField('negotiation_arrangement', StringType(), True)]), True), True), StructField('last_updated_on', StringType(), True), StructField('provider_references', ArrayType(StructType([StructField('provider_group_id', LongType(), True), StructField('provider_groups', ArrayType(StructType([StructField('npi', ArrayType(LongType(), True), True), StructField('tin', StructType([StructField('type', StringType(), True), StructField('value', StringType(), True)]), True)]), True), True)]), True), True), StructField('reporting_entity_name', StringType(), True), StructField('reporting_entity_type', StringType(), True), StructField('version', StringType(), True)])

# 直接读取JSON文件并指定Schema
df = spark.read.schema(schema).json("dbfs:/mnt/transparency/mkjl/in_frt_qSW.json")
df.show()

方案2:修复UDF并处理异常数据

如果必须用文本模式读取(比如原始数据是Python字典格式),修改UDF加入异常处理,跳过或标记错误行:

import ast
import json
from pyspark.sql import functions as F
from pyspark.sql.types import *

# 复用原定义的Schema
schema = StructType([StructField('in_network', ArrayType(StructType([StructField('billing_code', StringType(), True), StructField('billing_code_type', StringType(), True), StructField('billing_code_type_version', StringType(), True), StructField('description', StringType(), True), StructField('name', StringType(), True), StructField('negotiated_rates', ArrayType(StructType([StructField('negotiated_prices', ArrayType(StructType([StructField('additional_information', StringType(), True), StructField('billing_class', StringType(), True), StructField('billing_code_modifier', ArrayType(StringType(), True), True), StructField('expiration_date', StringType(), True), StructField('negotiated_rate', DoubleType(), True), StructField('negotiated_type', StringType(), True), StructField('service_code', ArrayType(StringType(), True), True)]), True), True), StructField('provider_references', ArrayType(LongType(), True), True)]), True), True), StructField('negotiation_arrangement', StringType(), True)]), True), True), StructField('last_updated_on', StringType(), True), StructField('provider_references', ArrayType(StructType([StructField('provider_group_id', LongType(), True), StructField('provider_groups', ArrayType(StructType([StructField('npi', ArrayType(LongType(), True), True), StructField('tin', StructType([StructField('type', StringType(), True), StructField('value', StringType(), True)]), True)]), True), True)]), True), True), StructField('reporting_entity_name', StringType(), True), StructField('reporting_entity_type', StringType(), True), StructField('version', StringType(), True)])

def convert_to_json(x):
    if not x:
        return None
    try:
        return json.dumps(ast.literal_eval(x))
    except (SyntaxError, ValueError):
        # 异常行返回None,后续过滤掉
        return None

dict_to_json = F.udf(convert_to_json, StringType())

df = spark.read.text("dbfs:/mnt/transparency/mkjl/in_frt_qSW.json") \
    .withColumn("json_str", dict_to_json("value")) \
    .filter(F.col("json_str").isNotNull()) \
    .withColumn("value", F.from_json("json_str", schema)) \
    .select("value.*")

df.show()

方案3:检查并清理源数据

先定位异常行,确认问题后清理源数据:

def check_valid(x):
    if not x:
        return "empty_line"
    try:
        ast.literal_eval(x)
        return "valid"
    except SyntaxError:
        return "invalid_syntax"
    except Exception as e:
        return f"error: {str(e)}"

check_udf = F.udf(check_valid, StringType())

# 标记每行数据的状态
df_check = spark.read.text("dbfs:/mnt/transparency/mkjl/in_frt_qSW.json") \
    .withColumn("status", check_udf("value"))

# 统计各状态行数
df_check.groupBy("status").count().show()

# 查看异常行的具体内容
df_check.filter(F.col("status") != "valid").show(truncate=False)

根据结果清理源数据(比如删除空行、修复格式错误行)后,再重新执行解析逻辑。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 05:50:37