PySpark如何编写基于两列条件填充新列的UDF
错误原因
- 你没有将自定义函数注册为PySpark UDF,直接传入Column对象到普通Python函数中运行,Python原生的if-else无法判断PySpark Column类型的布尔值,触发报错
- 你的逻辑里搞反了col1和col2的判断规则:需求是判断col2是否包含分号、根据col1的取值切分,你的原代码反过来判断col1是否包含分号,参数传递也顺序错误
推荐方案:使用PySpark内置函数实现(性能远高于自定义UDF)
直接用split+when组合即可实现需求,不需要自定义UDF:
import pyspark.sql.functions as F df = df.withColumn( "col3", # 先判断col2是否包含分号 F.when( F.col("col2").contains(";"), # 再根据col1的取值取分号前后的内容 F.when(F.col("col1") == 1, F.split(F.col("col2"), ";").getItem(0)) .when(F.col("col1") == 2, F.split(F.col("col2"), ";").getItem(1)) # 不包含分号直接返回col2 ).otherwise(F.col("col2")) )
自定义UDF实现方案(如果必须用UDF的场景)
如果确需使用自定义UDF,需要先注册UDF,且函数内直接处理每行的具体值而非Column对象:
import pyspark.sql.functions as F from pyspark.sql.types import StringType # 定义处理单行数据的Python函数,参数是每行对应字段的实际值 def split_by_semicolon(col1_val, col2_val): if ";" in col2_val: parts = col2_val.split(";") if col1_val == 1: return parts[0] elif col1_val == 2: return parts[1] return col2_val # 注册为PySpark UDF,指定返回值类型 split_udf = F.udf(split_by_semicolon, StringType()) # 调用UDF生成col3 df = df.withColumn("col3", split_udf(F.col("col1"), F.col("col2")))
内容的提问来源于stack exchange,提问作者chris_to
相关产品推荐
相关产品推荐

