Spark Scala按cid分组累加交易金额生成sumAmt新列实现问询
Spark Scala 按客户分组累加历史交易金额实现方案
需求说明
需要为DataFrame新增sumAmt列,按cid分组累加对应客户的历史交易金额,同日期的交易按顺序累计。
初始数据样例
cid transAmt trasnDate 1 10 2-Aug 1 20 3-Aug 1 30 3-Aug 2 40 2-Aug 2 50 3-Aug 3 60 4-Aug
预期输出样例
cid transAmt trasnDate sumAmt 1 10 2-Aug 10 1 20 3-Aug 30 1 30 3-Aug 60 2 40 2-Aug 40 2 50 3-Aug 90 3 60 4-Aug 60
实现代码
依赖导入
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions.sum // 若需要处理日期格式转换可额外导入to_date函数 // import org.apache.spark.sql.functions.to_date
核心逻辑
// 定义窗口规则:按cid分组,按交易日期排序,累加范围为分组内首行到当前行 val customerTransWindow = Window.partitionBy("cid") .orderBy("trasnDate") .rowsBetween(Window.unboundedPreceding, Window.currentRow) // 新增累计金额列 val resultDF = initialDF.withColumn("sumAmt", sum("transAmt").over(customerTransWindow)) // 输出验证结果 resultDF.show()
注意事项
如果实际场景中日期为字符串格式且排序不符合预期,建议先使用to_date函数将日期列转为日期类型后再参与窗口排序。
内容的提问来源于stack exchange,提问作者HEMANT PATEL
相关产品推荐
相关产品推荐

