1. 三个核心概念:时间、时间属性、Watermark

1.1 两种时间概念

Flink 里最常用的是这两个:

  • Processing Time(处理时间)
    用的是算子所在机器的系统时间(epoch time),例如 System.currentTimeMillis()。简单粗暴,不需要解析时间戳、不需要 Watermark,但结果不可复现、也无法精确处理乱序和迟到。

  • Event Time(事件时间)
    基于每条数据自带的时间戳来计算,表示“事件真正发生的时间”。
    优势:

    • 能在乱序、迟到场景下得到一致的结果;
    • 从 Kafka / 对象存储回放时,结果可重现;
    • 批流一个 SQL,语义一致(批里就是普通时间列)。

1.2 什么是“时间属性”

在 Table / SQL 里,我们会把某一列声明为时间属性。它本质上就是一个带语义的时间戳列:

  • 在 schema 中被标记为 Event TimeProcessing Time
  • 可以被窗口、Interval Join 等时间相关算子直接使用;
  • 只要不参与运算、只是透传,它一直是“时间属性”;
  • 一旦用于算术计算(比如加减),就会被物化成普通时间戳,之后就不再是时间属性了。

注意:普通 TIMESTAMP/TIMESTAMP_LTZ 不能直接“升级”为时间属性,必须在 DDL 或 DataStream→Table 时显式声明。

1.3 Watermark:事件时间的“进度指示器”

在 Event Time 场景下,为了分辨“准时”与“迟到”事件,Flink 需要知道事件时间已经推进到哪一刻
这个“进度”就是 Watermark

  • Watermark = “我相信时间戳 ≤ W 的数据基本到齐了”;
  • 窗口、Interval Join 等算子会根据 Watermark 决定何时出结果、何时关窗。

后面的内容都围绕“时间属性 + Watermark”展开。

2. Event Time:SQL 中的标准写法

2.1 时间字段类型怎么选?

Flink 支持在两种列类型上定义 Event Time 属性:

  • TIMESTAMP§:适合源数据是“年月日时分秒”的无时区字符串,例如
    2020-04-15 20:13:40.564
  • TIMESTAMP_LTZ§:适合源数据是 epoch 毫秒/秒的 long,或需要带时区语义时。

2.2 典型 DDL 模板

场景一:源数据是字符串时间戳

CREATE TABLE user_actions (
  user_name        STRING,
  data             STRING,
  user_action_time TIMESTAMP(3),
  -- 声明 user_action_time 为事件时间,并使用 5 秒延迟的 Watermark 策略
  WATERMARK FOR user_action_time AS user_action_time - INTERVAL '5' SECOND
) WITH (
  'connector' = 'kafka',
  'topic' = 'user_actions',
  'properties.bootstrap.servers' = 'localhost:9092',
  'format' = 'json'
);

-- 每 10 分钟一个滚动窗口统计活跃用户数
SELECT
  TUMBLE_START(user_action_time, INTERVAL '10' MINUTE) AS win_start,
  COUNT(DISTINCT user_name)                           AS uv
FROM user_actions
GROUP BY TUMBLE(user_action_time, INTERVAL '10' MINUTE);

场景二:源数据是 epoch 毫秒

CREATE TABLE user_actions (
  user_name STRING,
  data      STRING,
  ts        BIGINT,
  -- 计算列:把毫秒时间戳转成 TIMESTAMP_LTZ(3)
  time_ltz  AS TO_TIMESTAMP_LTZ(ts, 3),
  WATERMARK FOR time_ltz AS time_ltz - INTERVAL '5' SECOND
) WITH (...);

SELECT
  TUMBLE_START(time_ltz, INTERVAL '10' MINUTE) AS win_start,
  COUNT(DISTINCT user_name)                    AS uv
FROM user_actions
GROUP BY TUMBLE(time_ltz, INTERVAL '10' MINUTE);

建议规则可以简单记成:

字符串时间戳 ⇒ TIMESTAMP + WATERMARK
Epoch 毫秒 ⇒ TIMESTAMP_LTZ + 计算列 + WATERMARK

3. 1.18+ 高级 Watermark 能做什么?

在 DataStream API 里我们可以通过 WatermarkGenerator 玩出很多花样,而之前 SQL 较弱。
Flink 1.18 开始,SQL 也支持了三项关键能力:

前提:Source Connector 实现了 SupportsWatermarkPushDown(如 Kafka / Pulsar);否则参数不会生效。

3.1 发射策略:按周期还是按事件?

Flink 支持两类策略:

  • on-periodic:周期性发射 Watermark

    • SQL 默认策略;
    • 周期由 pipeline.auto-watermark-interval 控制,默认 200ms。
  • on-event每条事件发射一次 Watermark

    • 延迟更低,但对吞吐有影响。

在表级别配置示例:

CREATE TABLE user_actions (
  ...
  user_action_time TIMESTAMP(3),
  WATERMARK FOR user_action_time AS user_action_time - INTERVAL '5' SECOND
) WITH (
  'scan.watermark.emit.strategy' = 'on-event',
  ...
);

也可以在查询里用 OPTIONS hint 临时覆盖:

SELECT ...
FROM user_actions /*+ OPTIONS('scan.watermark.emit.strategy'='on-periodic') */;

3.2 idle-timeout:解决“某个分区停了拖累全局”的问题

Watermark 是取所有上游分区 Watermark 的最小值,这就会出现一个经典问题:

某个分区长时间没数据,新 Watermark 不往前推 → 下游窗口永远不开/不关。

解决办法:配置 空闲超时,让“沉默分区”在超时后被标记为 idle,下游计算 Watermark 时忽略它。

  • 全局配置:table.exec.source.idle-timeout(作用于所有 Source)。

  • 表级配置(优先级更高):

    CREATE TABLE user_actions (
      ...
      user_action_time TIMESTAMP(3),
      WATERMARK FOR user_action_time AS user_action_time - INTERVAL '5' SECOND
    ) WITH (
      'scan.watermark.idle-timeout'='1min',
      ...
    );
    

或者:

SELECT ...
FROM user_actions /*+ OPTIONS('scan.watermark.idle-timeout'='1min') */;

3.3 水位线对齐(Watermark Alignment)

问题场景:
不同分区 / 不同 Source 的消费速率差别很大;
快分区不断往前推进 Watermark,慢分区却很落后,导致:

  • 下游为了等慢分区,状态缓存越积越多;
  • 整体乱序变严重,窗口计算精度受影响。

Watermark Alignment 的目标是:

让“快分区”的 Watermark 不要比“慢分区”快太多,在配置的 max-drift 范围内对齐推进。

配置示例:

CREATE TABLE user_actions (
  ...
  user_action_time TIMESTAMP(3),
  WATERMARK FOR user_action_time AS user_action_time - INTERVAL '5' SECOND
) WITH (
  'scan.watermark.alignment.group'          = 'alignment-group-1',
  'scan.watermark.alignment.max-drift'      = '1min',
  'scan.watermark.alignment.update-interval'='1s',
  ...
);

参数说明:

  • scan.watermark.alignment.group:对齐组名,同组内的数据源一起对齐;
  • scan.watermark.alignment.max-drift:允许的最大偏移(快分区不能比对齐时间快太多);
  • scan.watermark.alignment.update-interval:多久重新计算一次对齐时间(默认 1s)。

注意:

  • 需要连接器落实 FLIP-217 对分片对齐的支持;
  • 否则可以配置 pipeline.watermark-alignment.allow-unaligned-source-splits: true 关闭分片级对齐(只在“分片数 = Source 并行度”时才工作正常)。

4. DataStream ↔ Table:如何带上时间属性

很多项目会同时用 DataStream 与 Table API,需要在两者之间互转,这时时间属性不能丢

4.1 DataStream → Table(Event Time)

步骤:

  1. 在 DataStream 上用 assignTimestampsAndWatermarks(...) 提前完成时间戳与 Watermark 分配;
  2. fromDataStream 的 schema 中用 .rowtime() 指明哪一列是事件时间属性。

Flink 会把 rowtime 推导成 TIMESTAMP WITHOUT TIME ZONE,并认为它是 UTC 时间(因为 DataStream 没有时区概念)。

Java 示例:

// Option 1:追加一列作为 rowtime
DataStream<Tuple2<String, String>> stream =
    inputStream.assignTimestampsAndWatermarks(...);

Table table = tEnv.fromDataStream(
    stream,
    $("user_name"),
    $("data"),
    $("user_action_time").rowtime()   // 新增 rowtime 列
);

// Option 2:用第一列做时间戳,转成 rowtime 替换掉
DataStream<Tuple3<Long, String, String>> stream2 =
    inputStream.assignTimestampsAndWatermarks(...);

Table table2 = tEnv.fromDataStream(
    stream2,
    $("user_action_time").rowtime(),  // 用 rowtime 覆盖第 1 列
    $("user_name"),
    $("data")
);

// 使用方式:直接做时间窗口
WindowedTable windowedTable = table2.window(
    Tumble.over(lit(10).minutes())
          .on($("user_action_time"))
          .as("userActionWindow"));

4.2 DataStream → Table(Processing Time)

更简单:在 schema 最后追加 .proctime() 字段即可:

DataStream<Tuple2<String, String>> stream = ...;

Table table = tEnv.fromDataStream(
    stream,
    $("user_name"),
    $("data"),
    $("user_action_time").proctime()   // 只能放在最后
);

5. Processing Time:什么时候可以偷懒?

Processing Time 是 Flink 里最省心的时间概念:

  • 不用解析业务时间戳;
  • 不需要 Watermark;
  • 只依赖“系统现在几点”。

代价也很明显:结果随机器时间、重启时机而变化,不可重放
适用场景:

  • 自监控类统计(比如“过去 10 分钟处理了多少条日志”);
  • 对绝对精度没那么敏感、只要“实时趋势”的业务;
  • 源头没有可靠事件时间的场景。

DDL 示例:

CREATE TABLE user_actions (
  user_name STRING,
  data STRING,
  user_action_time AS PROCTIME()
) WITH (...);

SELECT
  TUMBLE_START(user_action_time, INTERVAL '10' MINUTE) AS win_start,
  COUNT(DISTINCT user_name)                            AS uv
FROM user_actions
GROUP BY TUMBLE(user_action_time, INTERVAL '10' MINUTE);

6. 实战经验与易踩坑

最后总结一组在项目中非常实用的结论:

  1. 优先 Event Time,其次才是 Processing Time
    只要事件时间能拿到,尽量按 Event Time 建模,这样回放、补数、批流一体都好处理。

  2. TIMESTAMP vs TIMESTAMP_LTZ 不要乱用

    • 文本时间戳 ⇒ TIMESTAMP
    • Epoch long ⇒ TIMESTAMP_LTZ + 计算列
  3. 记得设置 Watermark 延迟和 idle-timeout

    • 延迟要根据业务的乱序程度来定;
    • 多分区 Source 几乎必配 scan.watermark.idle-timeout,避免单分区“卡进度”。
  4. 不要指望普通时间戳自动变成时间属性
    必须在 DDL(WATERMARKPROCTIME())或 DataStream 转表时显式声明。

  5. 注意时间属性被“用坏”的场景
    event_time + INTERVAL '1' HOUR 之后再拿去做窗口,就已经是普通 TIMESTAMP 了,
    此时需要单独保留原始时间属性列。

  6. Watermark 对齐不是银弹
    只在多源/多分片消费速率差异明显、且窗口结果对顺序敏感时开启,
    注意评估对吞吐和延迟的影响。

7. 小结

一句话把本文内容收拢起来:

“在 Flink 里,时间属性是带语义的时间戳;
Event Time + Watermark 保证精确与可重放,Processing Time 换来简单与高实时性;
1.18 开始,高级 Watermark 让 SQL 也能玩转发射策略、空闲超时和对齐。”

真正落地时,你只需要先回答三个问题:

  1. 我关心的是“事件发生时间”还是“处理到达时间”?
  2. 源数据的时间戳是什么格式(字符串 / epoch)?
  3. 我的 Source 分区多不多,速率差异大不大?

更多推荐