PySpark中when条件调用UDF报错:TypeError(NoneType相乘)
问题:UDF调用触发NoneType乘法错误
我定义了rate_calc UDF用于计算利率,在transformation_d1_gp函数中通过when...otherwise逻辑调用该UDF,预期满足when条件时返回0,否则执行UDF。但实际触发PythonException,报错TypeError: unsupported operand type(s) for *: 'NoneType' and 'NoneType'。若将otherwise设为0测试则无报错,求解决。
报错信息
PythonException: An exception was thrown from the Python worker. Please see the stack trace below. Traceback (most recent call last): File "<ipython-input-40-be527cb47980>", line 11, in rate_calc TypeError: unsupported operand type(s) for *: 'NoneType' and 'NoneType'
相关代码
自定义UDF代码
@udf (DoubleType()) def rate_calc(periodo ,valor_par,valor_ven, valor_cli): if (periodo * valor_par) == valor_ven : return 0.000000 else: return float(npf.rate(periodo, valor_par, valor_ven, valor_cli))
数据转换函数代码
def transformation_d1_gp(df_input: DataFrame) -> DataFrame: df_output = (df_input.withColumn('PROD_PLANO_VLR',sf.col('PARCELA_CLIENTE')*sf.col('VALOR_VENDA_PEDIDO')) .groupBy('GRUPO_PROMOCIONAL_NOVO').agg(sum( sf.col('VALOR_VENDA_PEDIDO')).alias("SUM_VALOR_VENDA"),sum(sf.col( 'VALOR_PG_CLIENTE')).alias("SUM_PG_CLIENTE"),sum(sf.col('PROD_PLANO_VLR')).alias('PROD_PLANO_VLR')) .withColumn('PLANO_MEDIO',(sf.col('PROD_PLANO_VLR')/sf.col('SUM_VALOR_VENDA')).cast(DoubleType())) .withColumn('NEGATIVE_SALES', -(sf.col('SUM_VALOR_VENDA')).cast(DoubleType())) .withColumn('PG_PERIODO', (sf.col('SUM_PG_CLIENTE')/sf.col('PLANO_MEDIO')).cast(DoubleType())) .withColumn('REBATE_R$', (sf.col('SUM_PG_CLIENTE')-sf.col('SUM_VALOR_VENDA')).cast(DoubleType())) .withColumn('REBATE_%', (sf.col('REBATE_R$')/sf.col('SUM_PG_CLIENTE')).cast(DoubleType())) .withColumn('MIX_VENDAS', (sf.col('SUM_VALOR_VENDA')/sf.sum('SUM_VALOR_VENDA').over(Window.partitionBy())).cast(DoubleType())) .withColumn('TAXA', when((sf.col('PLANO_MEDIO').isNull()) | (sf.col('PLANO_MEDIO') == 0),0) .when((sf.col('PG_PERIODO').isNull()) | (sf.col('PG_PERIODO') == 0),0) .when((sf.col('NEGATIVE_SALES') == -0) | (sf.col('NEGATIVE_SALES').isNull()), 0) .otherwise(rate_calc('PLANO_MEDIO','PG_PERIODO','NEGATIVE_SALES',sf.lit(0.0)))) .select('GRUPO_PROMOCIONAL_NOVO','SUM_VALOR_VENDA','MIX_VENDAS','REBATE_R$','PLANO_MEDIO','REBATE_%','TAXA','NEGATIVE_SALES')) return df_output
解决方案
1. 修正UDF传参方式
调用UDF时必须传递Column对象而非列名字符串,否则UDF会接收字符串常量而非列的实际值。修改TAXA列的UDF调用代码:
.otherwise(rate_calc(sf.col('PLANO_MEDIO'), sf.col('PG_PERIODO'), sf.col('NEGATIVE_SALES'), sf.lit(0.0)))
2. UDF内部添加参数合法性检查
即使when条件覆盖大部分场景,仍可能存在漏网的None或NaN值,在UDF内部提前判空并处理:
@udf(DoubleType()) def rate_calc(periodo, valor_par, valor_ven, valor_cli): # 检查关键参数是否为None if periodo is None or valor_par is None or valor_ven is None: return 0.0 # 检查参数是否为有效数值(排除NaN) if not isinstance(periodo, (int, float)) or not isinstance(valor_par, (int, float)) or not isinstance(valor_ven, (int, float)): return 0.0 if (periodo * valor_par) == valor_ven: return 0.000000 # 包裹npf.rate调用,捕获计算异常 try: return float(npf.rate(periodo, valor_par, valor_ven, valor_cli)) except Exception: return 0.0
3. 完善when条件,覆盖NaN场景
Spark中NaN不等于任何值(包括自身),isNull()无法识别NaN,需用isnan()补充判断:
.withColumn('TAXA', when((sf.col('PLANO_MEDIO').isNull()) | (sf.col('PLANO_MEDIO') == 0) | sf.isnan(sf.col('PLANO_MEDIO')), 0) .when((sf.col('PG_PERIODO').isNull()) | (sf.col('PG_PERIODO') == 0) | sf.isnan(sf.col('PG_PERIODO')), 0) .when((sf.col('NEGATIVE_SALES') == 0) | (sf.col('NEGATIVE_SALES').isNull()) | sf.isnan(sf.col('NEGATIVE_SALES')), 0) .otherwise(rate_calc(sf.col('PLANO_MEDIO'), sf.col('PG_PERIODO'), sf.col('NEGATIVE_SALES'), sf.lit(0.0))))
注:NEGATIVE_SALES == -0等价于NEGATIVE_SALES == 0,可简化写法。
内容的提问来源于stack exchange,提问作者Vivian
相关产品推荐
相关产品推荐

