如何将Pandas的Lambda自定义行处理逻辑迁移到Pyspark实现
PySpark实现目标列col6的两种解决方案
首先先完成左连接得到new_df的基础步骤,和Pandas逻辑对齐:
from pyspark.sql import functions as F # 关联键on参数根据实际业务替换即可,这里假设两张表的关联字段为join_key new_df = df1.join(df2, on="join_key", how="left")
方案1:Spark原生函数实现(推荐,性能远高于自定义UDF)
不需要写自定义UDF,直接用内置函数即可完全匹配需求逻辑,避免UDF的序列化和类型匹配问题:
new_df = new_df.withColumn( "col6", F.coalesce( F.when(F.col("col5") != 0, F.col("col5")), F.when(F.col("col4") != 0, F.col("col4")), F.when(F.col("col3") != 0, F.col("col3")), F.when(F.col("col2") != 0, F.col("col2")), F.lit("NOT FOUND") ).cast("string") # 统一转字符串类型避免数值和字符串混合报错 )
逻辑说明:when函数满足条件时返回对应列值,不满足则返回null,coalesce会按顺序返回第一个非null的值,正好匹配从col5到col2优先取第一个非0值的逻辑,四列全为0时所有when都返回null,最终返回固定值NOT FOUND。
方案2:自定义UDF的正确写法
如果你一定要用自定义函数实现,之前报错大概率是返回值类型不匹配或者传参方式错误,正确写法如下:
from pyspark.sql.types import StringType # 编写自定义逻辑函数 def get_final_au(col5, col4, col3, col2): for val in [col5, col4, col3, col2]: # 额外判空避免Spark空值触发运算报错 if val is not None and val != 0: return str(val) return "NOT FOUND" # 注册UDF并指定返回值类型为字符串 get_final_au_udf = F.udf(get_final_au, StringType()) # 调用withColumn按顺序传入多列参数 new_df = new_df.withColumn( "col6", get_final_au_udf(F.col("col5"), F.col("col4"), F.col("col3"), F.col("col2")) )
注意:如果col2~col5是数值类型,必须统一转成字符串返回,否则会和最终返回的NOT FOUND字符串类型冲突触发报错
内容的提问来源于stack exchange,提问作者santhosh
相关产品推荐
相关产品推荐

