Spark中如何在同一条withColumn链式调用里引用新列并检查Null?
问题解决:Spark DataFrame链式withColumn中引用新创建列的正确方式
错误原因分析
- 第一种写法错误:
"test".isNull()是对字符串调用isNull()方法,而非Spark的Column对象,因此会触发'str' object has no attribute 'isNull'错误。 - 第二种写法错误:
table_df_flat.test引用的是原始DataFrame的列,而test列是在链式调用的第一步才创建的,原始DF中不存在该列,所以报'DataFrame' object has no attribute 'test'错误。
正确写法
你需要用Spark的col()函数(或f.col(),取决于导入时的别名)来引用刚创建的test列,让Spark识别这是DataFrame中的列而非普通字符串。
写法一:使用col()引用新列
import pyspark.sql.functions as f table_df_flat \ .withColumn("test", f.lit(None)) \ .withColumn("yes", f.when(f.col("test").isNull(), "yes"))
写法二:使用expr()表达式(复杂逻辑场景适用)
如果需要更贴近SQL风格的写法,也可以用expr()实现:
table_df_flat \ .withColumn("test", f.lit(None)) \ .withColumn("yes", f.expr("CASE WHEN test IS NULL THEN 'yes' END"))
可选补充:设置非Null场景的默认值
如果要给test不为Null的情况指定默认值,可搭配otherwise():
table_df_flat \ .withColumn("test", f.lit(None)) \ .withColumn("yes", f.when(f.col("test").isNull(), "yes").otherwise("no"))
内容的提问来源于stack exchange,提问作者Blue Clouds
相关产品推荐
相关产品推荐

