Flink 源码与运行原理

流式计算 / 源码阅读 / LESSON 06

Watermark 与事件时间深度解析

掌握 Watermark 如何表达乱序边界,以及事件时间计算何时可以安全推进。

阅读时间
45 分钟
学习路径
Flink 源码与运行原理
内容来源
Flink 深度笔记

预计阅读时间: 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 合并

Watermark 与事件时间深度解析 图 01

源码实现

// 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 中顺序传输。


下一步