SQL转换为Spark DataFrame API时遇===报错及关联条件写法问题
错误原因
select中直接写"PENDING"是尝试读取DF中名为PENDING的字段,不是生成常量字符串,需要用lit函数传入常量值。- 你在
filter中直接将ORDER_DATE列和SHOPPING_DF.groupBy("CART_ID").agg(max("ORDER_DATE"))返回的DataFrame做等值判断,二者数据类型完全不匹配,所以触发了===方法的运行时异常。 - 原SQL中的关联子查询是按
cart_id分组取每个分组的最大order_date,再匹配对应行,你之前的写法没有处理分组后的cart_id关联逻辑。
正确实现方案
方案1:开窗函数实现(性能最优,避免二次扫描表)
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ // 定义按CART_ID分组的开窗规则 val winSpec = Window.partitionBy("CART_ID") val result = SHOPPING_DF.as("sp") // 新增列存储当前CART_ID对应的最大订单日期 .withColumn("MAX_ORDER_DATE", max(col("ORDER_DATE")).over(winSpec)) // 过滤出等于最大订单日期的行 .filter(col("ORDER_DATE") === col("MAX_ORDER_DATE")) // 和订单表内关联 .join(ORDER_DF.as("o"), Seq("CART_ID"), "inner") // 选取指定字段,常量用lit函数生成 .select( col("o.order_id"), lit("PENDING").alias("order_status") ) // 去重 .distinct()
方案2:完全对齐原SQL关联子查询逻辑
import org.apache.spark.sql.functions._ // 先预计算每个CART_ID的最大订单日期 val maxDateDf = SHOPPING_DF.groupBy("CART_ID").agg(max("ORDER_DATE").alias("MAX_ORDER_DATE")) val result = SHOPPING_DF // 关联拿到每个CART_ID的最大订单日期 .join(maxDateDf, Seq("CART_ID"), "inner") .filter(col("ORDER_DATE") === col("MAX_ORDER_DATE")) // 关联订单表 .join(ORDER_DF, Seq("CART_ID"), "inner") .select( col("order_id"), lit("PENDING").alias("order_status") ) .distinct()
内容的提问来源于stack exchange,提问作者Punter Vicky
相关产品推荐
相关产品推荐

