在 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 不动
  • 窗口不触发
  • 数据“卡住”
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐