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

Spark SQL中高效实现running sum:扫描量统计优化方案咨询

优化每日扫描数聚合计算的高效方案

我正在开发一套指标计算逻辑,需要从每日扫描数(DailyScan)计算TotalScan、Last5DayScan、Month2DayScan这三个指标。目前的实现是每日对全量数据集执行sum(dailyscan)计算,但随着数据量增长,计算压力陡增。我考虑改用滚动求和(running sum)方案,但不清楚如何基于历史TotalScan值计算当日TotalScan——即当日TotalScan = 历史TotalScan值 + 当日扫描数(历史值可追溯至1-2个月前)。

示例源数据

ProcessName          DailyScan    Date
 NewInsurance         8000        04/12/2024
 InsuranceRenewal     4500        04/12/2024
 Fraud Detection      28          04/12/2024
 Policy Withdrawn     100         04/01/2024
 NewInsurance         2100        04/13/2024
 New Insurance        400         04/14/2024
InsuranceRenewal     500         04/14/2024
InsuranceRenewal     500         04/18/2024
New Insurance        500         04/18/2024

期望输出(04/18/2024执行查询)

ProcessName   TotalScan Last5DayScan  Month2DayScan   DailyScan     Date 
NewInsurance    8000     8000          8000              8000        04/12/2024
NewInsurance    10100    10100         10100             2100        04/13/2024
NewInsurance    10500    10500         10500             400         04/14/2024
NewInsurance    11000     900          11000             500         04/18/2024

我当前的实现是将源表与日历表关联后,按ProcessName和CalendarDate分组,每日全量求和得到TotalScan。虽然能得到正确结果,但效率太低,希望得到更高效的实现思路。


高效实现思路

1. 窗口函数滚动求和(推荐)

利用SQL窗口函数直接计算各指标,避免全量扫描:

  • TotalScan:按ProcessName分组,日期升序计算累计求和
    SUM(DailyScan) OVER (
      PARTITION BY ProcessName 
      ORDER BY Date 
      ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
    ) AS TotalScan
    
  • Last5DayScan:计算当前日期往前5天内的扫描数总和(适配日期不连续场景)
    SUM(DailyScan) OVER (
      PARTITION BY ProcessName 
      ORDER BY Date 
      RANGE BETWEEN INTERVAL '5' DAY PRECEDING AND CURRENT ROW
    ) AS Last5DayScan
    
  • Month2DayScan:计算当月从月初到当前日期的累计扫描数
    SUM(DailyScan) OVER (
      PARTITION BY ProcessName, DATE_TRUNC('month', Date) 
      ORDER BY Date 
      ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
    ) AS Month2DayScan
    
    注:不同SQL方言日期函数语法有差异,比如MySQL用DATE_FORMAT(Date, '%Y-%m'),SQL Server用DATEFROMPARTS(YEAR(Date), MONTH(Date), 1)。

2. 预计算快照表(超大数据量场景)

如果窗口函数仍有性能压力,可维护每日快照表,仅计算增量:

  • 每日基于前一日快照+当日新增数据更新指标,避免全量计算:
    INSERT INTO process_scan_snapshot 
    (ProcessName, Date, TotalScan, Last5DayScan, Month2DayScan, DailyScan)
    SELECT 
      COALESCE(d.ProcessName, s.ProcessName),
      CURRENT_DATE,
      COALESCE(s.TotalScan, 0) + COALESCE(d.DailyScan, 0) AS TotalScan,
      (SELECT SUM(DailyScan) 
       FROM process_scan_snapshot 
       WHERE ProcessName = COALESCE(d.ProcessName, s.ProcessName) 
         AND Date >= CURRENT_DATE - INTERVAL '4' DAY) + COALESCE(d.DailyScan, 0) AS Last5DayScan,
      CASE 
        WHEN DATE_TRUNC('month', CURRENT_DATE) = CURRENT_DATE 
        THEN COALESCE(d.DailyScan, 0) 
        ELSE COALESCE(s.Month2DayScan, 0) + COALESCE(d.DailyScan, 0) 
      END AS Month2DayScan,
      COALESCE(d.DailyScan, 0) AS DailyScan
    FROM (
      SELECT ProcessName, SUM(DailyScan) AS DailyScan 
      FROM daily_scan WHERE Date = CURRENT_DATE GROUP BY ProcessName
    ) d
    FULL JOIN (
      SELECT ProcessName, TotalScan, Month2DayScan 
      FROM process_scan_snapshot WHERE Date = CURRENT_DATE - INTERVAL '1' DAY
    ) s ON d.ProcessName = s.ProcessName;
    
  • 查询时直接读取快照表,无需实时计算,性能大幅提升。

3. 源数据预处理

先清洗源数据,合并同一ProcessName、同一Date的重复记录(示例中NewInsurance和New Insurance疑似输入错误),减少计算量:

SELECT 
  TRIM(REPLACE(ProcessName, ' ', '')) AS ProcessName,
  SUM(DailyScan) AS DailyScan,
  Date
FROM daily_scan
GROUP BY TRIM(REPLACE(ProcessName, ' ', '')), Date;

内容的提问来源于stack exchange,提问作者mehtat_90

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 03:13:11