PySpark where子句可使用不存在列的异常行为技术问询
PySpark 3.1.2列重命名后原列名在where条件中可正常执行的原因解析
现象描述
在Azure Databricks上运行PySpark 3.1.2时,出现以下特殊行为:将DataFrame的列批量重命名(添加_new后缀)后,原列名已不在DataFrame的schema中(c in df.columns返回False),使用df.select(原列名)会触发列不存在的错误,但df.where(col(原列名).isNotNull())却能正常执行,且过滤结果与使用新列名完全一致。
代码与输出示例
代码示例
print(spark.version) df = spark.read.format("csv").option("header", True).load("abfss://some_abfs_path/df.csv") print(type(df), df.columns.__len__(), df.count()) c = df.columns[0] # 重命名前的列名 df = df.select(*[col(x).alias(f"{x}_new") for x in df.columns]) # 为列名添加后缀 print(c in df.columns) try: df.select(c) except: print("SO THIS DOESN'T WORK, WHICH MAKES SENSE.") # 为何此代码能运行: print(df.where(col(c).isNotNull()).count()) # 实际使用的是c_new列 print(df.where(col(f"{c}_new").isNotNull()).count())
输出结果
3.1.2 <class 'pyspark.sql.dataframe.DataFrame'> 102 1226791 False SO THIS DOESN'T WORK, WHICH MAKES SENSE. 1226791 1226791
技术解析
这个现象的核心是PySpark对列的两种解析逻辑差异,以及alias方法的底层实现特性:
alias方法的本质
使用col(x).alias(y)重命名列时,仅修改了列在DataFrame schema中的显示名称(即df.columns返回的名称),但并未改变底层表达式中对应Attribute的原始名称(即原列名x)。这个原始名称会被保留在列的表达式树结构中,不会被alias覆盖。两种列查找逻辑的区别
- 字符串直接查找(如
df.select(c)):这种方式严格匹配DataFrame schema中定义的列名(即重命名后的c_new),由于原列名c不在schema内,因此触发列不存在的错误。 col(c)表达式查找:col(c)会创建一个UnresolvedAttribute对象,Spark在生成执行计划时,会遍历DataFrame所有列的底层表达式,匹配Attribute的原始名称。由于重命名后的列的底层Attribute名称仍是c,因此col(c)会成功匹配到该列(即c_new列),进而正常执行过滤逻辑,结果与使用新列名完全一致。
- 字符串直接查找(如
版本特性说明
该行为是PySpark 3.1.2的设计特性,后续版本对表达式解析逻辑进行了优化,已避免这种容易引发混淆的情况。
总结
若要避免此类混淆,建议重命名列后始终使用新列名进行操作;或在引用列时直接从当前DataFrame的schema中获取列名,而非依赖重命名前的变量。
内容的提问来源于stack exchange,提问作者Roc
相关产品推荐
相关产品推荐

