PySpark中将带K/M单位的货币字符串转换为整数的报错解决
报错原因
你在普通Python自定义函数中混用了PySpark的列操作函数,expr()返回的是PySpark的Column对象,不能直接放在Python原生的if/elif判断里做布尔比较,所以触发了类型转换错误。
推荐解决方案(原生Spark函数实现,无UDF性能损耗)
直接使用Spark SQL的when/otherwise分支逻辑实现转换,不需要自定义函数,执行效率远高于UDF:
from pyspark.sql import functions as F from pyspark.sql.types import IntegerType # 替换df为你自己的DataFrame变量名 df = df.withColumn( "converted_amount", F.when(F.col("ReleaseClause").isNull(), 0) .when(F.substring(F.col("ReleaseClause"), -1, 1) == "K", (F.regexp_replace(F.col("ReleaseClause"), "[^0-9]", "").cast(IntegerType()) * 1000)) .when(F.substring(F.col("ReleaseClause"), -1, 1) == "M", (F.regexp_replace(F.col("ReleaseClause"), "[^0-9]", "").cast(IntegerType()) * 1000000)) .otherwise(F.regexp_replace(F.col("ReleaseClause"), "[^0-9]", "").cast(IntegerType())) )
代码中直接用regexp_replace过滤所有非数字字符(自动去除金额前缀符号、K/M后缀),不需要多次截取字符串,逻辑更简洁。
可选方案:UDF实现
如果需要使用自定义函数,要完全用Python原生的字符串逻辑处理单条数据,不要在UDF内调用Spark列函数:
from pyspark.sql.types import IntegerType from pyspark.sql.functions import udf @udf(returnType=IntegerType()) def convertall(release_clause): if not release_clause: return 0 # 过滤出数字、K、M字符,自动去掉前缀货币符号 clean_str = ''.join([c for c in release_clause if c.isdigit() or c in ('K', 'M')]) if clean_str.endswith('K'): return int(clean_str[:-1]) * 1000 elif clean_str.endswith('M'): return int(clean_str[:-1]) * 1000000 else: return int(clean_str) if clean_str.isdigit() else 0 # 调用UDF新增列 df = df.withColumn("converted_amount", convertall(F.col("ReleaseClause")))
内容的提问来源于stack exchange,提问作者Nishant
相关产品推荐
相关产品推荐

