Scala Spark子查询中使用聚合函数计算日期差并新增列的实现方法
实现方案
错误点说明
你原有代码存在几处问题:
- 语法错误:join条件的写法不规范,Spark 中列相等判断应使用
===,且列引用的引号不匹配,正确写法为$"a.name" === $"b.name" - 逻辑错误:不能在
withColumn中直接调用a.agg(),agg会返回聚合后的完整DataFrame,无法直接作为列值参与每行计算 - 拼写错误:你定义的表中日期列名为
openday,代码中误写为opday
实现方法
根据max(openday的来源不同,分两种常用场景实现:
场景1:计算a表自身按name分组的最大openday与当前行openday的差值
该场景不需要关联b表,直接用窗口函数实现即可:
// 导入依赖包 import org.apache.spark.sql.functions._ import org.apache.spark.sql.expressions.Window // 定义按name分区的窗口 val nameWindow = Window.partitionBy("name") // 计算差值生成新列 val result = a_table .withColumn("max_openday", max("openday").over(nameWindow)) .withColumn("difference", datediff($"max_openday", $"openday")) result.show()
场景2:取b表按name分组的最大openday计算差值
如果你的max(openday)取自b表,先对b表做预聚合再关联,避免join后数据膨胀:
// 先聚合b表得到每个name对应的最大openday val b_agg = b_table.groupBy("name").agg(max("openday").alias("b_max_openday")) // 左关联a表后计算差值 val result = a_table .join(b_agg, Seq("name"), "left") .withColumn("difference", datediff($"b_max_openday", $"openday")) result.show()
如果需要添加判断条件,直接在withColumn中嵌套when函数即可,示例:withColumn("difference", when(condition, datediff($"max_openday", $"openday")).otherwise(null))
内容的提问来源于stack exchange,提问作者Рамазан Кагерманов
相关产品推荐
相关产品推荐

