滑动窗口VS滚动窗口
·
滑动窗口 vs 滚动窗口:原理、差异与业务选型
本文档结合本项目(Spark Structured Streaming)的实际代码,深入讲解两种窗口的底层机制和业务场景选型。
一、一句话区分
| 窗口类型 | 核心特征 |
|---|---|
| 滚动窗口(Tumbling Window) | 窗口之间不重叠,每条数据只属于一个窗口 |
| 滑动窗口(Sliding Window) | 窗口之间有重叠,每条数据可能属于多个窗口 |
二、图解对比
2.1 滚动窗口(Tumbling Window)
窗口大小 = 30 秒,滑动步长 = 30 秒(即步长 == 窗口大小)
时间轴:
0s 30s 60s 90s 120s
├─────────┼─────────┼─────────┼─────────┤
│ 窗口 1 │ 窗口 2 │ 窗口 3 │ 窗口 4 │
│ [0, 30) │ [30,60) │ [60,90) │[90,120) │
└─────────┴─────────┴─────────┴─────────┘
事件归属:
event A (t=15s) → 只属于窗口 1
event B (t=30s) → 只属于窗口 2(左闭右开)
event C (t=45s) → 只属于窗口 2
event D (t=89s) → 只属于窗口 3
每条事件恰好属于 1 个窗口,不多不少。
2.2 滑动窗口(Sliding Window)
窗口大小 = 30 秒,滑动步长 = 10 秒
时间轴:
0s 10s 20s 30s 40s 50s 60s 70s 80s
├─────┼─────┼─────┼─────┼─────┼─────┼─────┼─────┤
│ 窗口1 [0, 30) │
│ │ 窗口2 [10, 40) │
│ │ 窗口3 [20, 50) │
│ │ 窗口4 [30, 60) │
│ │ 窗口5 [40, 70) │
事件归属:
event A (t=15s) → 属于窗口 1 [0,30) ✅
→ 属于窗口 2 [10,40) ✅
→ 不属于窗口 3 [20,50),因为 15 < 20 ❌
event B (t=25s) → 属于窗口 1 [0,30) ✅
→ 属于窗口 2 [10,40) ✅
→ 属于窗口 3 [20,50) ✅
→ 三个窗口!
每条事件属于 ⌈窗口大小 / 滑动步长⌉ = ⌈30/10⌉ = 3 个窗口。
2.3 重叠可视化
滚动窗口(无重叠):
┌──────┐┌──────┐┌──────┐┌──────┐
│ W1 ││ W2 ││ W3 ││ W4 │
└──────┘└──────┘└──────┘└──────┘
滑动窗口(有重叠):
┌──────────────┐
│ W1 │
└────┬─────────┘
┌──────────────┐
│ W2 │
└────┬─────────┘
┌──────────────┐
│ W3 │
└────┬─────────┘
┌──────────────┐
│ W4 │
└──────────────┘
三、Spark 代码对比
3.1 滚动窗口(本项目当前实现)
// BehaviorStreamProcessor.java 中的实现
parsed.withWatermark("eventTime", "1 minute")
.groupBy(
window(col("eventTime"), "30 seconds"), // 只传一个参数 = 滚动窗口
col("userId"),
col("behaviorType"))
.agg(count("*").as("eventCount"),
avg("rating").as("avgRating"));
window(timeColumn, windowDuration) — 只传窗口大小,步长自动等于窗口大小,即滚动窗口。
3.2 滑动窗口(对比写法)
// 如果改为滑动窗口,只需加一个参数
parsed.withWatermark("eventTime", "1 minute")
.groupBy(
window(col("eventTime"), "30 seconds", "10 seconds"), // 多了滑动步长
col("userId"),
col("behaviorType"))
.agg(count("*").as("eventCount"),
avg("rating").as("avgRating"));
window(timeColumn, windowDuration, slideDuration) — 第三个参数是滑动步长,窗口大小 > 步长 = 滑动窗口。
代码差异只有一个参数,但语义和计算成本完全不同。
四、底层计算差异
4.1 计算量对比
假设 1 分钟内到达 600 条事件:
| 维度 | 滚动窗口(30s) | 滑动窗口(30s 窗口, 10s 步长) |
|---|---|---|
| 窗口数量 | 2 个 | 6 个 |
| 每条事件参与的窗口 | 1 个 | 3 个 |
| 总计算次数 | 600 × 1 = 600 | 600 × 3 = 1800 |
| 内存占用 | 维护 2 个窗口状态 | 维护 6 个窗口状态 |
公式:
滑动窗口的计算膨胀系数 = ⌈窗口大小 / 滑动步长⌉
示例:
30s 窗口 / 10s 步长 = 3 倍计算量
1min 窗口 / 10s 步长 = 6 倍计算量
5min 窗口 / 30s 步长 = 10 倍计算量
4.2 输出频率对比
滚动窗口(30s):
时间轴: 0s ──── 30s ──── 60s ──── 90s
输出: ① ② ③
每 30 秒输出一次结果
滑动窗口(30s 窗口, 10s 步长):
时间轴: 0s ── 10s ── 20s ── 30s ── 40s ── 50s ── 60s
输出: ① ② ③ ④ ⑤ ⑥
每 10 秒输出一次结果 → 更新更频繁
4.3 结果平滑度对比
假设事件分布:0-10s 来了 100 条,10-30s 来了 0 条
滚动窗口 [0, 30s) 的结果:
count = 100(前 10 秒的突发全部计入)
滑动窗口的结果:
窗口 [0, 30s) → count = 100
窗口 [10, 40s) → count = 0 ← 突发已"滑出"窗口
窗口 [20, 50s) → count = 0
滑动窗口能更快地反映"突发已过去"这个事实。
五、业务场景选型
5.1 适合滚动窗口的业务
场景一:日报/小时报统计
需求:每小时统计一次各商品的销售额
窗口:1 hour,步长 = 1 hour(滚动)
为什么用滚动?
- 每笔订单只应该被统计一次,重复统计会导致财务数据错误
- "上午 10 点到 11 点卖了多少" 是一个精确的时间段问题
- 不需要平滑,需要的是精确切分
Spark 代码:
window(col("orderTime"), "1 hour")
场景二:本项目——用户行为窗口聚合
需求:每 30 秒统计一次用户行为计数和平均评分
窗口:30 seconds(滚动)
为什么用滚动?
- 目的是生成"这 30 秒内用户做了什么"的快照
- 每个事件只计入一个窗口,避免重复计数导致特征膨胀
- 下游推荐算法需要的是"离散时间片"的统计,不是平滑曲线
本项目代码:
window(col("eventTime"), "30 seconds")
场景三:批量数据落盘
需求:每 5 分钟将聚合结果写入数据库
窗口:5 minutes(滚动)
为什么用滚动?
- 数据库写入幂等性要求每条数据只写一次
- 如果用滑动窗口,同一条数据在多个窗口中出现,会导致重复写入
- 滚动窗口天然保证不重不漏
场景四:计费系统
需求:按自然小时统计 API 调用次数,用于计费
窗口:1 hour(滚动)
为什么用滚动?
- 计费必须精确,每次 API 调用只能收费一次
- "10:00-11:00 调用了 1000 次" → 收费 1000 次
- 如果用滑动窗口,同一次调用可能出现在多个窗口中 → 重复计费
滚动窗口的共性:数据不能重复计算,追求精确切分。
5.2 适合滑动窗口的业务
场景一:实时异常检测 / 限流
需求:如果某用户在任意 1 分钟内请求超过 100 次,触发限流
窗口:1 minute,步长 = 10 seconds(滑动)
为什么必须用滑动?
用滚动窗口会出现"边界盲区":
滚动窗口 [00:00, 01:00) 和 [01:00, 02:00)
00:50 ─ 01:10 之间的 20 秒内爆发了 120 次请求
但分散到两个窗口后:
窗口1 中只有 60 次 → 不触发
窗口2 中只有 60 次 → 不触发
→ 漏报!
滑动窗口 [00:50, 01:50) 可以捕获这 120 次
→ 触发限流 ✅
┌─────── 滚动窗口1 ───────┐┌─────── 滚动窗口2 ───────┐
│ 60次 ││ 60次 │
└─────────────────────────┘└─────────────────────────┘
├── 120次爆发 ──┤
↑
这段爆发横跨两个窗口,被各分一半
┌───── 滑动窗口 ─────────────────┐
│ 完整捕获 120 次 │
└────────────────────────────────┘
场景二:实时热搜 / 热门商品排行
需求:展示"最近 5 分钟最热门的 10 个搜索关键词",每 30 秒刷新
窗口:5 minutes,步长 = 30 seconds(滑动)
为什么用滑动?
- 用户看到的排行榜需要"平滑过渡",而不是每 5 分钟突变一次
- 如果用 5 分钟滚动窗口,排行榜每 5 分钟才更新一次,体验差
- 滑动窗口每 30 秒输出一次,但覆盖最近 5 分钟数据 → 平滑且实时
Spark 代码:
window(col("searchTime"), "5 minutes", "30 seconds")
场景三:移动平均(Moving Average)
需求:监控系统 CPU 使用率,展示"最近 10 分钟的平均值",每分钟刷新
窗口:10 minutes,步长 = 1 minute(滑动)
为什么用滑动?
- 运维看板需要"趋势曲线",不是阶梯图
- 滚动窗口产生的是阶梯状数据:
┌──┐ ┌──┐ ┌──┐
│ │ │ │ │ │
│ └─────┘ └──────────┘ │ ← 突变、跳跃
- 滑动窗口产生的是平滑曲线:
╱╲ ╱╲
╱ ╲╱╱ ╲ ← 平滑过渡
场景四:金融风控——可疑交易检测
需求:如果某账户在任意 1 小时内转账总额超过 50 万,触发审核
窗口:1 hour,步长 = 5 minutes(滑动)
为什么必须用滑动?
- 与异常检测同理:犯罪分子不会配合你的窗口边界操作
- 他可能在 10:55 转 30 万,11:05 转 25 万
- 滚动窗口各算各的,看不出来
- 滑动窗口 [10:05, 11:05) 能捕获完整的 55 万 → 触发
Spark 代码:
window(col("transferTime"), "1 hour", "5 minutes")
滑动窗口的共性:需要捕获"任意时间段"内的异常,或需要平滑输出。
六、决策流程图
当你拿到一个新的流式计算需求时,用这个流程判断:
你的需求是什么?
│
┌───────────┴───────────┐
│ │
"统计精确的时间段" "监控任意时间段"
(日报、计费、落盘) (异常检测、限流、风控)
│ │
▼ ▼
滚动窗口 滑动窗口
│
│
┌───────────┴───────────┐
│ │
"输出可以是阶梯状" "输出需要平滑过渡"
(后台统计、批处理) (排行榜、监控看板)
│ │
▼ ▼
滚动窗口 滑动窗口
│
│
┌───────────┴───────────┐
│ │
"每条数据只能算一次" "同一条数据可以参与多次计算"
(计费、去重、财务) (趋势分析、移动平均)
│ │
▼ ▼
滚动窗口 滑动窗口
七、性能影响与调优建议
7.1 滑动窗口的性能陷阱
计算膨胀系数 = ⌈窗口大小 / 滑动步长⌉
危险配置示例:
window("1 hour", "1 second") → 膨胀 3600 倍!
每条数据被复制到 3600 个窗口中 → 内存爆炸 + GC 风暴
7.2 调优原则
| 原则 | 说明 |
|---|---|
| 步长不要太小 | 步长 ≥ 窗口大小 / 10,避免膨胀超过 10 倍 |
| 窗口不要太大 | 窗口越大,Spark 需要维护的状态越多 |
| Watermark 配合 | Watermark 决定了旧窗口何时清理,设太大会占内存 |
| 能用滚动就用滚动 | 滚动窗口是滑动窗口的特例(步长=窗口),性能最优 |
7.3 为什么本项目选择滚动窗口
本项目的需求:
✅ 每 30 秒生成一份用户行为统计 → 精确时间段
✅ 每个事件只应被统计一次 → 不能重复
✅ 结果作为推荐算法特征 → 不需要平滑
✅ 学习项目,性能预算有限 → 计算膨胀系数 = 1
结论:滚动窗口是最优选择。
如果未来需要增加"实时异常检测"功能(如本项目 BehaviorAnalysisService 中的机器人检测),可以在 Spark 层额外添加一个滑动窗口流:
// 异常检测专用:1分钟窗口,10秒步长
Dataset<Row> anomalyStream = parsed
.withWatermark("eventTime", "2 minutes")
.groupBy(
window(col("eventTime"), "1 minute", "10 seconds"),
col("userId"))
.agg(count("*").as("eventCount"))
.filter(col("eventCount").gt(100)); // 任意1分钟超100次 → 可疑
八、总结速查表
| 维度 | 滚动窗口 | 滑动窗口 |
|---|---|---|
| 重叠 | 无 | 有 |
| 数据归属 | 每条数据属于 1 个窗口 | 每条数据属于多个窗口 |
| 输出频率 | 等于窗口大小 | 等于滑动步长 |
| 计算量 | 1x | ⌈窗口/步长⌉x |
| 输出形态 | 阶梯状 | 平滑曲线 |
| 边界盲区 | 有(跨窗口的突发可能漏检) | 无 |
| 典型场景 | 统计报表、计费、批量落盘 | 异常检测、限流、排行榜、风控 |
| Spark 写法 | window(col, "30s") | window(col, "30s", "10s") |
| 本项目选择 | ✅ 当前使用 | 未来异常检测可引入 |
更多推荐
所有评论(0)