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逻辑的步骤,因此能正常运行。
解决办法
- 临时方案:保留
cache()语句,确保中间结果被物化,这是当前环境下最直接的可行方案 - 长期方案:升级Spark版本至3.x系列,Spark 3.x对Pandas UDF和PyArrow的集成进行了全面优化,从根本上解决这类兼容性问题
内容的提问来源于stack exchange,提问作者Amadeus
相关产品推荐
相关产品推荐

