如何用PySpark按ID、时间戳实现col1=1的累计求和?
解决方案
你完全可以用PySpark的Window函数实现需求,核心是只对col1等于1的行进行累加,具体步骤如下:
1. 定义窗口规则
按id分区(保证每个ID单独计算累计值),按timestamp升序排序(确保时间顺序正确),窗口范围默认覆盖从分区起始到当前行,无需额外指定。
2. 计算累计求和
使用sum函数结合when条件判断,把col1不等于1的行替换为0,这样累加时只会统计col1=1的情况,0不会改变累计值。
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import sum, when # 初始化SparkSession spark = SparkSession.builder.appName("cumulative_sum").getOrCreate() # 创建示例DataFrame data = [ (1, "2022-01-01", 0), (1, "2022-01-02", 1), (1, "2022-01-03", 1), (1, "2022-01-04", 0), (2, "2022-01-01", 1), (2, "2022-01-02", 0), (2, "2022-01-03", 1) ] df = spark.createDataFrame(data, ["id", "timestamp", "col1"]) # 定义窗口 window_spec = Window.partitionBy("id").orderBy("timestamp") # 添加累计求和列 result_df = df.withColumn( "cum_sum", sum(when(df.col1 == 1, 1).otherwise(0)).over(window_spec) ) # 显示结果 result_df.show()
输出结果
+---+----------+----+-------+ | id| timestamp|col1|cum_sum| +---+----------+----+-------+ | 1|2022-01-01| 0| 0| | 1|2022-01-02| 1| 1| | 1|2022-01-03| 1| 2| | 1|2022-01-04| 0| 2| | 2|2022-01-01| 1| 1| | 2|2022-01-02| 0| 1| | 2|2022-01-03| 1| 2| +---+----------+----+-------+
本质是通过条件转换,只让col1=1的行贡献1到累计和中,0的行贡献0,累加后自然得到你需要的结果。
内容的提问来源于stack exchange,提问作者Marco
相关产品推荐
相关产品推荐

