PySpark JDBC读取JOIN结果创建DataFrame出现重复列错误如何解决
可行替代方案如下
- 方案1:使用
USING关键字替代ON做关联条件
MariaDB/MySQL原生支持USING语法,专门用于两张表关联字段同名的场景,会自动合并关联字段,仅保留1份,无需手动指定列别名。修改后的查询如下:
query = """ (SELECT * FROM persons AS pers LEFT JOIN addresses AS addr USING(id)) tmp """ # 其余读取代码保持不变即可
- 方案2:拆分读取逻辑,在Spark侧完成关联
先单独读取两张表的全量数据,再利用Spark原生的关联算子完成join,Spark对关联键为同名字段的场景会自动去重保留1份:
# 读取persons表 df_pers = spark.read \ .format("jdbc") \ .option("driver", "org.mariadb.jdbc.Driver") \ .option("url", mariadb_url) \ .option("dbtable", "persons") \ .option("user", mariadb_user) \ .option("password", mariadb_password) \ .load() # 读取addresses表 df_addr = spark.read \ .format("jdbc") \ .option("driver", "org.mariadb.jdbc.Driver") \ .option("url", mariadb_url) \ .option("dbtable", "addresses") \ .option("user", mariadb_user) \ .option("password", mariadb_password) \ .load() # Spark侧关联,自动保留唯一id列 df = df_pers.join(df_addr, on="id", how="left")
- 方案3:使用MariaDB的
EXCEPT语法排除重复列(要求MariaDB版本≥10.5)
无需逐个罗列列名,仅需在查询时排除addresses表的重复id列即可:
query = """ (SELECT pers.*, addr.* EXCEPT (id) FROM persons AS pers LEFT JOIN addresses AS addr ON pers.id = addr.id) tmp """ # 其余读取代码保持不变即可
内容的提问来源于stack exchange,提问作者Minura Punchihewa
相关产品推荐
相关产品推荐

