Databricks中PySpark DataFrame链式调用withColumn报错咨询
Spark DataFrame链式调用列不存在问题解析
问题原因
分步调用时,每次执行df1 = df1.withColumn(...)后,df1变量指向的是已经新增了列的新DataFrame,所以第二次调用时df1["new1"]能找到刚添加的列。
但链式调用里,df1.withColumn("new1", func1(df1["old1"])).withColumn("new2", func2(df1["new1"]))中,所有对df1的引用都是指向原始的未修改的DataFrame——因为整个表达式在执行时,df1["new1"]会先被求值,而此时withColumn("new1")返回的新DataFrame还没赋值给df1变量,自然找不到new1列。
正确的链式调用写法
要在链式调用中引用上一步新增的列,不能用原DataFrame变量(比如df1["new1"]),而是直接用列名的Column对象,具体写法如下:
from pyspark.sql.functions import col df1 = df1.withColumn("new1", func1(col("old1"))).withColumn("new2", func2(col("new1")))
补充说明
Spark的DataFrame是不可变对象,所有withColumn这类操作都会返回一个新的DataFrame,而不是修改原对象。分步调用通过重新赋值变量,让后续操作基于新的DataFrame;链式调用则需要直接引用新DataFrame中的列,而非原变量的引用。
内容的提问来源于stack exchange,提问作者Niehm
相关产品推荐
相关产品推荐

