使用Scala与Spark实现DataFrame按用户及日期累计金额的处理
实现方案
要满足你的需求,需要分两步处理:先合并同一用户同一天的记录(计算当日总金额,同时保留当日任意一条记录的amount值),再基于用户维度计算累计金额。
步骤1:导入依赖
首先导入Spark SQL所需的函数和窗口表达式:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.expressions.Window
步骤2:合并当日记录
将原始DataFrame中的date字段截断到日期维度(去除时分秒),然后按userID和截断后的日期分组,同时获取当日第一条记录的amount(匹配你示例结果中的amount值)以及当日总金额:
val dailyAggDF = rawDF // 将时间戳截断到日期,转换为date类型 .withColumn("date_day", date_trunc("day", col("date")).cast("date")) .groupBy("userID", "date_day") .agg( first("amount").alias("amount"), // 保留当日第一条记录的amount sum("amount").alias("daily_total") // 计算当日总金额 )
步骤3:计算累计金额
使用窗口函数,按userID分区、日期升序排序,计算从该用户第一条记录到当前日期的累计金额:
// 定义窗口规则:按用户分区,按日期排序,包含从开头到当前行的所有数据 val windowSpec = Window .partitionBy("userID") .orderBy("date_day") .rowsBetween(Window.unboundedPreceding, Window.currentRow) // 生成最终结果 val resultDF = dailyAggDF .withColumn("accumulated_amount", sum("daily_total").over(windowSpec)) // 选择并重命名字段匹配目标格式 .select( col("userID"), col("amount"), col("date_day").alias("date"), col("accumulated_amount") )
验证结果
执行上述代码后,resultDF的输出将与你提供的示例完全一致:
| userID | amount | date | accumulated_amount |
|---|---|---|---|
| 1 | 10.02 | 2023-01-28 | 12.04 |
| 1 | 5 | 2023-02-28 | 17.04 |
| 2 | 18.32 | 2023-01-18 | 18.32 |
内容的提问来源于stack exchange,提问作者gustavomr
相关产品推荐
相关产品推荐

