Spark 2.3/2.4与3.2同查询结果不一致的配置解决方案咨询
Spark 3.2 兼容Spark 2.x Join去重列配置方案
问题重现
以下代码在Spark 2.3/2.4和Spark 3.2中执行结果不一致:
spark-shell --master yarn --deploy-mode client val employeeDF = spark.createDataFrame(Seq((1,"Amit","2022","DAP"),(2,"Simran","2022","DAP"),(3,"Purna","2019","DAP"))).toDF("emp_id","name","year_joined","dept_name") val deptDF = spark.createDataFrame(Seq(("DAP",1),("DAP",2),("DAP",3))).toDF("dept_name", "emp") val joinedDF = employeeDF.join(deptDF, Seq("dept_name")).as("df1").select($"df1.*") joinedDF.printSchema
版本差异表现
- Spark 2.3/2.4:执行后
joinedDF仅保留唯一的dept_name列,无重复列 - Spark 3.2:执行后
joinedDF包含两个重复的dept_name列,后续分组操作会触发Ambiguous column name错误
执行计划对比
Spark 2.4 执行计划
*(2) Project [dept_name#11, emp_id#8, name#9, year_joined#10, emp#21] +- *(2) BroadcastHashJoin [dept_name#11], [dept_name#20], Inner, BuildRight :- *(2) Project [_1#0 AS emp_id#8, _2#1 AS name#9, _3#2 AS year_joined#10, _4#3 AS dept_name#11] : +- *(2) Filter isnotnull(_4#3) : +- LocalTableScan [_1#0, _2#1, _3#2, _4#3] +- BroadcastExchange HashedRelationBroadcastMode(List(input[0, string, true])) +- *(1) Project [_1#16 AS dept_name#20, _2#17 AS emp#21] +- *(1) Filter isnotnull(_1#16) +- LocalTableScan [_1#16, _2#17]
Spark 3.2 执行计划
== Physical Plan == AdaptiveSparkPlan isFinalPlan=false +- Project [dept_name#20, dept_name#11, emp_id#8, name#9, year_joined#10, emp#21] +- BroadcastHashJoin [dept_name#11], [dept_name#20], Inner, BuildRight, false :- LocalTableScan [emp_id#8, name#9, year_joined#10, dept_name#11] +- BroadcastExchange HashedRelationBroadcastMode(List(input[0, string, true]),false), [id=#16] +- LocalTableScan [dept_name#20, emp#21]
已尝试无效配置
- 关闭自适应执行计划(AQE)
- 将递归CTE模式设置为LEGACY
以上配置均无法消除重复列。
可行配置方案
Spark 3.x默认关闭了Join后自动去重连接键的行为,可通过以下配置恢复Spark 2.x的旧行为:
配置参数说明
启用spark.sql.legacy.join.removeDuplicateKeys参数,该参数控制Join操作后是否自动移除重复的连接键列,Spark 2.x默认值为true,Spark 3.x默认值为false。
配置方式
- 启动spark-shell时指定:
spark-shell --master yarn --deploy-mode client --conf spark.sql.legacy.join.removeDuplicateKeys=true
- 代码中动态设置:
spark.conf.set("spark.sql.legacy.join.removeDuplicateKeys", "true")
- 全局配置(spark-defaults.conf):
spark.sql.legacy.join.removeDuplicateKeys true
验证效果
配置生效后,重新执行原Join代码,joinedDF.printSchema将仅输出一个dept_name列,执行计划的Project阶段也会和Spark 2.x一致,仅保留单一连接键列,后续分组操作不再报错。
内容的提问来源于stack exchange,提问作者ASR
相关产品推荐
相关产品推荐

