PySpark筛选数组长度为1的行报错及字符串列转数组问题
报错原因
JSONDecodeError触发Python worker崩溃:lang列存在空值、空字符串或不符合JSON格式的脏数据,你封装json.loads的UDF没有做异常捕获,遇到脏数据直接抛出异常导致进程退出。- 筛选条件存在语法错误:
- 位运算符
&优先级高于比较运算符==,未加括号会导致运算顺序错误 array_contains函数参数传错,第二个参数是要匹配的值,不应放在F.col的入参中
- 位运算符
解决方案
步骤1:安全转换lang列为数组(优先使用Spark原生函数,性能远高于Python UDF)
不需要自己写UDF调用json.loads,直接用Spark内置的from_json函数即可自动处理异常:
from pyspark.sql import functions as F from pyspark.sql.types import ArrayType, StringType # 定义目标数组类型schema lang_schema = ArrayType(StringType()) # 转换lang列,无效JSON/空值会自动返回null,不会抛出异常 df = df.withColumn("lang_arr", F.from_json(F.col("lang"), lang_schema))
如果确实需要使用UDF,必须增加异常捕获逻辑:
import json from pyspark.sql.functions import udf from pyspark.sql.types import ArrayType, StringType @udf(returnType=ArrayType(StringType())) def safe_parse_lang(s): # 处理空值、空字符串场景 if not s: return None try: return json.loads(s) except json.JSONDecodeError: # 脏数据返回null,也可根据业务需求返回空数组[] return None df = df.withColumn("lang_arr", safe_parse_lang(F.col("lang")))
步骤2:正确编写筛选逻辑
给逻辑运算加括号、修正array_contains参数:
# 筛选仅含EN的行,同时可过滤转换失败的null值 result_df = df.filter( F.col("lang_arr").isNotNull() & (F.size(F.col("lang_arr")) == 1) & F.array_contains(F.col("lang_arr"), "EN") )
内容的提问来源于stack exchange,提问作者eun ji
相关产品推荐
相关产品推荐

