You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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的输出将与你提供的示例完全一致:

userIDamountdateaccumulated_amount
110.022023-01-2812.04
152023-02-2817.04
218.322023-01-1818.32

内容的提问来源于stack exchange,提问作者gustavomr

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 07:15:39