Flink Watermark 卡住的真正原因:keyBy 根本不是“时间分区”
在 Flink 事件时间(Event Time)开发中,
很多人都会遇到一个极其诡异的问题:👉 明明某些 key 的数据还在源源不断地产生,窗口却迟迟不触发,watermark 卡住不动。
归根结底,问题不在 API,而在于:
没有把 Flink 里的“两种分区”分清楚。
一、先给一句结论
Flink 里至少同时存在两种“分区”:
① Source 并行分区(subtask / 物理分区)
② keyBy 之后的逻辑分区(key)👉 Watermark 只和 ① 有关,和 ② 没关系。
只要这句话没理解透,Flink 事件时间几乎一定会踩坑。
二、把“分区”这个词拆成两种
① Source 分区(物理分区 / subtask)【Watermark 关心的】
这是 Flink 运行时层面的分区,也是 Watermark 的最小作用单位。
Source 并行度 = 3
Subtask-0
Subtask-1
Subtask-2
它通常对应:
- Kafka 的 partition
- RabbitMQ / MQ 的 consumer 分片
- DDS Reader
- Flink Source 的并行子任务
👉 一个 Source subtask = 一条独立的 Watermark 时间线
② keyBy 分区(逻辑分区 / key)【Watermark 不关心的】
这是 业务层面的分区,由你在代码中定义:
.keyBy(UavState::getUavId)
逻辑上你看到的是:
uav-001
uav-002
uav-003
uav-004
key 的作用是:
- 决定 state 的隔离
- 决定 同一 UAV 的数据会进同一个算子实例
- 方便你写“每个实体一条轨迹”的业务逻辑
但要特别注意:
- ❌ key 不产生 watermark
- ❌ key 不推进时间
三、结构图
Kafka / DDS / MQ
│
▼
Source 并行子任务(subtask)
├── Subtask-0 ← watermark-0
│ ├── uav-001
│ ├── uav-002
│ └── uav-003
│
├── Subtask-1 ← watermark-1
│ ├── uav-004
│ ├── uav-005
│ └── uav-006
│
└── Subtask-2 ← watermark-2
├── uav-007
└── uav-008
关键认知点:
- 一个 subtask 里 可以同时处理很多 key
- watermark 是 subtask 级别的
- key 只是“住在 subtask 里的住户”
四、彻底讲清“同一分区 / 不同分区”
情况 1️⃣ 同一个 key + 同一个 subtask(最简单)
uav-001 → subtask-0
- 有数据 → subtask-0 的 watermark 推进
- 没数据 → subtask-0 可能 idle
✅ 很好理解
情况 2️⃣ 不同 key,但在 同一个 subtask(重点)
uav-001 → subtask-0
uav-002 → subtask-0
uav-002 掉线了
- subtask-0 仍然在接收 uav-001 的数据
- subtask-0 不是 idle
- watermark 继续推进
👉 某个 key 掉线 ≠ watermark 一定会卡住
这是很多人第一次理解错误的地方。
情况 3️⃣ 不同 key,在 不同 subtask(最容易卡死)
uav-001 → subtask-0
uav-002 → subtask-1
uav-002 掉线了
- subtask-1 长时间没数据
- subtask-1 的 watermark 停在旧时间
- 下游 watermark = min(subtask-0, subtask-1)
- 被 subtask-1 拖死
👉 这时候:
- 即使 uav-001 数据正常
- 窗口也可能一直不触发
情况 4️⃣ 加了 withIdleness 之后
.withIdleness(Duration.ofSeconds(10))
subtask-1 连续 10 秒没数据
- 被 Flink 标记为 idle
- 不再参与 watermark 最小值计算
结果是:
全局 watermark = watermark(subtask-0)
👉 时间继续向前推进,窗口正常触发。
五、一个非常重要的“公式”
Watermark 会被“最慢的 subtask”拖住,
而不是被“最慢的 key”拖住。
六、窗口卡住时,自检的 3 个问题
当你发现窗口迟迟不触发时,直接问自己:
1️⃣ 掉线的是 key 还是 整个 subtask?
2️⃣ 掉线的 key 是否是该 subtask 唯一活跃的 key?
3️⃣ 是否配置了 withIdleness?
判断结论:
- 掉的是 key,但 subtask 还有其他 key → 不会卡
- 掉的是整个 subtask → 一定会卡(除非 withIdleness)
七、总结
keyBy 是“业务分区”,
subtask 是“时间分区”,
watermark 只认后者。
很多 Flink 事件时间的问题,看起来是:
- watermark 不动
- 窗口不触发
- 数据“卡住”
更多推荐

所有评论(0)