PySpark动态传递列名至when条件,判断列空值的问题
解决PySpark多列空值判断的报错问题
错误原因
你触发的TypeError: col() takes 1 positional argument but 2 were given,是因为col()函数仅支持传入单个列名作为参数,而你通过*primary_key_columns解包列表一次性传入了多个列名,导致参数数量不匹配。
正确实现方式
要实现「列表中任意列存在空值则赋值为delete,否则赋值为no delete」的逻辑,需要动态拼接多列的空值判断条件,以下是两种可行方案:
方案1:适配任意长度的列列表(推荐)
借助functools.reduce动态拼接所有列的isNull()判断,用逻辑或(|)连接,适合列数量不确定的场景:
from pyspark.sql.functions import col, when from functools import reduce primary_key_columns = ['id', 'label'] # 生成所有列的空值判断条件,并拼接为逻辑或关系 null_condition = reduce( lambda acc, col_name: acc | col(col_name).isNull(), primary_key_columns[1:], # 从第二个列开始迭代 col(primary_key_columns[0]).isNull() # 初始条件为第一个列的空值判断 ) df_new = df_new.withColumn("action", when(null_condition, "delete").otherwise("no delete"))
方案2:固定列数量的直接写法
如果列列表长度固定且数量较少,可直接手动拼接条件:
from pyspark.sql.functions import col, when primary_key_columns = ['id', 'label'] df_new = df_new.withColumn( "action", when(col("id").isNull() | col("label").isNull(), "delete").otherwise("no delete") )
内容的提问来源于stack exchange,提问作者sharma_re
相关产品推荐
相关产品推荐

