Spark Scala Join后withColumn报类型不匹配(需Boolean得Column)问题
错误原因
报错found : org.apache.spark.sql.Column required: Boolean的核心问题是相等判断方法使用错误:
- 代码中
col("input_file_crt_dt").eq(col("output_file_crt_dt"))调用的是Scala原生的eq方法,该方法用于判断两个对象的引用是否相等,返回值是普通ScalaBoolean类型,只能在Driver端对已知对象做判断,无法接受SparkColumn类型作为分布式计算的表达式参数。 - Spark SQL中判断两个列的值相等,需要使用
===操作符或者Column.equalTo()方法,这两个API返回的才是Spark可执行的Column类型布尔表达式,可正常传入when函数作为判断条件。
修正后的完整代码
import org.apache.spark.sql.functions._ val demandDF1 = Seq(("DAO","2022-06-30","1"), ("DAO","2022-06-30","1"), ("CCC","2022-06-30","1"), ("APCC","2022-06-30","1"), ("ODM","2022-06-29","3"), ("EMF","2022-06-30","1"), ("T2Region","2022-06-29","4"), ("BCC","2022-06-30","1"), ("EMF","2022-07-01","1")).toDF("rgn_nm","file_crt_dt","file_vrsn") .withColumn("file_crt_dt", col("file_crt_dt").cast("date")) .withColumn("file_vrsn", col("file_vrsn").cast("int")) val outputDistinctDF = Seq(("DAO","2022-06-30","1"), ("CCC","2022-06-29","1"), ("APCC","2022-06-30","1"), ("ODM","2022-06-29","2"), ("EMF","2022-06-30","1"), ("BCC","2022-06-30","1")).toDF("region","file_crt_dt","file_vrsn") .withColumn("file_crt_dt", col("file_crt_dt").cast("date")) .withColumn("file_vrsn", col("file_vrsn").cast("int")) val inputDistinctDF = demandDF1.select(col("rgn_nm"), col("file_crt_dt"), col("file_vrsn")).distinct() val resultantDF = inputDistinctDF.join(outputDistinctDF, inputDistinctDF.col("rgn_nm") === outputDistinctDF.col("region") , "left_outer") .select( inputDistinctDF.col("rgn_nm") as "input_region", inputDistinctDF.col("file_crt_dt") as "input_file_crt_dt", inputDistinctDF.col("file_vrsn") as "input_file_vrsn", outputDistinctDF.col("region") as "output_region", outputDistinctDF.col("file_crt_dt") as "output_file_crt_dt", outputDistinctDF.col("file_vrsn") as "output_file_vrsn" ) .withColumn("flag", when( col("output_region").isNull || col("input_file_crt_dt") > col("output_file_crt_dt") || (col("input_file_crt_dt") === col("output_file_crt_dt") && col("input_file_vrsn") > col("output_file_vrsn")) , lit(1)).otherwise(lit(0)) )
注:代码中额外将flag值从字符串类型改为整型,更符合标识列的使用习惯,如果需要字符串类型改回
lit("1")/lit("0")即可;同时将gt方法替换为更直观的>操作符,效果完全一致。
运行结果
执行resultantDF.show()后输出与预期完全匹配:
+------------+-----------------+---------------+-------------+------------------+----------------+----+ |input_region|input_file_crt_dt|input_file_vrsn|output_region|output_file_crt_dt|output_file_vrsn|flag| +------------+-----------------+---------------+-------------+------------------+----------------+----+ | DAO| 2022-06-30| 1| DAO| 2022-06-30| 1| 0| | CCC| 2022-06-30| 1| CCC| 2022-06-29| 1| 1| | T2Region| 2022-06-29| 4| null| null| null| 1| | ODM| 2022-06-29| 3| ODM| 2022-06-29| 2| 1| | APCC| 2022-06-30| 1| APCC| 2022-06-30| 1| 0| | EMF| 2022-07-01| 1| EMF| 2022-06-30| 1| 1| | EMF| 2022-06-30| 1| EMF| 2022-06-30| 1| 0| | BCC| 2022-06-30| 1| BCC| 2022-06-30| 1| 0| +------------+-----------------+---------------+-------------+------------------+----------------+----+
内容的提问来源于stack exchange,提问作者Svk
相关产品推荐
相关产品推荐

