Spark 2.2.0执行unionByName报错无该属性的问题咨询
报错根因
你遇到的报错本质是Spark版本不兼容导致,和Python 2.7版本无关:
- PySpark的
unionByName方法从Spark 2.3版本才正式上线,你当前使用的Spark 2.2.0版本原生没有提供这个API - 你代码中用到的
allowMissingColumns参数到Spark 3.1版本才新增,即使升级到Spark 2.3/2.4也无法使用该参数
方案1:现有Spark 2.2版本实现等价功能
场景1:两个DataFrame列名完全一致,仅列顺序不同
对齐列顺序后直接调用union即可实现原生unionByName的基础效果:
# 按第一个df的列顺序重排第二个df的列,再合并 df_new = old.union(old2.select(old.columns))
场景2:需要支持缺失列自动补null(对应allowMissingColumns=True效果)
先提取两个DataFrame所有列的并集,给两个df分别补全不存在的列并赋值为null,对齐列顺序后再合并:
from pyspark.sql.functions import lit # 取两个df的所有列的并集 all_cols = list(set(old.columns + old2.columns)) # 补全df的缺失列,按统一顺序排列列 def align_columns(df, all_cols): for col_name in all_cols: if col_name not in df.columns: df = df.withColumn(col_name, lit(None)) return df.select(sorted(all_cols)) old_aligned = align_columns(old, all_cols) old2_aligned = align_columns(old2, all_cols) # 合并对齐后的两个df df_new = old_aligned.union(old2_aligned)
方案2:升级Spark到支持对应功能的版本
- 仅需要原生支持
unionByName基础功能:升级Spark到2.3及以上版本即可 - 需要原生支持
allowMissingColumns参数:升级Spark到3.1及以上版本
注意:升级Spark版本需要适配你当前使用的Cloudera发行版支持的组件版本,避免和现有集群服务出现兼容性冲突;Spark 3.x系列推荐搭配Python 3.6+版本使用。
内容的提问来源于stack exchange,提问作者Niels
相关产品推荐
相关产品推荐

