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:计算当月从月初到当前日期的累计扫描数
注:不同SQL方言日期函数语法有差异,比如MySQL用SUM(DailyScan) OVER ( PARTITION BY ProcessName, DATE_TRUNC('month', Date) ORDER BY Date ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS Month2DayScanDATE_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
相关产品推荐
相关产品推荐

