Pandas UDF运行报错:返回DataFrame列数与指定Schema不匹配求解决
Pandas UDF列数不匹配错误排查与解决
问题代码
Spark UDF定义
def some_udf(df, keys = IDS_COLS, cols_to_keep = COLS_TO_KEEP): INT_ID_COLUMN = '__iid' df_keys = df.select(keys).distinct().withColumn(INT_ID_COLUMN, F.monotonically_increasing_id()) df = df.join(df_keys, keys) pred_df_schema = df.select(INT_ID_COLUMN, *keys, 'tx_id', F.lit(0.0).alias('score'), 'class').schema pred_with_some_udf = F.pandas_udf( some_func, returnType=pred_df_schema, functionType=F.PandasUDFType.GROUPED_MAP ) prediction_df = df.select(INT_ID_COLUMN, 'tx_id', 'class', *keys ,*cols_to_keep) \ .groupby(INT_ID_COLUMN) \ .apply(pred_with_some_udf) \ .drop(INT_ID_COLUMN) return prediction_df
对应的Pandas处理函数
def some_func(df): ... return df[[*keys, 'tx_id', 'score', 'class']]
错误信息
'RuntimeError: Number of columns of the returned pandas.DataFrame doesn't match specified schema. Expected: 6 Actual: 5'
问题原因
你定义的pred_df_schema包含了分组键INT_ID_COLUMN,但some_func返回的DataFrame里并没有这个列,导致返回列数比schema要求少1(预期6列,实际返回5列)。
解决方法
方法一:调整schema,移除分组键
GROUPED_MAP类型的Pandas UDF不需要在返回schema中包含分组键,Spark会自动将分组列合并到最终结果中。修改schema定义:
# 去掉INT_ID_COLUMN,只保留需要返回的业务列 pred_df_schema = df.select(*keys, 'tx_id', F.lit(0.0).alias('score'), 'class').schema
方法二:在Pandas函数中返回分组键
确保some_func返回的DataFrame包含INT_ID_COLUMN,与schema列数匹配:
def some_func(df): ... # 加入INT_ID_COLUMN到返回列中 return df[[INT_ID_COLUMN, *keys, 'tx_id', 'score', 'class']]
额外注意事项
- 确认
keys的列数:如果IDS_COLS包含2列,那么*keys会扩展为2列,加上tx_id、score、class共5列,加上INT_ID_COLUMN正好是6列,与预期一致。 - 确保返回列的顺序与schema定义的顺序完全一致,Spark会严格按顺序匹配列,即使列名相同顺序不同也可能引发异常。
内容的提问来源于stack exchange,提问作者Lior T
相关产品推荐
相关产品推荐

