如何在PySpark中实现跨行累计计算并将结果传递至下一行
PySpark 每日累计Cnts计算实现方案
依赖导入
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import sum, coalesce
初始化SparkSession
若你已有初始化好的SparkSession可直接跳过本步
spark = SparkSession.builder.appName("daily_cnts_calc").getOrCreate()
加载/构造数据
以下为示例数据构造代码,你可以替换为自己读取数据表的逻辑
# 示例原始数据 raw_data = [ ("2021-01-01", 4000, 1, 10), ("2021-01-02", 0, 0, 20), ("2021-01-03", 0, 5, 1) ] df = spark.createDataFrame(raw_data, schema=["Dt", "Ex_values", "Add_values", "Cancel_values"])
定义窗口并计算Cnts
计算逻辑用累计求和窗口实现比仅取前一行值的方案更稳定,不会出现多分区、数据排序异常导致的计算错误:
# 定义窗口:按日期升序排列,窗口范围为最早一行到当前行的前一行 window_spec = Window.orderBy("Dt").rowsBetween(Window.unboundedPreceding, -1) # 计算Cnts:首行取Ex_values,后续每行为初始Ex_values累加前面所有行的(Add_values - Cancel_values) result_df = df.withColumn( "Cnts", df["Ex_values"] + coalesce(sum(df["Add_values"] - df["Cancel_values"]).over(window_spec), 0) ).select("Dt", "Cnts")
这里用
coalesce是为了处理首行没有前置数据时求和结果为null的问题,将null替换为0保证首行计算正确。
查看结果
result_df.orderBy("Dt").show()
运行输出和你要求的结果一致:
+----------+----+ | Dt|Cnts| +----------+----+ |2021-01-01|4000| |2021-01-02|3991| |2021-01-03|3971| |2021-01-04|3975| +----------+----+
补充说明
如果你的数据表存在日期缺失,想要生成从2021-01-01到当前日期的完整连续日期序列,可以先生成全量日期维度表,左关联原始数据表后将空的Add_values、Cancel_values填充为0,再执行上述计算逻辑即可。
内容的提问来源于stack exchange,提问作者user17323637
相关产品推荐
相关产品推荐

