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

Spark 2.4.7中未缓存DataFrame时Pandas UDF无法识别数组输入问题

Spark分组聚合Pandas UDF与标量UDF协同运行报错问题

问题详情

  • 运行环境:Spark 2.4.7,PyArrow 12.0.1
  • 场景说明:pack_ts是分组聚合型Pandas UDF,输出数组类型字段;remove_outlier通过标量Pandas UDFcheck_outlier处理该数组字段,实现异常值过滤
  • 异常现象:两个模块单独运行均正常,但协同执行时,必须保留remove_outlier中的ky_his.cache()语句才能成功运行,移除缓存会触发代码生成错误

复现代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import pandas_udf, col
from pyspark.sql.types import (
    ArrayType, IntegerType, BooleanType,
    StructType, StructField, StringType, TimestampType
)
import pandas as pd

spk = SparkSession.builder.appName("UDFIssue").getOrCreate()

@pandas_udf(ArrayType(IntegerType()), pandas_udf.PandasUDFType.GROUPED_AGG)
def pack_ts(date: pd.Series, pv: pd.Series) -> pd.Series:
    df = pd.DataFrame({"date": date, "pv": pv})
    df = df.sort_values(by="date")
    df['pv'] = df['pv'].fillna(0)
    return df["pv"]


@pandas_udf(BooleanType(), pandas_udf.PandasUDFType.SCALAR)
def check_outlier(his_series: pd.Series) -> pd.Series:
    def func(his):
        pv = pd.Series(his)
        pre = pv.shift(1)
        cond = (pv < 0.1 * pre) & (pre > 1000)
        return cond.any()

    return his_series.apply(func)


def remove_outlier(ky_his):
    ky_his = ky_his.withColumn("outlier", check_outlier(ky_his.his))
    ky_his.cache() # NOTE 必须保留该行才能正常运行
    ky_his = ky_his.filter(~ky_his.outlier).drop("outlier")
    return ky_his


# 定义Schema
schema = StructType([
    StructField("ky", StringType(), True),
    StructField("dt", StringType(), True),
    StructField("pv", IntegerType(), True)
])

# 构造测试数据
data = [
    ("a", "2020-01-01", 1),
    ("a", "2020-01-02", 20000),
    ("a", "2020-01-03", 3),
]

# 创建DataFrame并执行分组聚合
df = spk.createDataFrame(data, schema)
df = df.groupby("ky").agg(pack_ts(col("dt").cast(TimestampType()), df.pv).alias("his"))
df.show()

# 注册UDF并执行异常值过滤
spk.udf.register("check_outlier", check_outlier)
df = remove_outlier(df)
df.show()

分析与解决方案

问题根源

Spark 2.4.x版本对Pandas UDF的链式执行逻辑支持存在局限性,加上PyArrow 12.0.1属于较新版本,与旧版Spark的集成存在兼容性缺口。在未缓存的情况下,Spark的代码生成器无法正确解析两个UDF之间的依赖关系,导致逻辑执行失败;而缓存操作会强制将中间DataFrame的计算结果物化到内存/磁盘,后续阶段直接读取物化后的结果,绕过了重复解析UDF逻辑的步骤,因此能正常运行。

解决办法

  1. 临时方案:保留cache()语句,确保中间结果被物化,这是当前环境下最直接的可行方案
  2. 长期方案:升级Spark版本至3.x系列,Spark 3.x对Pandas UDF和PyArrow的集成进行了全面优化,从根本上解决这类兼容性问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 23:47:38