如何使用PySpark UDF根据time字段规则批量新增3个判断列
PySpark多判断列新增实现方案
这里分两种实现方式,你可以根据实际场景选择:
方案1:修改自定义UDF一次返回多列
调整UDF返回值为三元组,同时指定结构化的返回类型,调用后拆分即可一次性得到3个新列:
from pyspark.sql.types import StructType, StructField, BooleanType, StringType # 定义返回三个判断结果的UDF函数 def time_check(value): col1_res = value > 10 col2_res = value < 0 col3_res = 0 < value < 12 # 如果需要返回字符串格式的'True'/'False',直接把三个值转成str即可 return (col1_res, col2_res, col3_res) # 注册UDF,指定返回类型为包含三个布尔字段的结构体 udf_time_check = udf(time_check, StructType([ StructField("col1", BooleanType(), nullable=False), StructField("col2", BooleanType(), nullable=False), StructField("col3", BooleanType(), nullable=False) ])) # 调用UDF并拆分出三个新列 df = df.withColumn("temp_res", udf_time_check("time")) df = df.select("*", "temp_res.*").drop("temp_res")
方案2:使用PySpark内置函数实现(优先推荐)
自定义UDF存在跨进程序列化开销,数据量大的场景下直接用内置算子性能高5~10倍,无需编写UDF:
from pyspark.sql.functions import col df = df.withColumn("col1", col("time") > 10) \ .withColumn("col2", col("time") < 0) \ .withColumn("col3", (col("time") > 0) & (col("time") < 12))
如果需要返回字符串格式的判断结果,给判断逻辑套上.cast("string")即可,示例:df.withColumn("col1", (col("time") > 10).cast("string"))
内容的提问来源于stack exchange,提问作者sam
相关产品推荐
相关产品推荐

