本节摘要:流处理有三只钟:事件时间看数据自己带的出生证明,处理时间看处理机器的挂钟,摄入时间折中二者。本节讲清三只钟各自的成立前提与失效场景,给出"选错钟"的经典翻车案例,并演示在代码里声明时间语义与时间戳的正确姿势。
第 2 章末尾我们说数据是一笔糊涂账,糊涂就糊涂在这里。一条日志在客户端生成时是十点整,经过采集、转发、入队,抵达 Flink 时已经十点过四分——用哪个时刻参与计算,得出的完全是两个世界的结果。处理时间派说"以我收到为准",简单可靠;事件时间派说"以发生为准",严谨但复杂。这一岔路是流计算的第一道分水岭,本节就把它讲透。
事件时间(Event Time):数据在业务世界发生的时刻,由产生端打进数据本身,是一条自带出生证明的时间。它的好处是结果可复现——无论数据何时到达、以什么顺序到达、甚至隔天回放历史数据,窗口聚合结果都一致。它的代价是引擎必须处理"时钟错乱":数据乱序、迟到、甚至时间戳本身缺失或错误,都需要一整套机制去兜底(水位线机制,下一节的主角)。
处理时间(Processing Time):执行算子的机器挂钟。它的好处是极简——没有乱序问题,数据到了就算,延迟最低;坏处是结果不可复现:同一份数据重跑一遍,因为到达时刻不同,结果就不同。它适合"只关心此刻、不必对账"的场景,比如内存监控告警。
摄入时间(Ingestion Time):数据进入 Flink 源算子那一刻的时间戳。它是折中方案:比事件时间省心(到源就定了,后续不再变),比处理时间多一层稳定(同一数据源内的先后顺序固定)。但在事件发生到摄入之间发生的乱序它无能为力,实际生产用得少,多数团队直接在"事件时间或处理时间"二选一。
| 维度 | 事件时间 | 处理时间 | 摄入时间 |
|---|---|---|---|
| 时间来源 | 数据自带时间戳 | 算子机器挂钟 | 进入源算子的时刻 |
| 可复现性 | 完全可复现 | 不可复现 | 重放源时基本可复现 |
| 乱序处理 | 需要(水位线兜底) | 天然免疫 | 部分免疫 |
| 延迟开销 | 等待乱序、水位线推进 | 最低 | 较低 |
| 典型场景 | 对账、计费、离线一致口径 | 实时告警、监控 | 中间态分析 |
翻车一:用处理时间算大屏,回补数据全部失踪。 某团队的大屏作业用处理时间窗口统计每分钟成交额。运营要求把昨天漏采的订单补进去,结果大屏纹丝不动——处理时间的窗口只认"此刻到达"的数据,昨天的数据到达时落进的是"补数那一刻"的窗口,早被滚过去了,而历史窗口早已触发完毕。换事件时间后,补的数据带着昨天的时间戳,落进昨天的窗口重新触发,问题消失。凡是要回放、要对账的口径,必须事件时间。
翻车二:用事件时间却从不推进水位线,窗口永远不触发。 另一团队按第一条经验换了事件时间,结果窗口一个都不出数。原因是源算子忘了分配时间戳与生成水位线,引擎的时间"冻"在起点。这个案例说明:事件时间不是免费的,它要求你显式回答"数据乱序到什么程度",这正是下一节水位线要回答的问题。
事件时间的正确打开方式分两步:第一步,源或转换处为每条记录分配时间戳;第二步,声明水位线生成策略。下面的代码把"每分钟成交额"作业的骨架立了起来:
KafkaSource<String> source = KafkaSource.<String>builder() .setBootstrapServers("broker1:9092") .setTopics("payments") .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); // WatermarkStrategy 声明:时间戳从记录的 payTime 字段取,容忍五分钟乱序 WatermarkStrategy<Payment> strategy = WatermarkStrategy .<Payment>forBoundedOutOfOrderness(Duration.ofMinutes(5)) .withTimestampAssigner((payment, ts) -> payment.getPayTimeMillis()); DataStream<Payment> payments = env .fromSource(source, strategy, "payments") .map(json -> JSON.parseObject(json, Payment.class)); payments.keyBy(Payment::getShopId) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .sum("amount");
注意三点:时间戳单位是毫秒;容忍五分钟乱序意味着窗口至少要等多五分钟才触发;处理时间语义则简单得多,把窗口类型换成处理时间变体、水位线策略传无即可,但请先想清楚上一节翻车案例里的代价。
💡 关键直觉:三只钟没有绝对优劣,只有场景匹配。判断口诀——要复现和对账,用事件时间;只要此刻的反应速度,用处理时间;两头都沾的,先想清楚业务真正要什么再选。
除了两次经典翻车,还有一类更隐蔽的病灶值得提前点名:同一作业内混用两只钟。某团队的风控作业用事件时间做支付聚合,却用处理时间做定时器回调里的特征刷新——平时相安无事,回溯历史数据时特征刷新的时间轴与聚合的时间轴彻底脱节,算出的"当时特征"实际上是"回放时刻的特征"。这类混用不会报错,只会让结果在特定场景悄悄变味,是时间类工单里最难缠的一支。诊断要诀:画出作业里所有"与时间有关的行为"(窗口、定时器、TTL、水印),逐一标注用的是哪只钟——一张作业里的钟越多,口径漂移的面就越大,能统一就统一。
选定了时钟,下一个问题是:引擎凭什么判断"该等的数据都到齐了"?这就进入全册最重要的机制之一——水位线。