The source notes for this track are currently maintained in Chinese.
预计阅读时间: 45 分钟 前置阅读: doc-05(Window & Timer) 下一次阅读: doc-08(反压)
1. 三种时间语义
| 语义 | 来源 | 确定性 | 延迟 | 适用 |
|---|---|---|---|---|
| Event Time | 数据字段中的时间戳 | 确定 (相同数据→相同结果) | 可能高 (等 Watermark) | 准确分析 |
| Processing Time | 算子执行时的系统时钟 | 不确定 (重启后结果不同) | 极低 | 大致统计 |
| Ingestion Time | 数据进入 Flink Source 的时间 | 部分确定 | 低 | 折中方案, 几乎不用 |
2. Watermark 的生成与语义
2.1 WatermarkGenerator
// WatermarkStrategy.java — 两种生成器模式
// 模式 1: Periodic (周期性)
new BoundedOutOfOrdernessWatermarks<>(Duration.ofSeconds(5));
// 每次 onPeriodicEmit():
// Watermark = maxTimestamp - 5s - 1ms
// 模式 2: Punctuated (逐事件)
new PunctuatedAssigner() {
public Watermark checkAndGetNextWatermark(Event event, long ts) {
// 某些特殊事件标记 Watermark (如 "END_OF_BATCH" 事件)
return event.isEndOfBatch() ? new Watermark(ts) : null;
}
};
2.2 BoundedOutOfOrderness 生成器详解
// BoundedOutOfOrdernessWatermarks.java
public void onEvent(T event, long eventTimestamp, WatermarkOutput output) {
// 跟踪最大时间戳
maxTimestamp = Math.max(maxTimestamp, eventTimestamp);
}
public void onPeriodicEmit(WatermarkOutput output) {
// 周期性输出 Watermark = maxTimestamp - outOfOrdernessMillis - 1
// -1 是为了保证 "Watermark > T" 意味着 "所有时间 ≤ T" 的数据都已到达
// (等于号的特殊处理)
output.emitWatermark(new Watermark(maxTimestamp - outOfOrdernessMillis - 1));
}
3. Watermark 传播与合并
3.1 多输入算子的 Watermark 合并

源码实现:
// StreamInputProcessor.java — Watermark 合并
private void processWatermark(Watermark watermark, int channelIndex) {
// 1. 更新该 Channel 的 Watermark
channelWatermarks[channelIndex] = watermark.getTimestamp();
// 2. 取所有 Channel 的最小值
long newMinWatermark = min(channelWatermarks);
// 3. 只有当新 Watermark > 当前 Watermark 时才推进
if (newMinWatermark > currentWatermark) {
currentWatermark = newMinWatermark;
output.emitWatermark(new Watermark(currentWatermark));
}
}
3.2 Idleness 处理
// StreamInputProcessor — Idle 检测
WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withIdleness(Duration.ofMinutes(1));
// 源码: 如果某个 Channel 超过 1 分钟无数据 → 标记为 IDLE
// IDLE Channel 不参与 min() 计算
// 例如: Channel-1 WM=100, Channel-2 IDLE
// → 合并后 WM = 100 (不是 min(100, IDLE))
4. Watermark 与 Window 的精确关系
Window [10:00, 10:01):
- 接收事件: 事件时间在 [10:00, 10:01) 的文件加入此窗口
- Trigger: EventTimeTrigger → Watermark ≥ 10:01 → 触发
- allowedLateness (如 30s): Watermark < 10:01:30 时延迟到达的元素
仍可加入窗口并触发重新计算
- Cleanup: Watermark ≥ 10:01:30 → window state + timer 被清除
5. 常见问题 / 面试题
Q1: 为什么多输入 Watermark = min(all inputs)?
A: 语义正确性。假设 Source-1 (最快) WM=100, Source-2 (最慢) WM=98。如果取 max=100,则 Source-2 的 Channel 仍有时间戳 ≤ 98 的数据未到达 → 窗口 [10:00,10:01) 被错误地触发 → 后续 Source-2 到达的时间戳 ≤ 98 的数据成为迟到数据(可能丢失)。取 min 保证所有上游都确认 "≤ T 的数据已到达"。
Q2: Idle Source 如何导致 Watermark 停滞?
A: Source Partition 无数据 → 无新元素 → maxTimestamp 不更新 → onPeriodicEmit() 输出的 Watermark 不推进 → 下游的 min(all Watermarks) 停滞 → 窗口永不触发 → 数据堆积。withIdleness() 将空闲 Channel 排除在 min 计算之外。
Q3: Watermark 和 Checkpoint Barrier 有什么区别?
A: 虽然两者都是流中注入的特殊事件,但语义完全不同。Watermark 是时间进度信号("≤ T 的数据到了"),Barrier 是快照触发信号("请做 Checkpoint")。Watermark 触发 Window 计算,Barrier 触发状态快照。两者相互独立但都在同一个 Channel 中顺序传输。