{T}

时间维度聚合计算:如何在长时间窗口上实时计算聚合值?

0. 引言

流数据操作中的 Reduce 聚合(第 9 课)采用滑动窗口内全量计算,依赖窗口内保存全部原始数据。但"时间维度聚合值计算"有两个严格限制:实时返回(每条数据毫秒级)与长时间窗口 + 大数据量。这两个限制使全量窗口方案失效,必须换一种思路——本文讲解"寄存器"法及其 Redis 实现。

1. 难点:为什么全量窗口方案不可行

按时间维度聚合是最常见的计算问题:计数(count)、求和(sum)、均值(avg)、方差(variance)、最小(min)、最大(max)。以风控场景为例:

sql
# 过去一周在相同设备上交易次数
SELECT COUNT(*) FROM stream
WHERE event_type = "transaction"
AND timestamp >= 1530547200000 and timestamp < 1531152000000
GROUP BY device_id;

关系型数据库未建索引时会遍历全表过滤再分组聚合。若流计算也把每条消息保存到缓冲区、窗口结束时遍历全量计算,当窗口较长、数据量较大时:内存占用巨大、计算资源消耗高、耗时不可控——无法满足实时性。

2. 寄存器法:只保留必要的聚合信息

核心优化思想:降低计算复杂度,只保留必要的聚合信息,不保存所有原始数据。幸运的是,各类聚合运算都能用少数几个"寄存器"记录中间结果:

聚合所需寄存器寄存器含义
count1 个记录数
sum1 个总和
avg2 个总和 + 记录数
min / max各 1 个最小值 / 最大值
variance3 个总和、记录数、平方和

更复杂的聚合(偏度、峰度)通过数学公式转化,同样能找到对应寄存器。

2.1 原理:以小窗口累加替代全量扫描

以"过去一周在相同设备上交易次数"为例,将 7 天窗口划分为 7 个 1 天小窗口:

图表渲染中…
  • 写入:每个(设备, 窗口)分配一个 count 寄存器,事件到来时寄存器 +1;
  • 查询:读取连续 7 个窗口的寄存器值累加,即得"过去一周"结果。

3. 实现:Redis INCR 计数

计数类聚合非常适合 Redis 的 INCR 指令(原子加一)。key 设计为:

text
$event_type.$device_id.$window_unit.$window_index
python
# 更新:只更新当前小窗口
$event_type = transaction
$device_id = d000001
$window_unit = 86400000            # 1 天(毫秒)
$window_index = 1532496076032 // $window_unit   # 时间戳除以窗口单元 → 窗口索引
$key = f"{$event_type}.{$device_id}.{$window_unit}.{$window_index}"
redis.incr($key)

# 查询:汇总过去 7 个子窗口
sum = 0
for i in range(0, 7):
    k = f"{$event_type}.{$device_id}.{$window_unit}.{$window_index - i}"
    sum += redis.get(k)
return sum

更新与查询分离:写入只需更新一个小窗口计数,查询时汇总 7 个 key 即可。

4. 寄存器方案的不足:状态存储

寄存器方案大幅减少内存与计算量,但引入了状态存储问题:

  • 分组变量的通常远高于预期——"统计每个 IP 访问次数"有四十多亿 IP,"统计用户登录次数"有十四亿用户;复合变量(用户 × IP C 段)更是笛卡尔积量级;
  • 因此状态不能全部放本地内存,需保存到外部存储(Redis、Apache Ignite 或本地磁盘);
  • 同时必须为状态设置 TTL 过期时间,清理过期状态、避免空间无限增长。

这正是流计算框架(Flink、Spark Streaming)专门引入状态存储功能并提供对应 API 的根本原因——状态问题是后续所有算法的共性基础。

5. 小结

  • 两个前提(实时返回 + 长周期大窗口)共同决定了必须采用寄存器法优化时间维度聚合;
  • 寄存器法以"小窗口累加"替代"全量扫描",count/sum/avg/min/max 只需 1-3 个寄存器;
  • Redis INCR + 分窗口 key 是计数聚合的经典工程实现;
  • 寄存器的引入使流计算成为有状态系统,也催生了框架级的状态存储与 TTL 机制。

下一章讲解关联图谱分析:如何用 Lambda 架构实现实时的社交网络分析。