PySpark连接不同列名DataFrame报错,求正确连接方法
PySpark DataFrame连接报错原因及解决方法
错误原因
你的代码存在两个核心问题:
- 提前筛选列导致连接条件丢失:你在
join的第一个参数中使用regions.select("region_id"),这会生成仅包含region_id列的新DataFrame,原regions中的name列被直接丢弃,导致连接条件city_df.region_name == regions_df.name找不到name列,触发报错。 - 变量名不一致:代码中第一个参数是
regions.select(...),但连接条件里用的是regions_df.name,如果你的地区数据集实际变量名是regions而非regions_df,这也会导致列引用错误。
正确写法
场景1:仅需要regions中的region_id列
先保留连接所需的name列,完成连接后再按需筛选列:
# 保留name列用于连接,同时选择需要的region_id city_df.join(regions.select("name", "region_id"), on=city_df.region_name == regions.name, how='inner') # 用别名提升代码可读性(推荐) city_df.alias("city").join(regions.alias("region").select("name", "region_id"), on=city.region_name == region.name, how='inner')
场景2:需要regions中的多列
直接使用原DataFrame连接,无需提前筛选:
city_df.join(regions_df, on=city_df.region_name == regions_df.name, how='inner') # 别名写法 city_df.alias("c").join(regions_df.alias("r"), on=c.region_name == r.name, how='inner')
错误信息解读
报错信息Resolved attribute(s) name#32 missing from ... region_id#31L明确说明:当前参与连接的右侧DataFrame只有region_id列,找不到连接条件需要的name列,这就是因为你提前用select("region_id")丢弃了该列。
内容的提问来源于stack exchange,提问作者user453575457
相关产品推荐
相关产品推荐

