Flink 窗口机制
概述
窗口是流处理聚合的核心:把无限流按时间或数量切分为有限区间再计算。Flink 提供滚动、滑动、会话、全局四种窗口,配合触发器(何时计算)与驱逐器(保留哪些数据)。本文逐一讲透。
一、窗口类型
1.1 时间窗口 vs 计数窗口
| 维度 | 时间窗口 | 计数窗口 |
|---|---|---|
| 划分依据 | 时间区间 | 数据条数 |
| 窗口算子 | Tumbling/Sliding/Session | Count |
| 典型场景 | 分钟级聚合 | 固定条数处理 |
1.2 四种窗口
| 窗口 | 特征 | 示例 |
|---|---|---|
| Tumbling(滚动) | 固定长度、无重叠 | 每 1 分钟统计 |
| Sliding(滑动) | 固定长度、可重叠 | 每 5 秒统计近 1 分钟 |
| Session(会话) | 按活跃间隙切分 | 用户访问会话 |
| Global(全局) | 所有数据一个窗口 | 配合触发器 |
1.3 示意图
滚动窗口(1 分钟):
[0-60) [60-120) [120-180)
滑动窗口(长度 60s,滑动 30s):
[0-60)
[30-90)
[60-120)
会话窗口(间隙 5s):
数据:t0 t1 ... t58 t65 t70 t90...
└─ 会话1 ─┘ └会话2┘二、窗口使用
2.1 滚动窗口
java
stream
.keyBy(...)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.aggregate(new CountAgg());2.2 滑动窗口
java
.window(SlidingEventTimeWindows.of(
Time.minutes(1), Time.seconds(30))) // 长度,滑动2.3 会话窗口
java
.window(EventTimeSessionWindows.withGap(Time.seconds(5)))2.4 计数窗口
java
.keyBy(...).countWindow(100) // 滚动 100 条
.keyBy(...).countWindow(100, 10) // 滑动 100 条,步长 10三、窗口函数
| 函数 | 类型 | 说明 |
|---|---|---|
apply | 全量 | 窗口内全数据参与 |
reduce | 增量 | 逐条增量归约 |
aggregate | 增量 | 自定义增量聚合 |
process | 底层 | 全窗口 + 上下文 |
| 对比 | 增量 | 全量 |
|---|---|---|
| 内存 | 只存状态 | 缓存全部数据 |
| 性能 | 高 | 低 |
| 适用 | 求和/计数 | 复杂计算 |
java
// 增量聚合(高效)
.window(...).reduce((a, b) -> a + b);
// 底层处理(可访问时间/状态/输出迟到)
.window(...).process(new ProcessWindowFunction<>() {...});四、触发器 Triggers
4.1 作用
触发器决定窗口何时计算结果并发出。
4.2 内置触发器
| 触发器 | 触发时机 |
|---|---|
| EventTimeTrigger | 水位线越过窗口结束时间 |
| ProcessingTimeTrigger | 处理时间到达结束 |
| CountTrigger | 条数达到阈值 |
| ContinuousEventTimeTrigger | 周期性触发 |
| 默认 | 时间窗口按时间,计数窗口按条数 |
4.3 自定义触发器
java
window.trigger(new Trigger<...>() {
@Override
public TriggerResult onElement(...) {
if (条件) return TriggerResult.FIRE_AND_PURGE;
return TriggerResult.CONTINUE;
}
@Override
public TriggerResult onProcessingTime(...) {...}
@Override
public TriggerResult onEventTime(...) {...}
@Override
public TriggerResult clear(...) {...}
});| 返回 | 行为 |
|---|---|
| CONTINUE | 不触发 |
| FIRE | 计算并保留窗口 |
| PURGE | 清除不计算 |
| FIRE_AND_PURGE | 计算并清除 |
五、驱逐器 Evictors
5.1 作用
窗口触发计算前,从窗口中移除指定数据。
5.2 内置驱逐器
| 驱逐器 | 行为 |
|---|---|
| CountEvictor | 保留最近 N 条 |
| DeltaEvictor | 按差值删除 |
| TimeEvictor | 保留最近时间段 |
java
.window(...)
.evictor(CountEvictor.of(100)) // 只保留最近 100 条5.3 与触发器的配合
| 组合 | 场景 |
|---|---|
| Trigger + Evictor | 滚动窗口 + 保留子集 |
| 只用 Trigger | 自定义触发时机 |
六、窗口生命周期
1. 首条数据到达 → 窗口创建(状态注册)
2. 数据按 key 进入窗口,累积
3. 触发器触发 → 调用窗口函数
4. 窗口结束时间 + 延迟 → 窗口清除(状态删除)| 事件 | 行为 |
|---|---|
| 窗口创建 | 惰性(有数据才建) |
| 窗口结束 | 水位线越过结束时间 |
| 清理 | 清除状态与定时器 |
七、常见问题
7.1 乱序与窗口
| 问题 | 处理 |
|---|---|
| 数据晚到 | 水位线容忍 + 延迟数据处理(侧输出) |
| 窗口结果晚出 | 水位线策略与延迟时间权衡 |
| 重叠窗口重复计算 | 滑动窗口天然重叠,结果含重复区间 |
7.2 性能
| 优化 | 说明 |
|---|---|
| 增量聚合 | 避免全量缓存 |
| 减小窗口状态 | 及时清理 |
| 并行度 | 窗口计算按 key 并行 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 窗口不触发 | 无水位线或水位线未推进 |
| 结果延迟 | 水位线太保守,收紧或降延迟 |
| 会话窗口分裂 | 间隙参数过大导致窗口合并 |
| 状态堆积 | 窗口未清理,检查触发器与驱逐器 |
| 计数窗口内存大 | 窗口过大,改增量聚合 |