Flink 时间属性Event Time、Processing Time 与高级 Watermark 实战
1. 三个核心概念:时间、时间属性、Watermark
1.1 两种时间概念
Flink 里最常用的是这两个:
-
Processing Time(处理时间)
用的是算子所在机器的系统时间(epoch time),例如System.currentTimeMillis()。简单粗暴,不需要解析时间戳、不需要 Watermark,但结果不可复现、也无法精确处理乱序和迟到。 -
Event Time(事件时间)
基于每条数据自带的时间戳来计算,表示“事件真正发生的时间”。
优势:- 能在乱序、迟到场景下得到一致的结果;
- 从 Kafka / 对象存储回放时,结果可重现;
- 批流一个 SQL,语义一致(批里就是普通时间列)。
1.2 什么是“时间属性”
在 Table / SQL 里,我们会把某一列声明为时间属性。它本质上就是一个带语义的时间戳列:
- 在 schema 中被标记为 Event Time 或 Processing 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)
步骤:
- 在 DataStream 上用
assignTimestampsAndWatermarks(...)提前完成时间戳与 Watermark 分配; - 在
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. 实战经验与易踩坑
最后总结一组在项目中非常实用的结论:
-
优先 Event Time,其次才是 Processing Time
只要事件时间能拿到,尽量按 Event Time 建模,这样回放、补数、批流一体都好处理。 -
TIMESTAMP vs TIMESTAMP_LTZ 不要乱用
- 文本时间戳 ⇒
TIMESTAMP; - Epoch long ⇒
TIMESTAMP_LTZ + 计算列。
- 文本时间戳 ⇒
-
记得设置 Watermark 延迟和 idle-timeout
- 延迟要根据业务的乱序程度来定;
- 多分区 Source 几乎必配
scan.watermark.idle-timeout,避免单分区“卡进度”。
-
不要指望普通时间戳自动变成时间属性
必须在 DDL(WATERMARK、PROCTIME())或 DataStream 转表时显式声明。 -
注意时间属性被“用坏”的场景
event_time + INTERVAL '1' HOUR之后再拿去做窗口,就已经是普通 TIMESTAMP 了,
此时需要单独保留原始时间属性列。 -
Watermark 对齐不是银弹
只在多源/多分片消费速率差异明显、且窗口结果对顺序敏感时开启,
注意评估对吞吐和延迟的影响。
7. 小结
一句话把本文内容收拢起来:
“在 Flink 里,时间属性是带语义的时间戳;
Event Time + Watermark 保证精确与可重放,Processing Time 换来简单与高实时性;
1.18 开始,高级 Watermark 让 SQL 也能玩转发射策略、空闲超时和对齐。”
真正落地时,你只需要先回答三个问题:
- 我关心的是“事件发生时间”还是“处理到达时间”?
- 源数据的时间戳是什么格式(字符串 / epoch)?
- 我的 Source 分区多不多,速率差异大不大?
更多推荐
所有评论(0)