Spark 2.3中连接含同名列的JDBC表报错问题咨询
解决Spark 2.3中JDBC Join查询的重复列错误
这个问题我之前也碰到过,本质是Spark 2.3对JDBC数据源的Schema处理逻辑做了调整,和早期版本的行为不一样:
- 在旧版Spark里,你直接把Join SQL丢给
spark.read.jdbc()时,Spark完全是“甩手掌柜”——它会把这条SQL原封不动发给数据库执行,只取返回结果的列(也就是你指定的name和body),根本不会去解析t1、t2这两张原表的Schema。 - 但Spark 2.3开始,JDBC模块新增了一个逻辑:它会先去解析你SQL里提到的每张表的Schema(t1有
id,body,t2有id,name),发现两张表都有id列,直接就抛出了AnalysisException,哪怕你的最终查询根本没把id列捞回来。
两种靠谱的解决办法
办法1:用子查询包装你的Join语句
把原来的Join SQL套进一个子查询里,这样Spark只会解析子查询返回的结果Schema,不会再去碰原表的Schema:
// 注意这里的table参数是一个子查询,要给它起个别名 spark.read.jdbc( url = "你的JDBC连接地址", table = "(select name, body from t1 inner join t2 on t1.id = t2.id) as temp_join", connectionProperties = 你的连接配置 )
数据库会先执行这个子查询,返回只有name和body的结果集,Spark拿到后就不会有重复列的问题了。
办法2:用DataFrame API分步实现Join
如果不想在JDBC里写复杂SQL,也可以分别读取两张表,然后在Spark里做Join,最后选需要的列:
// 先分别读取两张表 val df_t1 = spark.read.jdbc("你的JDBC地址", "t1", 连接配置) val df_t2 = spark.read.jdbc("你的JDBC地址", "t2", 连接配置) // 执行Join并选择指定列 val result_df = df_t1.join(df_t2, df_t1("id") === df_t2("id"), "inner") .select(df_t2("name"), df_t1("body"))
这种方式虽然会分两次读表,但Spark会帮你处理列名冲突,你可以明确指定要取的列,完全避开重复列的问题。
额外提醒
这个问题算是Spark 2.3优化Schema推断逻辑带来的小“坑”,本意是提前发现潜在的列冲突,但在直接写Join SQL的场景下就误判了。两种方法都能解决问题,推荐根据实际场景选:如果Join逻辑复杂(比如有过滤、聚合),用子查询让数据库来处理效率更高;如果逻辑简单,用DataFrame API更灵活易维护。
内容的提问来源于stack exchange,提问作者Fletcher Stump Smith
相关产品推荐
相关产品推荐

