You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于BigQuery/PySpark实现带上下限的累计求和逻辑(千万级数据)

带上下限的累计求和实现(大数据量适配)

需求背景

现有一张含1000万唯一Product_ID的表,hr_status字段对应三种状态:good(-1)、neutral(0)、bad(1)。需实现带上下限的累计求和逻辑:

  • 上限规则:累计和达到5后不再累加,保持5;
  • 下限规则:累计和最低不低于0,低于时保持0;
    通过该累计和判断产品是否需要抽检。

此前Python嵌套循环的POC无法适配大数据量,row_number()分区方式尝试失败,需用BigQuery或PySpark实现。

示例数据

Product_IDtimehr_status
AA2024-06-10 01:10:10-1
AA2024-06-10 02:10:100
AA2024-06-10 03:10:101
AA2024-06-10 04:10:101
AA2024-06-10 05:10:101
AA2024-06-10 06:10:101
AA2024-06-10 07:10:101
AA2024-06-10 08:10:101
AA2024-06-10 09:10:101

预期结果(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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 01:44:59