PySpark中聚合数据的滚动累积乘积计算方案
PySpark 按ID分组并对连续相同值子组计算累积乘积
要实现需求中按id分组,对每个id下连续相同values的子组计算累积乘积,可通过以下步骤完成:
步骤说明
- 标记连续相同值的子组:在每个
id内,识别连续相同的values并分配唯一组号,确保后续累积计算按子组顺序进行。 - 计算子组累积乘积:在每个
id内,按子组的出现顺序计算累积乘积。 - 关联乘积结果到原数据:将累积乘积结果映射到对应子组的所有行。
代码实现
from pyspark.sql import Window from pyspark.sql import functions as F # 定义窗口:按id分区,用自增ID保证原数据顺序 window_id = Window.partitionBy("id").orderBy(F.monotonically_increasing_id()) # 标记连续相同values的子组:当前行与前一行values不同则标记为新组,累加生成组号 df_with_group = df.withColumn( "is_new_group", F.when(F.lag("values").over(window_id) != F.col("values"), 1).otherwise(0) ).withColumn( "group_id", F.sum("is_new_group").over(window_id.rangeBetween(Window.unboundedPreceding, Window.currentRow)) ) # 计算每个id内按组号顺序的累积乘积 window_product = Window.partitionBy("id").orderBy("group_id").rowsBetween(Window.unboundedPreceding, Window.currentRow) df_result = df_with_group.withColumn( "calculated_product", F.product("values").over(window_product) ).drop("is_new_group", "group_id") # 查看结果 df_result.show()
结果验证
运行代码后,calculated_product列将与预期的product列完全匹配:
- 例如
id=5的子组顺序为12→4→0.5,累积乘积依次为12、124=48、480.5=24,与预期一致。 - 同一子组内的所有行将共享相同的累积乘积值。
内容的提问来源于stack exchange,提问作者cnns
相关产品推荐
相关产品推荐

