欢迎光临
我们一直在努力

一文搞懂Flink 窗口设计深入详解(含源码分析)

Flink 窗口设计深入详解(含源码分析)

    • 1. 窗口概述
      • 1.1 什么是窗口
      • 1.2 窗口的核心组件
    • 2. 窗口类型详解
      • 2.1 滚动窗口(Tumbling Window)
        • DataStream API 示例
        • SQL 示例
      • 2.2 滑动窗口(Sliding Window / Hop Window)
        • DataStream API 示例
        • SQL 示例
      • 2.3 会话窗口(Session Window)
        • DataStream API 示例
        • SQL 示例
      • 2.4 全局窗口(Global Window)
    • 3. 窗口分配器源码
      • 3.1 WindowAssigner 接口
      • 3.2 TumblingEventTimeWindows 源码
        • 窗口起始时间计算详解
      • 3.3 SlidingEventTimeWindows 源码
      • 3.4 EventTimeSessionWindows 源码
    • 4. 窗口触发器源码
      • 4.1 Trigger 接口
        • TriggerResult 枚举
      • 4.2 EventTimeTrigger 源码
      • 4.3 ProcessingTimeTrigger 源码
      • 4.4 CountTrigger 源码
      • 4.5 ContinuousEventTimeTrigger 源码(早期触发)
    • 5. 窗口驱逐器源码
      • 5.1 Evictor 接口
      • 5.2 CountEvictor 源码
      • 5.3 TimeEvictor 源码
    • 6. 窗口状态管理
      • 6.1 窗口状态存储结构
      • 6.2 窗口状态类型
      • 6.3 窗口状态生命周期
      • 6.4 窗口状态大小估算
    • 7. Watermark 与窗口触发
      • 7.1 Watermark 生成策略
      • 7.2 Watermark 传播机制
      • 7.3 窗口触发时机详解
      • 7.4 迟到数据处理
    • 8. 窗口计算流程源码分析
      • 8.1 WindowOperator 核心方法
      • 8.2 窗口计算执行流程
      • 8.3 窗口清理流程
      • 8.4 增量聚合 vs 全量聚合
        • 增量聚合(ReduceFunction / AggregateFunction)
        • 全量聚合(ProcessWindowFunction)
        • 混合模式(增量 + 全量)
    • 9. SQL 窗口实现原理
      • 9.1 SQL 窗口函数映射
      • 9.2 TUMBLE 窗口 SQL → DataStream 转换
        • SQL 代码
        • 等价的 DataStream 代码
      • 9.3 OVER 窗口实现原理
        • SQL 代码
        • 底层实现原理
      • 9.4 SQL 窗口聚合优化
    • 10. 性能优化最佳实践
      • 10.1 选择合适的窗口类型
      • 10.2 减小窗口状态大小
        • ❌ 避免:全量存储 + 大窗口
        • ✅ 推荐:增量聚合 + 合理窗口大小
      • 10.3 滑动窗口优化
        • ❌ 避免:小步长导致大量重叠
        • ✅ 推荐:调整步长或改用其他方案
      • 10.4 State TTL 配置
      • 10.5 早期触发(Early Firing)
      • 10.6 Checkpoint 优化
      • 10.7 水位线对齐优化
    • 11. 窗口性能监控
      • 11.1 关键指标
      • 11.2 问题排查
    • 12. 总结
      • 12.1 窗口设计核心要点
      • 12.2 性能优化金字塔
      • 12.3 技术选型决策树

1. 窗口概述

1.1 什么是窗口

窗口(Window)是处理无界流的核心机制,将无限的数据流切分为有限的"桶"(bucket),在每个桶上执行计算。

无界流数据: ———[1]———[2]———[3]———[4]———[5]———[6]———[7]———→
↓ ↓ ↓ ↓
滚动窗口: [1,2,3] [4,5,6] [7,8,9] [10,11,12]

1.2 窗口的核心组件

Flink 窗口由四个核心组件构成:

DataStream<T> input = ...;
input
.keyBy(<key selector>) // 1. KeyBy(可选)
.window(<window assigner>) // 2. Window Assigner(窗口分配器)
.trigger(<trigger>) // 3. Trigger(触发器,可选)
.evictor(<evictor>) // 4. Evictor(驱逐器,可选)
.process(<window function>); // 5. Window Function(窗口函数)

组件作用是否必需
Window Assigner 决定数据分配到哪个窗口 ✅ 必需
Trigger 决定何时触发窗口计算 ⚠️ 可选(有默认值)
Evictor 决定哪些数据从窗口移除 ⚠️ 可选
Window Function 定义窗口计算逻辑 ✅ 必需

2. 窗口类型详解

2.1 滚动窗口(Tumbling Window)

特点:固定大小,无重叠

时间轴: 0 1 2 3 4 5 6 7 8 9 10
[ 窗口1 ] [ 窗口2 ] [ 窗口3 ]

DataStream API 示例

// 1. 时间滚动窗口(5 秒)
DataStream<Tuple2<String, Long>> input = ...;
input
.keyBy(value -> value.f0)
.window(TumblingEventTimeWindows.of(Time.seconds(5)))
.sum(1);

// 2. 计数滚动窗口(每 100 条触发)
input
.keyBy(value -> value.f0)
.countWindow(100);

SQL 示例

— TUMBLE 函数:[窗口开始, 窗口结束)
SELECT
user_id,
TUMBLE_START(event_time, INTERVAL '5' SECOND) AS window_start,
TUMBLE_END(event_time, INTERVAL '5' SECOND) AS window_end,
COUNT(*) AS cnt,
SUM(amount) AS total_amount
FROM orders
GROUP BY
user_id,
TUMBLE(event_time, INTERVAL '5' SECOND);

输出示例:

user_id | window_start | window_end | cnt | total_amount
——–|———————|———————|—–|————-
user1 | 2025-01-12 10:00:00 | 2025-01-12 10:00:05 | 3 | 150.0
user1 | 2025-01-12 10:00:05 | 2025-01-12 10:00:10 | 5 | 280.0


2.2 滑动窗口(Sliding Window / Hop Window)

特点:固定大小,可重叠

窗口大小: 6 秒,滑动步长: 3 秒

时间轴: 0 1 2 3 4 5 6 7 8 9 10
[ 窗口1 ]
[ 窗口2 ]
[ 窗口3 ]

DataStream API 示例

// 1. 时间滑动窗口(窗口 10 秒,滑动 5 秒)
input
.keyBy(value -> value.f0)
.window(SlidingEventTimeWindows.of(
Time.seconds(10), // 窗口大小
Time.seconds(5) // 滑动步长
))
.sum(1);

// 2. 计数滑动窗口(窗口 100 条,滑动 10 条)
input
.keyBy(value -> value.f0)
.countWindow(100, 10);

SQL 示例

— HOP 函数:窗口大小 10 分钟,滑动步长 5 分钟
SELECT
user_id,
HOP_START(event_time, INTERVAL '5' MINUTE, INTERVAL '10' MINUTE) AS window_start,
HOP_END(event_time, INTERVAL '5' MINUTE, INTERVAL '10' MINUTE) AS window_end,
COUNT(*) AS cnt
FROM orders
GROUP BY
user_id,
HOP(event_time, INTERVAL '5' MINUTE, INTERVAL '10' MINUTE);

数据重复说明:

  • 滑动窗口会导致数据被计算多次
  • 例如:时间戳为 10:00:03 的数据会同时出现在 [10:00:00, 10:00:06) 和 [10:00:03, 10:00:09) 两个窗口中

2.3 会话窗口(Session Window)

特点:动态大小,根据活动间隔(gap)分组

数据到达: [1] [2][3] [4] [5][6][7]
↓ ↓ ↓ ↓ ↓ ↓ ↓
会话窗口: [ 会话1 ] (gap) [ 会话2 ]
(gap > 30s) (gap < 30s,在同一会话)

DataStream API 示例

// 1. 固定 gap 会话窗口(30 秒不活动则关闭会话)
input
.keyBy(value -> value.f0)
.window(EventTimeSessionWindows.withGap(Time.seconds(30)))
.sum(1);

// 2. 动态 gap 会话窗口(根据数据内容决定 gap)
input
.keyBy(value -> value.f0)
.window(EventTimeSessionWindows.withDynamicGap(
new SessionWindowTimeGapExtractor<Tuple2<String, Long>>() {
@Override
public long extract(Tuple2<String, Long> element) {
// VIP 用户 60 秒,普通用户 30 秒
return element.f0.contains("vip") ? 60000L : 30000L;
}
}
))
.sum(1);

SQL 示例

— SESSION 函数:30 分钟不活动则结束会话
SELECT
user_id,
SESSION_START(event_time, INTERVAL '30' MINUTE) AS session_start,
SESSION_END(event_time, INTERVAL '30' MINUTE) AS session_end,
COUNT(*) AS page_views,
MAX(event_time) AS last_active_time
FROM user_events
GROUP BY
user_id,
SESSION(event_time, INTERVAL '30' MINUTE);

适用场景:

  • 用户会话分析(网站浏览行为)
  • IoT 设备间歇性数据上报
  • 客服对话分组

2.4 全局窗口(Global Window)

特点:所有数据分配到同一个窗口,需要自定义 Trigger

// 全局窗口 + 自定义触发器
input
.keyBy(value -> value.f0)
.window(GlobalWindows.create())
.trigger(CountTrigger.of(100)) // 每 100 条触发一次
.sum(1);


3. 窗口分配器源码

3.1 WindowAssigner 接口

/**
* 窗口分配器接口:决定数据元素分配到哪个窗口
*/

public abstract class WindowAssigner<T, W extends Window> implements Serializable {

/**
* 为元素分配窗口(一个元素可能分配到多个窗口,如滑动窗口)
* @param element 数据元素
* @param timestamp 元素时间戳
* @param context 分配器上下文
* @return 该元素所属的窗口集合
*/

public abstract Collection<W> assignWindows(
T element,
long timestamp,
WindowAssignerContext context
);

/**
* 返回默认触发器
*/

public abstract Trigger<T, W> getDefaultTrigger(StreamExecutionEnvironment env);

/**
* 返回窗口序列化器
*/

public abstract TypeSerializer<W> getWindowSerializer(ExecutionConfig executionConfig);

/**
* 是否基于事件时间
*/

public abstract boolean isEventTime();
}


3.2 TumblingEventTimeWindows 源码

/**
* 事件时间滚动窗口分配器
*/

public class TumblingEventTimeWindows extends WindowAssigner<Object, TimeWindow> {

private static final long serialVersionUID = 1L;

private final long size; // 窗口大小
private final long offset; // 窗口偏移量(用于对齐时区)

private TumblingEventTimeWindows(long size, long offset) {
if (Math.abs(offset) >= size) {
throw new IllegalArgumentException("offset 必须小于窗口大小");
}
this.size = size;
this.offset = offset;
}

/**
* 创建滚动窗口(无偏移)
*/

public static TumblingEventTimeWindows of(Time size) {
return new TumblingEventTimeWindows(size.toMilliseconds(), 0);
}

/**
* 创建滚动窗口(带偏移,用于时区对齐)
* 例如:中国时区偏移 +8 小时
*/

public static TumblingEventTimeWindows of(Time size, Time offset) {
return new TumblingEventTimeWindows(
size.toMilliseconds(),
offset.toMilliseconds()
);
}

@Override
public Collection<TimeWindow> assignWindows(
Object element,
long timestamp,
WindowAssignerContext context
) {
if (timestamp > Long.MIN_VALUE) {
// 核心算法:计算元素所属窗口的起始时间
long start = TimeWindow.getWindowStartWithOffset(
timestamp,
offset,
size
);
return Collections.singletonList(new TimeWindow(start, start + size));
} else {
throw new RuntimeException("数据时间戳无效,无法分配窗口");
}
}

@Override
public Trigger<Object, TimeWindow> getDefaultTrigger(StreamExecutionEnvironment env) {
return EventTimeTrigger.create(); // 事件时间触发器
}

@Override
public boolean isEventTime() {
return true;
}
}

窗口起始时间计算详解

/**
* TimeWindow 类中的窗口起始时间计算方法
*/

public static long getWindowStartWithOffset(long timestamp, long offset, long windowSize) {
// 核心公式:窗口起始时间 = timestamp – (timestamp – offset + windowSize) % windowSize
return timestamp (timestamp offset + windowSize) % windowSize;
}

计算示例:

假设:
– timestamp = 12345 ms
– windowSize = 5000 ms(5 秒)
– offset = 0

计算过程:
1. timestamp – offset + windowSize = 12345 – 0 + 5000 = 17345
2. 17345 % 5000 = 2345(余数)
3. start = 12345 – 2345 = 10000
4. end = start + windowSize = 10000 + 5000 = 15000

结果:该元素分配到窗口 [10000, 15000)


3.3 SlidingEventTimeWindows 源码

/**
* 事件时间滑动窗口分配器
*/

public class SlidingEventTimeWindows extends WindowAssigner<Object, TimeWindow> {

private final long size; // 窗口大小
private final long slide; // 滑动步长
private final long offset;

private SlidingEventTimeWindows(long size, long slide, long offset) {
if (Math.abs(offset) >= slide || size <= 0) {
throw new IllegalArgumentException("窗口参数非法");
}
this.size = size;
this.slide = slide;
this.offset = offset;
}

public static SlidingEventTimeWindows of(Time size, Time slide) {
return new SlidingEventTimeWindows(
size.toMilliseconds(),
slide.toMilliseconds(),
0
);
}

@Override
public Collection<TimeWindow> assignWindows(
Object element,
long timestamp,
WindowAssignerContext context
) {
if (timestamp > Long.MIN_VALUE) {
List<TimeWindow> windows = new ArrayList<>((int) (size / slide));

// 计算最后一个窗口的起始时间
long lastStart = TimeWindow.getWindowStartWithOffset(
timestamp,
offset,
slide
);

// 核心:向前计算所有包含该元素的窗口
for (long start = lastStart;
start > timestamp size;
start -= slide) {
windows.add(new TimeWindow(start, start + size));
}

return windows;
} else {
throw new RuntimeException("数据时间戳无效");
}
}

@Override
public Trigger<Object, TimeWindow> getDefaultTrigger(StreamExecutionEnvironment env) {
return EventTimeTrigger.create();
}
}

滑动窗口分配示例:

配置:size = 10 秒,slide = 5 秒
timestamp = 12 秒的数据

计算过程:
1. lastStart = 10 秒(最后一个窗口起始)
2. 循环计算:
– start = 10: 窗口 [10, 20) ✅ 包含 12 秒
– start = 5: 窗口 [5, 15) ✅ 包含 12 秒
– start = 0: 窗口 [0, 10) ❌ 不包含 12 秒(终止)

结果:该元素分配到 [10,20) 和 [5,15) 两个窗口


3.4 EventTimeSessionWindows 源码

/**
* 事件时间会话窗口分配器
*/

public class EventTimeSessionWindows extends MergingWindowAssigner<Object, TimeWindow> {

private final long sessionTimeout; // 会话超时时间

private EventTimeSessionWindows(long sessionTimeout) {
if (sessionTimeout <= 0) {
throw new IllegalArgumentException("会话超时时间必须 > 0");
}
this.sessionTimeout = sessionTimeout;
}

public static EventTimeSessionWindows withGap(Time gap) {
return new EventTimeSessionWindows(gap.toMilliseconds());
}

@Override
public Collection<TimeWindow> assignWindows(
Object element,
long timestamp,
WindowAssignerContext context
) {
// 会话窗口初始分配:为每个元素创建独立窗口
// 窗口范围:[timestamp, timestamp + sessionTimeout)
return Collections.singletonList(
new TimeWindow(timestamp, timestamp + sessionTimeout)
);
}

/**
* 核心方法:合并重叠的会话窗口
*/

@Override
public void mergeWindows(
Collection<TimeWindow> windows,
MergeCallback<TimeWindow> callback
) {
// 1. 按窗口起始时间排序
TimeWindow[] sortedWindows = windows.toArray(new TimeWindow[0]);
Arrays.sort(sortedWindows, new Comparator<TimeWindow>() {
@Override
public int compare(TimeWindow o1, TimeWindow o2) {
return Long.compare(o1.getStart(), o2.getStart());
}
});

// 2. 合并重叠窗口
List<TimeWindow> merged = new ArrayList<>();
TimeWindow currentMerge = sortedWindows[0];

for (int i = 1; i < sortedWindows.length; i++) {
TimeWindow next = sortedWindows[i];

// 如果下一个窗口的起始时间 <= 当前合并窗口的结束时间,则合并
if (next.getStart() <= currentMerge.getEnd()) {
currentMerge = new TimeWindow(
currentMerge.getStart(),
Math.max(currentMerge.getEnd(), next.getEnd())
);
} else {
merged.add(currentMerge);
currentMerge = next;
}
}
merged.add(currentMerge);

// 3. 回调合并结果
for (TimeWindow mergedWindow : merged) {
callback.merge(mergedWindow, mergedWindow);
}
}

@Override
public Trigger<Object, TimeWindow> getDefaultTrigger(StreamExecutionEnvironment env) {
return EventTimeTrigger.create();
}
}

会话窗口合并示例:

会话超时:30 秒

数据到达时间轴:
10:00:00 → 创建窗口 [10:00:00, 10:00:30)
10:00:15 → 创建窗口 [10:00:15, 10:00:45)
10:00:50 → 创建窗口 [10:00:50, 10:01:20)

合并过程:
1. 窗口1 和窗口2 重叠(15 < 30),合并为 [10:00:00, 10:00:45)
2. 窗口3 不重叠(50 > 45),独立会话 [10:00:50, 10:01:20)

最终会话:
– 会话1:[10:00:00, 10:00:45)
– 会话2:[10:00:50, 10:01:20)


4. 窗口触发器源码

4.1 Trigger 接口

/**
* 触发器接口:决定何时触发窗口计算、何时清理窗口
*/

public abstract class Trigger<T, W extends Window> implements Serializable {

/**
* 每个元素加入窗口时调用
* @return TriggerResult(触发结果)
*/

public abstract TriggerResult onElement(
T element,
long timestamp,
W window,
TriggerContext ctx
) throws Exception;

/**
* 处理时间定时器触发时调用
*/

public abstract TriggerResult onProcessingTime(
long time,
W window,
TriggerContext ctx
) throws Exception;

/**
* 事件时间定时器触发时调用
*/

public abstract TriggerResult onEventTime(
long time,
W window,
TriggerContext ctx
) throws Exception;

/**
* 窗口合并时调用(仅用于会话窗口)
*/

public void onMerge(W window, OnMergeContext ctx) throws Exception {
throw new UnsupportedOperationException("该触发器不支持合并");
}

/**
* 窗口清理时调用
*/

public abstract void clear(W window, TriggerContext ctx) throws Exception;
}

TriggerResult 枚举

/**
* 触发器返回结果
*/

public enum TriggerResult {
CONTINUE, // 什么都不做,继续等待
FIRE, // 触发计算,但保留窗口数据
PURGE, // 清理窗口数据,但不触发计算
FIRE_AND_PURGE // 触发计算并清理窗口数据
}


4.2 EventTimeTrigger 源码

/**
* 事件时间触发器:当 Watermark >= 窗口结束时间时触发
*/

public class EventTimeTrigger extends Trigger<Object, TimeWindow> {

private EventTimeTrigger() {}

public static EventTimeTrigger create() {
return new EventTimeTrigger();
}

@Override
public TriggerResult onElement(
Object element,
long timestamp,
TimeWindow window,
TriggerContext ctx
) {
if (window.maxTimestamp() <= ctx.getCurrentWatermark()) {
// Watermark 已超过窗口结束时间,立即触发
return TriggerResult.FIRE;
} else {
// 注册事件时间定时器:在窗口结束时间触发
ctx.registerEventTimeTimer(window.maxTimestamp());
return TriggerResult.CONTINUE;
}
}

@Override
public TriggerResult onEventTime(
long time,
TimeWindow window,
TriggerContext ctx
) {
// 定时器触发时,检查是否达到窗口结束时间
return time == window.maxTimestamp()
? TriggerResult.FIRE
: TriggerResult.CONTINUE;
}

@Override
public TriggerResult onProcessingTime(
long time,
TimeWindow window,
TriggerContext ctx
) {
return TriggerResult.CONTINUE; // 事件时间触发器不响应处理时间
}

@Override
public void clear(TimeWindow window, TriggerContext ctx) {
// 清理注册的定时器
ctx.deleteEventTimeTimer(window.maxTimestamp());
}

@Override
public String toString() {
return "EventTimeTrigger()";
}
}

关键点:

  • window.maxTimestamp() = window.getEnd() – 1(窗口最大时间戳)
  • Flink 的窗口是左闭右开:[start, end)
  • 触发条件:Watermark >= window.end – 1

4.3 ProcessingTimeTrigger 源码

/**
* 处理时间触发器:当处理时间 >= 窗口结束时间时触发
*/

public class ProcessingTimeTrigger extends Trigger<Object, TimeWindow> {

private ProcessingTimeTrigger() {}

public static ProcessingTimeTrigger create() {
return new ProcessingTimeTrigger();
}

@Override
public TriggerResult onElement(
Object element,
long timestamp,
TimeWindow window,
TriggerContext ctx
) {
// 注册处理时间定时器:在窗口结束时间触发
ctx.registerProcessingTimeTimer(window.maxTimestamp());
return TriggerResult.CONTINUE;
}

@Override
public TriggerResult onProcessingTime(
long time,
TimeWindow window,
TriggerContext ctx
) {
return TriggerResult.FIRE; // 处理时间到达,立即触发
}

@Override
public TriggerResult onEventTime(
long time,
TimeWindow window,
TriggerContext ctx
) {
return TriggerResult.CONTINUE; // 处理时间触发器不响应事件时间
}

@Override
public void clear(TimeWindow window, TriggerContext ctx) {
ctx.deleteProcessingTimeTimer(window.maxTimestamp());
}
}


4.4 CountTrigger 源码

/**
* 计数触发器:窗口内元素数量达到阈值时触发
*/

public class CountTrigger<W extends Window> extends Trigger<Object, W> {

private final long maxCount; // 触发阈值

// 用于存储窗口内当前元素数量的 State
private final ReducingStateDescriptor<Long> stateDesc =
new ReducingStateDescriptor<>(
"count",
new Sum(),
LongSerializer.INSTANCE
);

private CountTrigger(long maxCount) {
this.maxCount = maxCount;
}

public static <W extends Window> CountTrigger<W> of(long maxCount) {
return new CountTrigger<>(maxCount);
}

@Override
public TriggerResult onElement(
Object element,
long timestamp,
W window,
TriggerContext ctx
) throws Exception {
// 获取窗口内元素计数器
ReducingState<Long> count = ctx.getPartitionedState(stateDesc);
count.add(1L); // 计数 +1

if (count.get() >= maxCount) {
count.clear(); // 重置计数器
return TriggerResult.FIRE; // 触发计算
}
return TriggerResult.CONTINUE;
}

@Override
public TriggerResult onProcessingTime(long time, W window, TriggerContext ctx) {
return TriggerResult.CONTINUE;
}

@Override
public TriggerResult onEventTime(long time, W window, TriggerContext ctx) {
return TriggerResult.CONTINUE;
}

@Override
public void clear(W window, TriggerContext ctx) throws Exception {
ctx.getPartitionedState(stateDesc).clear();
}

/**
* 求和函数
*/

private static class Sum implements ReduceFunction<Long> {
@Override
public Long reduce(Long value1, Long value2) {
return value1 + value2;
}
}
}


4.5 ContinuousEventTimeTrigger 源码(早期触发)

/**
* 连续事件时间触发器:周期性触发 + 窗口结束时触发
* 适用场景:需要看到窗口中间结果(Early Firing)
*/

public class ContinuousEventTimeTrigger extends Trigger<Object, TimeWindow> {

private final long interval; // 触发间隔

private ContinuousEventTimeTrigger(long interval) {
this.interval = interval;
}

public static ContinuousEventTimeTrigger of(Time interval) {
return new ContinuousEventTimeTrigger(interval.toMilliseconds());
}

@Override
public TriggerResult onElement(
Object element,
long timestamp,
TimeWindow window,
TriggerContext ctx
) {
if (window.maxTimestamp() <= ctx.getCurrentWatermark()) {
return TriggerResult.FIRE;
}

// 注册周期性定时器
long nextFireTimestamp = Math.max(
timestamp (timestamp % interval),
window.getStart()
);

while (nextFireTimestamp < window.maxTimestamp()) {
ctx.registerEventTimeTimer(nextFireTimestamp);
nextFireTimestamp += interval;
}

// 注册窗口结束定时器
ctx.registerEventTimeTimer(window.maxTimestamp());

return TriggerResult.CONTINUE;
}

@Override
public TriggerResult onEventTime(
long time,
TimeWindow window,
TriggerContext ctx
) {
if (time == window.maxTimestamp()) {
return TriggerResult.FIRE_AND_PURGE; // 窗口结束,触发并清理
} else {
// 周期性触发,但保留数据
long nextFireTimestamp = time + interval;
if (nextFireTimestamp < window.maxTimestamp()) {
ctx.registerEventTimeTimer(nextFireTimestamp);
}
return TriggerResult.FIRE; // 触发但不清理
}
}

@Override
public void clear(TimeWindow window, TriggerContext ctx) {
// 清理所有定时器(删除所有注册的定时器)
long timestamp = window.getStart();
while (timestamp <= window.maxTimestamp()) {
ctx.deleteEventTimeTimer(timestamp);
timestamp += interval;
}
}
}

使用示例:

// 每 5 秒输出一次窗口中间结果,窗口大小 1 分钟
input
.keyBy(...)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.trigger(ContinuousEventTimeTrigger.of(Time.seconds(5)))
.process(new MyWindowFunction());


5. 窗口驱逐器源码

5.1 Evictor 接口

/**
* 驱逐器接口:决定哪些元素应该从窗口中移除
*/

public interface Evictor<T, W extends Window> extends Serializable {

/**
* 窗口计算前驱逐元素
* @param elements 窗口内所有元素
* @param size 窗口元素数量
* @param window 窗口
* @param evictorContext 驱逐器上下文
*/

void evictBefore(
Iterable<TimestampedValue<T>> elements,
int size,
W window,
EvictorContext evictorContext
);

/**
* 窗口计算后驱逐元素
*/

void evictAfter(
Iterable<TimestampedValue<T>> elements,
int size,
W window,
EvictorContext evictorContext
);
}


5.2 CountEvictor 源码

/**
* 计数驱逐器:保留窗口内最新的 N 个元素
*/

public class CountEvictor<W extends Window> implements Evictor<Object, W> {

private final long maxCount; // 最大保留数量
private final boolean doEvictAfter; // 是否在计算后驱逐

private CountEvictor(long maxCount, boolean doEvictAfter) {
this.maxCount = maxCount;
this.doEvictAfter = doEvictAfter;
}

public static <W extends Window> CountEvictor<W> of(long maxCount) {
return new CountEvictor<>(maxCount, false);
}

@Override
public void evictBefore(
Iterable<TimestampedValue<Object>> elements,
int size,
W window,
EvictorContext ctx
) {
if (!doEvictAfter) {
evict(elements, size, ctx);
}
}

@Override
public void evictAfter(
Iterable<TimestampedValue<Object>> elements,
int size,
W window,
EvictorContext ctx
) {
if (doEvictAfter) {
evict(elements, size, ctx);
}
}

/**
* 驱逐逻辑:移除最旧的元素,保留最新的 maxCount 个
*/

private void evict(
Iterable<TimestampedValue<Object>> elements,
int size,
EvictorContext ctx
) {
if (size <= maxCount) {
return; // 元素数量未超过阈值,不驱逐
}

int evictedCount = 0;
int toEvict = (int) (size maxCount);

for (Iterator<TimestampedValue<Object>> iterator = elements.iterator();
iterator.hasNext(); ) {
iterator.next();
if (evictedCount < toEvict) {
iterator.remove(); // 移除最旧的元素
evictedCount++;
} else {
break;
}
}
}
}

使用示例:

// 滑动窗口 + 计数驱逐器:只保留最新 100 条数据
input
.keyBy(...)
.window(SlidingEventTimeWindows.of(Time.minutes(10), Time.minutes(1)))
.evictor(CountEvictor.of(100))
.sum(1);


5.3 TimeEvictor 源码

/**
* 时间驱逐器:保留窗口内指定时间范围内的元素
*/

public class TimeEvictor<W extends Window> implements Evictor<Object, W> {

private final long windowSize; // 保留时间窗口大小
private final boolean doEvictAfter;

private TimeEvictor(long windowSize, boolean doEvictAfter) {
this.windowSize = windowSize;
this.doEvictAfter = doEvictAfter;
}

public static <W extends Window> TimeEvictor<W> of(Time windowSize) {
return new TimeEvictor<>(windowSize.toMilliseconds(), false);
}

@Override
public void evictBefore(
Iterable<TimestampedValue<Object>> elements,
int size,
W window,
EvictorContext ctx
) {
if (!doEvictAfter) {
evict(elements, size, ctx);
}
}

@Override
public void evictAfter(
Iterable<TimestampedValue<Object>> elements,
int size,
W window,
EvictorContext ctx
) {
if (doEvictAfter) {
evict(elements, size, ctx);
}
}

/**
* 驱逐逻辑:移除时间戳 < (maxTimestamp – windowSize) 的元素
*/

private void evict(
Iterable<TimestampedValue<Object>> elements,
int size,
EvictorContext ctx
) {
if (!elements.iterator().hasNext()) {
return;
}

// 找到最大时间戳
long maxTimestamp = Long.MIN_VALUE;
for (TimestampedValue<Object> element : elements) {
maxTimestamp = Math.max(maxTimestamp, element.getTimestamp());
}

// 驱逐时间戳小于 (maxTimestamp – windowSize) 的元素
long evictCutoff = maxTimestamp windowSize;
for (Iterator<TimestampedValue<Object>> iterator = elements.iterator();
iterator.hasNext(); ) {
TimestampedValue<Object> element = iterator.next();
if (element.getTimestamp() <= evictCutoff) {
iterator.remove();
}
}
}
}


6. 窗口状态管理

6.1 窗口状态存储结构

Flink 窗口状态存储在 WindowOperator 中,主要包括:

/**
* WindowOperator 核心状态(简化版源码)
*/

public class WindowOperator<K, IN, OUT, W extends Window>
extends AbstractUdfStreamOperator<OUT, InternalWindowFunction<IN, OUT, K, W>> {

// 1. 窗口状态:存储窗口内的所有元素
private transient InternalAppendingState<K, W, IN, ACC, ACC> windowState;

// 2. 合并状态(仅用于会话窗口)
private transient InternalMergingState<K, W, IN, ACC, ACC> windowMergingState;

// 3. 触发器状态
private transient InternalValueState<K, W, MergingWindowSet<W>> mergingSetsState;

// 4. 窗口定时器 Service
private transient InternalTimerService<W> internalTimerService;
}


6.2 窗口状态类型

状态类型适用场景State 实现
ListState 需要访问窗口内所有元素(如 ProcessWindowFunction) InternalListState
ReducingState 增量聚合(如 sum, min, max) InternalReducingState
AggregatingState 自定义聚合逻辑 InternalAggregatingState
FoldingState 折叠聚合(已废弃) InternalFoldingState

6.3 窗口状态生命周期

/**
* WindowOperator 中的窗口触发和清理逻辑(简化版)
*/

@Override
public void onEventTime(InternalTimer<K, W> timer) throws Exception {
W window = timer.getNamespace(); // 获取窗口

// 1. 触发窗口计算
TriggerResult triggerResult = context.onEventTime(timer.getTimestamp());

if (triggerResult.isFire()) {
// 2. 获取窗口状态
ACC contents = windowState.get();

if (contents != null) {
// 3. 执行窗口函数
emitWindowContents(window, contents);
}
}

if (triggerResult.isPurge()) {
// 4. 清理窗口状态
windowState.clear();
}

// 5. 如果窗口结束,注册清理定时器
if (isCleanupTime(window, timer.getTimestamp())) {
clearAllState(window);
}
}

/**
* 清理所有窗口相关状态
*/

private void clearAllState(W window) throws Exception {
// 清理窗口数据
windowState.clear();

// 清理触发器状态
trigger.clear(window, triggerContext);

// 删除定时器
context.deleteEventTimeTimer(window.maxTimestamp());
context.deleteCleanupTimer(window);
}


6.4 窗口状态大小估算

/**
* 窗口状态大小计算公式
*/

窗口状态大小 =
并行度
× key 数量
× 并发窗口数
× 每个窗口的元素数量
× 每个元素的大小

/**
* 示例计算
*/

假设:
并行度 = 4
活跃 key 数量 = 10,000
滑动窗口:size = 1 小时,slide = 5 分钟(并发窗口数 = 12
每秒每个 key 产生 10 条数据
每条数据 1 KB

计算:
单个窗口元素数 = 60 分钟 × 60 秒 × 10/= 36,000
单个窗口大小 = 36,000 × 1 KB = 36 MB
总状态大小 = 4 × 10,000 × 12 × 36 MB / 4(分布在 4 个并行度上)
= 10,000 × 12 × 36 MB
= 4,320 GB = 4.3 TB // 非常大!

优化建议:

  • 使用增量聚合(ReduceFunction, AggregateFunction)替代全量存储
  • 减小窗口大小或增大滑动步长
  • 使用 RocksDB 状态后端(支持大状态)

7. Watermark 与窗口触发

7.1 Watermark 生成策略

/**
* 周期性 Watermark 生成器(Bounded Out Of Orderness)
*/

public class BoundedOutOfOrdernessGenerator
implements WatermarkGenerator<Event> {

private final long maxOutOfOrderness = 3500; // 最大乱序时间 3.5 秒
private long currentMaxTimestamp;

@Override
public void onEvent(Event event, long eventTimestamp, WatermarkOutput output) {
// 更新最大时间戳
currentMaxTimestamp = Math.max(currentMaxTimestamp, eventTimestamp);
}

@Override
public void onPeriodicEmit(WatermarkOutput output) {
// 周期性发射 Watermark(默认 200ms 一次)
output.emitWatermark(
new Watermark(currentMaxTimestamp maxOutOfOrderness 1)
);
}
}

/**
* 使用示例
*/

DataStream<Event> input = ...;
input
.assignTimestampsAndWatermarks(
WatermarkStrategy
.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(3))
.withTimestampAssigner((event, timestamp) -> event.getTimestamp())
)
.keyBy(Event::getUserId)
.window(TumblingEventTimeWindows.of(Time.seconds(10)))
.sum("amount");


7.2 Watermark 传播机制

/**
* Watermark 在算子间的传播(简化版源码)
*/

public class AbstractStreamOperator<OUT> {

/**
* 处理 Watermark
*/

@Override
public void processWatermark(Watermark mark) throws Exception {
// 1. 更新当前算子的 Watermark
if (timeServiceManager != null) {
timeServiceManager.advanceWatermark(mark);
}

// 2. 向下游传播 Watermark
output.emitWatermark(mark);
}
}

/**
* 多输入算子的 Watermark 对齐(如 Join 算子)
*/

public class TwoInputStreamOperator {

private long input1Watermark = Long.MIN_VALUE;
private long input2Watermark = Long.MIN_VALUE;

@Override
public void processWatermark1(Watermark mark) {
input1Watermark = mark.getTimestamp();
// 取两个输入的最小 Watermark
long minWatermark = Math.min(input1Watermark, input2Watermark);
output.emitWatermark(new Watermark(minWatermark));
}

@Override
public void processWatermark2(Watermark mark) {
input2Watermark = mark.getTimestamp();
long minWatermark = Math.min(input1Watermark, input2Watermark);
output.emitWatermark(new Watermark(minWatermark));
}
}


7.3 窗口触发时机详解

事件时间轴: 0 1 2 3 4 5 6 7 8 9 10 11 12
窗口: [ 窗口1 ] [ 窗口2 ]
start=0, end=5 start=5, end=10

数据到达:
– timestamp=1, watermark=0 → 分配到窗口1,不触发
– timestamp=3, watermark=1 → 分配到窗口1,不触发
– timestamp=4, watermark=3 → 分配到窗口1,不触发
– timestamp=6, watermark=4 → 分配到窗口2,触发窗口1(watermark=4 >= 窗口1.maxTimestamp=4)
– timestamp=9, watermark=7 → 分配到窗口2,不触发
– timestamp=11, watermark=9 → 分配到窗口3,触发窗口2(watermark=9 >= 窗口2.maxTimestamp=9)

关键公式:

窗口触发条件: Watermark >= window.maxTimestamp()
window.maxTimestamp() = window.getEnd() – 1

例如:窗口 [0, 5)
– window.maxTimestamp() = 5 – 1 = 4
– 当 Watermark >= 4 时触发


7.4 迟到数据处理

/**
* 允许延迟数据 + 侧输出流
*/

final OutputTag<Event> lateOutputTag = new OutputTag<Event>("late-data"){};

DataStream<Event> input = ...;
SingleOutputStreamOperator<Result> result = input
.keyBy(Event::getUserId)
.window(TumblingEventTimeWindows.of(Time.seconds(10)))
.allowedLateness(Time.seconds(5)) // 允许 5 秒延迟
.sideOutputLateData(lateOutputTag) // 超过延迟的数据输出到侧输出流
.sum("amount");

// 获取迟到数据流
DataStream<Event> lateStream = result.getSideOutput(lateOutputTag);
lateStream.print("LATE");

延迟容忍时间轴:

窗口: [0, 10)
窗口触发时间: Watermark >= 9

1. 正常触发: Watermark = 9 时触发窗口计算
2. 延迟容忍: Watermark = 9 到 14 之间,迟到数据仍可触发窗口重新计算
3. 彻底丢弃: Watermark >= 14 时,窗口被清理,迟到数据输出到侧输出流


8. 窗口计算流程源码分析

8.1 WindowOperator 核心方法

/**
* WindowOperator 处理元素的核心流程
*/

@Override
public void processElement(StreamRecord<IN> element) throws Exception {
// 1. 提取元素的时间戳
final long timestamp = element.getTimestamp();

// 2. 获取元素的 key
final K key = this.<K>getKeyedStateBackend().getCurrentKey();

// 3. 窗口分配:确定元素属于哪些窗口
final Collection<W> elementWindows = windowAssigner.assignWindows(
element.getValue(),
timestamp,
windowAssignerContext
);

// 4. 遍历每个窗口,将元素加入窗口状态
for (W window : elementWindows) {

// 4.1 判断窗口是否已过期
if (isWindowLate(window)) {
continue; // 跳过过期窗口
}

// 4.2 设置窗口上下文
triggerContext.window = window;

// 4.3 将元素加入窗口状态
windowState.setCurrentNamespace(window);
windowState.add(element.getValue());

// 4.4 触发器判断:是否触发计算
TriggerResult triggerResult = triggerContext.onElement(element);

// 4.5 根据触发结果执行相应操作
if (triggerResult.isFire()) {
// 触发计算
ACC contents = windowState.get();
emitWindowContents(window, contents);
}

if (triggerResult.isPurge()) {
// 清理窗口状态
windowState.clear();
}

// 4.6 注册窗口清理定时器
registerCleanupTimer(window);
}
}


8.2 窗口计算执行流程

/**
* 执行窗口函数计算并发射结果
*/

private void emitWindowContents(W window, ACC contents) throws Exception {
// 设置时间戳为窗口最大时间戳
timestampedCollector.setAbsoluteTimestamp(window.maxTimestamp());

// 执行用户定义的窗口函数
userFunction.process(
triggerContext.key, // key
window, // 窗口
processContext, // 上下文
contents, // 窗口内容(聚合结果或元素列表)
timestampedCollector // 输出收集器
);
}


8.3 窗口清理流程

/**
* 清理过期窗口的定时器回调
*/

@Override
public void onEventTime(InternalTimer<K, W> timer) throws Exception {
W window = timer.getNamespace();

// 1. 触发器判断
TriggerResult triggerResult = triggerContext.onEventTime(timer.getTimestamp());

if (triggerResult.isFire()) {
// 2. 触发计算
ACC contents = windowState.get();
if (contents != null) {
emitWindowContents(window, contents);
}
}

if (triggerResult.isPurge()) {
// 3. 清理窗口数据
windowState.clear();
}

// 4. 如果是清理定时器,清理所有状态
if (timer.getTimestamp() == cleanupTime(window)) {
clearAllState(window);
}
}

/**
* 计算窗口清理时间
*/

private long cleanupTime(W window) {
if (windowAssigner.isEventTime()) {
// 事件时间:窗口结束时间 + 允许延迟时间
long cleanupTime = window.maxTimestamp() + allowedLateness;
return cleanupTime >= window.maxTimestamp() ? cleanupTime : Long.MAX_VALUE;
} else {
// 处理时间:窗口结束时间
return window.maxTimestamp();
}
}


8.4 增量聚合 vs 全量聚合

增量聚合(ReduceFunction / AggregateFunction)

/**
* 使用 ReduceFunction 进行增量聚合
*/

input
.keyBy(Event::getUserId)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.reduce(new ReduceFunction<Event>() {
@Override
public Event reduce(Event value1, Event value2) {
// 每次新元素到达时执行聚合
return new Event(
value1.getUserId(),
value1.getAmount() + value2.getAmount()
);
}
});

状态存储:只存储聚合结果(1 条记录)

窗口状态: Event(userId=1, amount=350) // 只存储累加结果


全量聚合(ProcessWindowFunction)

/**
* 使用 ProcessWindowFunction 进行全量聚合
*/

input
.keyBy(Event::getUserId)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.process(new ProcessWindowFunction<Event, Result, String, TimeWindow>() {
@Override
public void process(
String key,
Context context,
Iterable<Event> elements,
Collector<Result> out
) {
// 窗口触发时访问所有元素
int count = 0;
double sum = 0;
for (Event event : elements) {
count++;
sum += event.getAmount();
}
out.collect(new Result(key, count, sum, sum / count));
}
});

状态存储:存储所有元素(N 条记录)

窗口状态: [Event1, Event2, Event3, …, EventN] // 存储所有数据


混合模式(增量 + 全量)

/**
* 结合 AggregateFunction 和 ProcessWindowFunction
* 既有增量聚合的性能,又有全量聚合的灵活性
*/

input
.keyBy(Event::getUserId)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.aggregate(
// 增量聚合:计算 sum 和 count
new AggregateFunction<Event, Tuple2<Double, Long>, Double>() {
@Override
public Tuple2<Double, Long> createAccumulator() {
return new Tuple2<>(0.0, 0L);
}

@Override
public Tuple2<Double, Long> add(Event value, Tuple2<Double, Long> acc) {
return new Tuple2<>(acc.f0 + value.getAmount(), acc.f1 + 1);
}

@Override
public Double getResult(Tuple2<Double, Long> acc) {
return acc.f0 / acc.f1; // 平均值
}

@Override
public Tuple2<Double, Long> merge(Tuple2<Double, Long> a, Tuple2<Double, Long> b) {
return new Tuple2<>(a.f0 + b.f0, a.f1 + b.f1);
}
},
// 全量聚合:访问窗口元数据
new ProcessWindowFunction<Double, Result, String, TimeWindow>() {
@Override
public void process(
String key,
Context context,
Iterable<Double> elements,
Collector<Result> out
) {
Double avgAmount = elements.iterator().next();
out.collect(new Result(
key,
context.window().getStart(),
context.window().getEnd(),
avgAmount
));
}
}
);

性能对比:

模式状态大小性能功能
纯增量 1 条记录 ⭐⭐⭐⭐⭐ 有限(只能做简单聚合)
纯全量 N 条记录 ⭐⭐ 强大(可访问所有数据和窗口元数据)
混合模式 1 条记录 ⭐⭐⭐⭐⭐ 强大(兼具性能和灵活性)

9. SQL 窗口实现原理

9.1 SQL 窗口函数映射

SQL 窗口函数对应的 DataStream APIWindow Assigner
TUMBLE() TumblingEventTimeWindows TumblingEventTimeWindows
HOP() SlidingEventTimeWindows SlidingEventTimeWindows
SESSION() EventTimeSessionWindows EventTimeSessionWindows
OVER() 不是窗口,是状态累加 无窗口(使用 State)

9.2 TUMBLE 窗口 SQL → DataStream 转换

SQL 代码

SELECT
user_id,
TUMBLE_START(event_time, INTERVAL '1' MINUTE) AS window_start,
TUMBLE_END(event_time, INTERVAL '1' MINUTE) AS window_end,
SUM(amount) AS total_amount
FROM orders
GROUP BY
user_id,
TUMBLE(event_time, INTERVAL '1' MINUTE);

等价的 DataStream 代码

DataStream<Order> orders = ...;

orders
.assignTimestampsAndWatermarks(...)
.keyBy(Order::getUserId)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.aggregate(
new AggregateFunction<Order, Double, Double>() {
@Override
public Double createAccumulator() { return 0.0; }

@Override
public Double add(Order value, Double acc) {
return acc + value.getAmount();
}

@Override
public Double getResult(Double acc) { return acc; }

@Override
public Double merge(Double a, Double b) { return a + b; }
},
new ProcessWindowFunction<Double, Result, String, TimeWindow>() {
@Override
public void process(
String key,
Context context,
Iterable<Double> elements,
Collector<Result> out
) {
out.collect(new Result(
key,
context.window().getStart(), // TUMBLE_START
context.window().getEnd(), // TUMBLE_END
elements.iterator().next()
));
}
}
);


9.3 OVER 窗口实现原理

重要区别:OVER 窗口不是真正的窗口,而是基于 State 的累加计算。

SQL 代码

SELECT
user_id,
event_time,
amount,
SUM(amount) OVER (
PARTITION BY user_id
ORDER BY event_time
RANGE BETWEEN INTERVAL '1' HOUR PRECEDING AND CURRENT ROW
) AS rolling_sum_1h
FROM orders;

底层实现原理

/**
* OVER 窗口底层使用 MapState 存储历史数据
*/

public class OverWindowOperator extends AbstractStreamOperator<Row> {

// 存储每个 key 的历史数据(时间戳 → 数据)
private transient MapState<Long, Double> historyState;

@Override
public void processElement(StreamRecord<Row> element) throws Exception {
Row row = element.getValue();
String userId = row.getField(0);
long eventTime = row.getField(1);
double amount = row.getField(2);

// 1. 将当前数据加入历史状态
historyState.put(eventTime, amount);

// 2. 清理过期数据(超过 1 小时)
long expireTime = eventTime 3600000; // 1 小时前
Iterator<Map.Entry<Long, Double>> iterator = historyState.iterator();
while (iterator.hasNext()) {
Map.Entry<Long, Double> entry = iterator.next();
if (entry.getKey() < expireTime) {
iterator.remove();
}
}

// 3. 计算 1 小时内的累加和
double sum = 0;
for (Map.Entry<Long, Double> entry : historyState.entries()) {
sum += entry.getValue();
}

// 4. 输出结果(每条数据都输出)
output.collect(new Row(userId, eventTime, amount, sum));
}
}

OVER 窗口 vs 滚动窗口的关键区别:

特性OVER 窗口滚动窗口(TUMBLE)
输出时机 每条数据到达时输出 窗口结束时输出
状态存储 存储历史数据(MapState) 存储窗口内数据(WindowState)
状态清理 依赖 State TTL 或手动清理 Watermark 触发自动清理
数据重复 每条数据都输出(流式输出) 每个窗口只输出一次

9.4 SQL 窗口聚合优化

— 低效写法:使用 ProcessWindowFunction(全量聚合)
SELECT
user_id,
window_start,
window_end,
COUNT(*) AS cnt,
SUM(amount) AS total
FROM TABLE(
TUMBLE(TABLE orders, DESCRIPTOR(event_time), INTERVAL '1' MINUTE)
)
GROUP BY user_id, window_start, window_end;

— 高效写法:Flink 自动使用增量聚合
— 无需修改 SQL,Planner 会自动优化为 ReduceFunction

Flink SQL 优化器自动优化规则:

  • SUM, COUNT, MIN, MAX, AVG → 增量聚合
  • COLLECT, LISTAGG, 复杂 UDF → 全量聚合

10. 性能优化最佳实践

10.1 选择合适的窗口类型

场景推荐窗口类型原因
固定时间段统计(如每分钟统计) 滚动窗口 无重叠,状态最小
平滑统计(如移动平均) 滑动窗口 提供连续视图
用户会话分析 会话窗口 动态适应活动模式
实时大屏(仪表盘) 滚动窗口 + 短周期 低延迟
需要所有历史数据 OVER 窗口或 State 但需注意状态膨胀

10.2 减小窗口状态大小

❌ 避免:全量存储 + 大窗口

// 不推荐:1 小时窗口 + 全量存储
input
.keyBy(...)
.window(TumblingEventTimeWindows.of(Time.hours(1)))
.process(new ProcessWindowFunction<Event, Result, String, TimeWindow>() {
@Override
public void process(...) {
// 存储 1 小时内所有数据
}
});

✅ 推荐:增量聚合 + 合理窗口大小

// 推荐:增量聚合 + 适当窗口大小
input
.keyBy(...)
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.aggregate(new MyAggregateFunction());


10.3 滑动窗口优化

❌ 避免:小步长导致大量重叠

// 不推荐:窗口 1 小时,步长 1 分钟 → 60 个并发窗口
input
.window(SlidingEventTimeWindows.of(
Time.hours(1), // 窗口大小
Time.minutes(1) // 步长 → 60 倍状态!
));

✅ 推荐:调整步长或改用其他方案

// 方案 1:增大步长
input
.window(SlidingEventTimeWindows.of(
Time.hours(1),
Time.minutes(10) // 步长改为 10 分钟 → 只有 6 个并发窗口
));

// 方案 2:改用滚动窗口 + 后处理
input
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.aggregate(...)
.keyBy(...)
.flatMap(new SlidingAggregator()); // 在内存中计算滑动聚合


10.4 State TTL 配置

/**
* 为窗口状态配置 TTL(防止状态无限增长)
*/

StreamExecutionEnvironment env = ...;
env.getConfig().enableObjectReuse(); // 开启对象重用

// 配置全局 State TTL
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.hours(24))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.build();

// 注意:窗口状态的 TTL 由窗口自身管理,无需手动配置


10.5 早期触发(Early Firing)

/**
* 需要实时性:每 5 秒输出中间结果
*/

input
.keyBy(...)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.trigger(ContinuousEventTimeTrigger.of(Time.seconds(5)))
.process(...);

适用场景:

  • 实时大屏(Dashboard)
  • 需要看到进度的长时间窗口

10.6 Checkpoint 优化

/**
* 窗口作业的 Checkpoint 配置
*/

StreamExecutionEnvironment env = ...;

// 1. 启用增量 Checkpoint(RocksDB)
env.setStateBackend(new EmbeddedRocksDBStateBackend(true));

// 2. 设置合理的 Checkpoint 间隔
env.enableCheckpointing(60000); // 1 分钟

// 3. 配置 Checkpoint 超时
env.getCheckpointConfig().setCheckpointTimeout(300000); // 5 分钟

// 4. 允许并发 Checkpoint
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);

// 5. 配置 Checkpoint 失败容忍次数
env.getCheckpointConfig().setTolerableCheckpointFailureNumber(3);


10.7 水位线对齐优化

/**
* 禁用水位线对齐(提高吞吐,但可能丢失数据)
*/

env.getCheckpointConfig()
.enableUnalignedCheckpoints(); // 非对齐 Checkpoint

权衡:

  • ✅ 提高吞吐量(减少反压)
  • ❌ 可能导致 Checkpoint 包含乱序数据

11. 窗口性能监控

11.1 关键指标

# 1. 窗口状态大小
flink.taskmanager.<tm_id>.<job_name>.operators.<operator_name>.state_size

# 2. 窗口触发延迟
flink.taskmanager.<tm_id>.<job_name>.operators.<operator_name>.window_trigger_latency

# 3. 窗口计算耗时
flink.taskmanager.<tm_id>.<job_name>.operators.<operator_name>.window_computation_time

# 4. 迟到数据数量
flink.taskmanager.<tm_id>.<job_name>.operators.<operator_name>.numLateRecordsDropped


11.2 问题排查

现象可能原因解决方案
窗口不触发 Watermark 未推进 检查 Watermark 生成器
状态过大 窗口范围太大或使用全量存储 使用增量聚合
Checkpoint 超时 状态过大 增大 Checkpoint 超时时间、使用 RocksDB
大量迟到数据 乱序时间设置过小 增大 maxOutOfOrderness
内存溢出 滑动窗口步长过小 增大步长或改用滚动窗口

12. 总结

12.1 窗口设计核心要点

  • 窗口分配器(Window Assigner):决定数据分配到哪个窗口
  • 触发器(Trigger):决定何时触发窗口计算
  • 驱逐器(Evictor):决定哪些数据从窗口移除(可选)
  • 窗口函数(Window Function):定义窗口计算逻辑
  • 12.2 性能优化金字塔

    增量聚合(最优)
    / \\
    合理窗口大小 State TTL
    / | | \\
    RocksDB 早期触发 延迟容忍 监控告警

    12.3 技术选型决策树

    是否需要全部历史数据?
    ├─ 是 → 使用 OVER 窗口或 State(注意配置 TTL)
    └─ 否 → 是否需要窗口元数据(如窗口起止时间)?
    ├─ 是 → 使用 aggregate + ProcessWindowFunction(混合模式)
    └─ 否 → 使用纯增量聚合(ReduceFunction / AggregateFunction)


    赞(0)
    未经允许不得转载:171主机测评 » 一文搞懂Flink 窗口设计深入详解(含源码分析)
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址