使用groupBy().applyInPandas()时遇INVALID_PANDAS_UDF错误排查
问题原因与解决方案
你的错误根源是混用了pandas_udf(GROUPED_MAP类型)和applyInPandas的API用法,这两个接口不能搭配使用,具体分析和修复如下:
核心问题
groupBy().applyInPandas() 不需要用@pandas_udf装饰器,它直接接受普通的Python函数;而用@pandas_udf(..., PandasUDFType.GROUPED_MAP)装饰的函数,应该搭配groupBy().apply()调用。你把两种API的用法混在一起,导致Spark的参数校验失败,抛出INVALID_PANDAS_UDF错误。
修复方案(二选一)
方案1:使用applyInPandas(推荐Spark 3.0+)
直接去掉@pandas_udf装饰器,保留函数的参数和返回结构即可:
# 移除@pandas_udf装饰器 def train_and_forecast_prophet_multi_metric(pdf: pd.DataFrame) -> pd.DataFrame: print(f"Type of input pdf: {type(pdf)}") print(f"Columns of input pdf: {pdf.columns.tolist()}") # 后续添加Prophet训练逻辑示例 model = Prophet() model.fit(pdf) # 生成预测周期数据 future = model.make_future_dataframe(periods=7) forecast = model.predict(future) # 合并原数据与预测结果,确保输出符合schema要求 result = pdf.merge(forecast[['ds', 'yhat', 'yhat_lower', 'yhat_upper']], on='ds', how='outer') # 保留分组列attribute result['attribute'] = pdf['attribute'].iloc[0] return result[['attribute', 'y', 'ds', 'yhat', 'yhat_lower', 'yhat_upper']] # 调用applyInPandas并传入输出schema df_result = df_train.groupBy("attribute").applyInPandas(train_and_forecast_prophet_multi_metric, schema=forecast_schema)
方案2:使用groupBy().apply()搭配GROUPED_MAP类型的pandas_udf
保留@pandas_udf装饰器,改用apply()方法调用:
@pandas_udf(forecast_schema, PandasUDFType.GROUPED_MAP) def train_and_forecast_prophet_multi_metric(pdf: pd.DataFrame) -> pd.DataFrame: print(f"Type of input pdf: {type(pdf)}") print(f"Columns of input pdf: {pdf.columns.tolist()}") return pdf # 改用apply()而非applyInPandas df_result = df_train.groupBy("attribute").apply(train_and_forecast_prophet_multi_metric)
额外注意事项
- Schema对齐:如果后续要返回Prophet的预测结果(如
yhat、yhat_lower等),需要提前修改forecast_schema,添加对应的字段(建议用DoubleType类型存储预测值)。 - 数据类型兼容:确保传入Prophet的
ds列是Pandas的datetime64类型,y列是数值类型(原数据的long类型会被Pandas自动转为int64,可直接用于Prophet训练)。 - 分布式环境依赖:EC2集群的每个Worker节点都需要安装Prophet库,否则会出现
ModuleNotFoundError。
内容的提问来源于stack exchange,提问作者Arnab Sinha
相关产品推荐
相关产品推荐

