Spark 2.3.0中withColumn调用callUDF触发AnalysisException求助
解决Spark 2.3.0升级后UDF调用的AnalysisException问题
这个问题是Spark 2.3.0对SQL分析器的属性解析逻辑做了严格优化导致的——虽然你的svWithCoordsTzAndDistancesDF里确实存在latitude、longitude、mean_lat、mean_lon这些列,但在执行计划解析阶段,Spark无法正确关联UDF引用的列到join后的DataFrame对应的物理列上(因为这些列分别来自childDF1(聚合结果)和childDF2(过滤后的原表),元数据来源不同)。
下面是几种可行的解决方案:
方案1:给Join的表添加别名,明确指定列来源
在join操作时给两个子DataFrame设置别名,然后在调用UDF时通过别名指定列的归属,让Spark分析器能精准定位列:
// 给join的两个DataFrame设置别名 val childDF1WithAlias = childDF1.as("agg_df") val childDF2WithAlias = childDF2.as("raw_df") // 基于hpan字段join,使用别名区分列来源 val svWithCoordsTzAndDistancesDF = childDF2WithAlias.join( childDF1WithAlias, childDF2WithAlias("hpan") === childDF1WithAlias("hpan"), "inner" // 根据你的业务需求选择join类型 ) // 调用UDF时明确指定列来自哪个别名表 val finalDF = svWithCoordsTzAndDistancesDF.withColumn( "distance", callUDF("calcDistance", col("raw_df.latitude"), col("raw_df.longitude"), col("agg_df.mean_lat"), col("agg_df.mean_lon") ) )
方案2:重命名Join后的列,消除元数据歧义
如果不想用表别名,可以在join后给相关列重命名,让每个列的引用更明确:
val joinedDF = childDF2.join(childDF1, "hpan") // 重命名原表的经纬度列 .withColumnRenamed("latitude", "user_latitude") .withColumnRenamed("longitude", "user_longitude") // 重命名聚合结果的均值列 .withColumnRenamed("mean_lat", "avg_latitude") .withColumnRenamed("mean_lon", "avg_longitude") // 使用新列名调用UDF val finalDF = joinedDF.withColumn( "distance", callUDF("calcDistance", col("user_latitude"), col("user_longitude"), col("avg_latitude"), col("avg_longitude") ) )
方案3:使用expr()替代callUDF()调用函数
有时候直接用SQL表达式的方式调用UDF,能绕过Spark分析器的列解析问题:
val finalDF = svWithCoordsTzAndDistancesDF.withColumn( "distance", expr("calcDistance(latitude, longitude, mean_lat, mean_lon)") )
为什么Spark 2.2.0没问题?
Spark 2.2.0的分析器对列的解析逻辑更宽松,即使列来自不同的父节点,只要列名唯一就能匹配到。但2.3.0为了优化执行计划的准确性,加强了对列元数据来源的校验,导致这种跨节点的列引用出现解析错误。
内容的提问来源于stack exchange,提问作者Danila Zharenkov
相关产品推荐
相关产品推荐

