基于BigQuery/PySpark实现带上下限的累计求和逻辑(千万级数据)
带上下限的累计求和实现(大数据量适配)
需求背景
现有一张含1000万唯一Product_ID的表,hr_status字段对应三种状态:good(-1)、neutral(0)、bad(1)。需实现带上下限的累计求和逻辑:
- 上限规则:累计和达到5后不再累加,保持5;
- 下限规则:累计和最低不低于0,低于时保持0;
通过该累计和判断产品是否需要抽检。
此前Python嵌套循环的POC无法适配大数据量,row_number()分区方式尝试失败,需用BigQuery或PySpark实现。
示例数据
| Product_ID | time | hr_status |
|---|---|---|
| AA | 2024-06-10 01:10:10 | -1 |
| AA | 2024-06-10 02:10:10 | 0 |
| AA | 2024-06-10 03:10:10 | 1 |
| AA | 2024-06-10 04:10:10 | 1 |
| AA | 2024-06-10 05:10:10 | 1 |
| AA | 2024-06-10 06:10:10 | 1 |
| AA | 2024-06-10 07:10:10 | 1 |
| AA | 2024-06-10 08:10:10 | 1 |
| AA | 2024-06-10 09:10:10 | 1 |
预期结果(CUM_SUM初始值为0)
| CUM_SUM |
|---|
| 0 |
| 0 |
| 1 |
| 2 |
| 3 |
| 4 |
| 5 |
| 5 |
| 5 |
解决方案
1. BigQuery实现(递归CTE)
利用BigQuery的递归CTE按Product_ID分区、时间排序后逐行计算带上下限的累计值,适配大数据量场景。
WITH ranked_data AS ( SELECT Product_ID, time, hr_status, ROW_NUMBER() OVER(PARTITION BY Product_ID ORDER BY time) AS rn FROM your_table_name ), recursive_cum_sum AS ( -- 初始化每个Product_ID的第一行数据 SELECT Product_ID, time, hr_status, rn, GREATEST(0, LEAST(5, 0 + hr_status)) AS CUM_SUM FROM ranked_data WHERE rn = 1 UNION ALL -- 递归计算后续行的累计值 SELECT rd.Product_ID, rd.time, rd.hr_status, rd.rn, GREATEST(0, LEAST(5, rcs.CUM_SUM + rd.hr_status)) AS CUM_SUM FROM ranked_data rd JOIN recursive_cum_sum rcs ON rd.Product_ID = rcs.Product_ID AND rd.rn = rcs.rn + 1 ) SELECT Product_ID, time, hr_status, CUM_SUM FROM recursive_cum_sum ORDER BY Product_ID, time;
说明
- 先通过
ranked_data给每个Product_ID下的行按时间排序并编号; - 递归CTE初始部分处理每个产品的第一行,基于初始值0计算首次累计和并做上下限截断;
- 递归部分逐行关联上一行结果,计算当前行累计值后同样应用上下限规则;
- BigQuery会自动优化递归CTE的执行计划,避免嵌套循环的性能瓶颈。
2. PySpark实现(分布式分组计算)
利用Spark的分布式计算能力,按Product_ID分组后逐行计算带上下限的累计和,适配千万级数据量。
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, TimestampType, IntegerType # 初始化SparkSession spark = SparkSession.builder.appName("CumSumWithBounds").getOrCreate() # 定义表结构 schema = StructType([ StructField("Product_ID", StringType(), True), StructField("time", TimestampType(), True), StructField("hr_status", IntegerType(), True) ]) # 读取数据(替换为你的实际数据源) df = spark.read.schema(schema).table("your_table_name") # 按Product_ID分组后计算带上下限的累计和 def calculate_cum_sum(rows): cum_sum = 0 # 按时间排序当前分组内的行 for row in sorted(rows, key=lambda x: x.time): new_sum = cum_sum + row.hr_status # 应用上下限规则 cum_sum = max(0, min(5, new_sum)) yield (row.Product_ID, row.time, row.hr_status, cum_sum) # 转换为RDD处理后转回DataFrame result_rdd = df.rdd.groupBy(lambda x: x.Product_ID).flatMap(lambda x: calculate_cum_sum(x[1])) result_df = result_rdd.toDF(["Product_ID", "time", "hr_status", "CUM_SUM"]) # 展示或保存结果 result_df.orderBy("Product_ID", "time").show()
说明
- 将DataFrame转为RDD后按
Product_ID分组,利用Spark分布式计算能力拆分计算压力; - 每个分组内按时间排序,逐行计算累计和并应用上下限规则;
- 该实现直观易懂,也可替换为自定义UDAF(用户自定义聚合函数)进一步优化性能。
内容的提问来源于stack exchange,提问作者Victor Franco
相关产品推荐
相关产品推荐

