PySpark关联含多同名列的表后val3过滤报错的解决办法
解决Spark关联后列名歧义的过滤问题
问题场景
你有两个DataFrame df_1 和 df_2,包含同名列val1、val2、val3,数据如下:
df_1:
+---+----+----+----+----+ | id|val1|val2|val3|val4| +---+----+----+----+----+ | 1| 2| 2| a| 54| | 1| 2| 2| b| 56| | 2| 3| 4| c| 12| +---+----+----+----+----+
df_2:
+---+----+----+----+----+ | id|val1|val2|val3|val5| +---+----+----+----+----+ | 1| 2| 2| aa| %%| | 1| 2| 2| bb| $#| | 2| 3| 4| cc| @&| +---+----+----+----+----+
执行左关联后:
df = df_1.join(df_2, ["val1", "val2"], "left")
生成的新DataFrame包含两个val3列:
+---+----+----+----+----+----+----+ | id|val1|val2|val3|val4|val3|val5| +---+----+----+----+----+----+----+ | 1| 2| 2| a| 54| aa| %%| | 1| 2| 2| b| 56| bb| $#| | 2| 3| 4| c| 12| cc| @&| +---+----+----+----+----+----+----+
此时执行过滤:
df = (df.where(F.col(val3) == "b"))
会抛出错误:
AnalysisException: Reference 'val3' is ambiguous, could be: val3, val3.
解决方案
方案1:关联前重命名重复列
提前将其中一个DataFrame的重复列改名,从根源避免歧义:
# 将df_2的val3重命名为val3_df2 df_2_renamed = df_2.withColumnRenamed("val3", "val3_df2") # 执行关联 df = df_1.join(df_2_renamed, ["val1", "val2"], "left") # 过滤时直接引用明确的列名 df = df.where(F.col("val3") == "b")
方案2:关联时通过原始DataFrame引用列
如果原始DataFrame仍在会话中,可以直接通过原始DataFrame的引用指定要过滤的列:
df = df_1.join(df_2, ["val1", "val2"], "left") # 明确指定使用df_1中的val3列 df = df.where(df_1.val3 == "b")
方案3:关联时显式选择列并别名
关联时只选择需要的列,同时给重复列加别名:
from pyspark.sql import functions as F # 对df_2的val3加别名后再关联 df = df_1.join( df_2.select("val1", "val2", "val5", F.col("val3").alias("val3_df2")), ["val1", "val2"], "left" ) # 过滤时直接使用df_1的val3列 df = df.where(F.col("val3") == "b")
原因说明
出现歧义错误是因为关联后的DataFrame存在两个完全同名的val3列,Spark无法判断你要引用哪一列,因此必须通过重命名或明确数据源的方式指定目标列。
内容的提问来源于stack exchange,提问作者Duc Vu
相关产品推荐
相关产品推荐

