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

PySpark筛选数组长度为1的行报错及字符串列转数组问题

报错原因

  1. JSONDecodeError触发Python worker崩溃:lang列存在空值、空字符串或不符合JSON格式的脏数据,你封装json.loads的UDF没有做异常捕获,遇到脏数据直接抛出异常导致进程退出。
  2. 筛选条件存在语法错误:
    • 位运算符&优先级高于比较运算符==,未加括号会导致运算顺序错误
    • 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 18:45:03