做实时计算的人,迟早都会撞上一个问题:数据明明是按时发生的,到你手里却乱得不成样子。比如业务系统记录了事件发生时间,但经过网络、消息队列、重试机制,到达 Flink 时先后顺序早就打乱了。如果这个时候直接按窗口聚合,统计结果必然失真。Flink 里专门解决这个问题的机制就是 Watermark,它不是业务字段,而是一种告诉引擎“哪些数据已经到齐了”的信号。这篇文章适合正在写 Flink 作业、遇到乱序数据不知道如何取舍的读者,我会从原理讲到实战,把 Watermark 怎么设计、怎么调参、怎么排查一次说清楚。
1. 为什么需要 Watermark:乱序数据背后的三个真实困境
1.1 事件时间与处理时间:别让统计口径错位
先分清两个最基础的概念:处理时间(Processing Time)是指数据到达 Flink 机器的时间,也就是系统当前时间;事件时间(Event Time)是数据里携带的业务发生时间,是真实世界里事件产生的那一刻。
举个例子。用户在一分钟窗口的最后 1 秒点击了页面,这条点击日志的事件时间是 12:00:59。但因为网络抖动,数据被 Kafka 缓冲了一会儿,直到 12:01:05 才到达 Flink。如果按处理时间开窗口,这条点击会被划入 12:01 那一分钟,统计结果就会整体偏移。更麻烦的是,如果 12:00 的窗口已经关闭,这条数据要么被丢弃,要么被错误地记到下一分钟。
处理时间最大的问题就是“统计口径不可复现”。同一批数据,重跑一遍作业,得到的结果可能完全不同,因为机器处理的快慢、网络状态都会影响数据进入哪个窗口。而事件时间是以数据自身携带的时间戳为准,无论数据什么时候到,最终都应该归属到它应该属于的那个窗口。
所以在实时数仓、金融风控、用户行为分析这类对准确性要求高的场景里,我们几乎都是站在事件时间这一侧来处理窗口。Watermark 就是配合事件时间窗口工作的关键机制。
1.2 乱序从哪里来:网络、缓冲、Producer 重试
很多人一听到“乱序”两个字,第一反应是消息队列的问题。其实乱序是多层原因叠加的结果。
第一层是网络。数据从客户端发出后要经过网卡、交换机、负载均衡,任何一个环节出现抖动,都会导致部分数据晚到几毫秒甚至几秒。第二层是生产者侧的批量发送和重试。Kafka Producer 为了提高吞吐,会攒一批数据再发送,一旦发送失败还会重试,重试成功的消息可能就排到后面了。第三层是消息队列内部的分区机制。Kafka 只能保证单分区内有序,多分区之间的全局顺序是不保证的。Flink 从多个分区消费时,即使每个分区内部有序,不同分区的数据交叉到达后仍然可能是乱序的。
还有一个很隐蔽的乱序来源:上游系统处理耗时不同。比如订单服务和支付服务都在发事件,订单事件处理了 10 毫秒就发出,支付事件却处理了 500 毫秒才发出,两条事件在业务上本来支付发生在订单之后,但到达下游时订单事件反而先到,这就产生了乱序。
明白了这些来源,你就能理解,完全消除乱序是不可能的。我们能做到的,是在知道了“乱序程度大概是多少”之后,通过 Watermark 机制把迟到的数据兜住,让它依然能进入正确的窗口。
1.3 没有 Watermark 时窗口计算的尴尬:滞后数据被丢弃、统计失真
如果不用 Watermark,直接以处理时间开窗口,典型的尴尬场景是这样的。
假设你统计每五分钟的支付金额。12:00 到 12:05 这一批交易中,有一笔交易实际发生在 12:04:58,但因为支付回调网络延迟,数据在 12:05:06 才到达 Flink。处理时间窗口在 12:05:00 就触发了,这笔交易没有赶上,它会被算进下一个五分钟窗口。但如果你用事件时间配合 Watermark,这笔交易的事件时间为 12:04:58,只要 Watermark 还没有越过 12:05:00,它依然能进入 12:00-12:05 这个窗口参与聚合。
另一种尴尬是“数据到了,但是不敢触发窗口”,这是很多新手写 Flink 作业时的真实感受。如果在代码里只设置了事件时间,却没有设置 Watermark,窗口可能永远都不会触发。因为 Flink 需要一个信号来知道“当前事件时间走到哪了”,没有这个信号,窗口就只能无限等待。Watermark 就是推动窗口触发的那只手。
所以 Watermark 的价值非常明确:它把“物理到达顺序”和“逻辑事件顺序”解耦,让引擎可以根据业务时间而不是数据到达时间做窗口决策,既不会为了等迟到数据卡住整个计算,也不会因为数据晚到几秒就让统计结果彻底失真。
2. Watermark 的核心概念与生成方式:从定义到参数估算
2.1 Watermark 到底是什么:插入数据流里的“时间刻度尺”
你可以把 Watermark 理解成一条特殊的消息,它混在普通数据里一起流动,但它不参与业务计算,只负责传递一个时间信息:“当前时间已经推进到了这个时间戳,所有事件时间小于等于这个时间戳的数据,理论上都已经到达了。”
比如说有一条数据的事件时间是 12:00:30,它经过的路径上生成了一个 Watermark = 12:00:30。这个 Watermark 意味着:事件时间小于等于 12:00:30 的数据,都已经到过了。那么 Flink 看到这个 Watermark 后,就可以放心地触发结束时间不超过 12:00:30 的窗口。
当然,“理论上都已经到达”对应的是现实里的一个策略。如果实际乱序很严重,数据实际到达时间比事件时间晚很久,而 Watermark 却推进得很快,那么晚到的数据就会被错过。所以 Watermark 本质上是一个“预期最大延迟”的约定:我预计最多有 N 秒的乱序,所以我让 Watermark 始终保持在“已见到最大事件时间减去 N 秒”这个位置,给晚到的数据留出缓冲。
Watermark 还有一个关键特性:单调递增。它只能往前走,不能往后退。这是为了保证窗口触发逻辑是确定的。一旦某个窗口触发了,后面再来更小的 Watermark 也不会影响已经触发过的窗口。
2.2 两种生成器:周期性(Periodic)与逐条(Punctuated)
Flink 里生成 Watermark 的方式有两种,对应两类接口或者说两类场景。
周期性生成器是最常用的。它会每隔一段时间(默认是 200ms)根据当前已经看到的所有数据中的最大事件时间,计算一次新的 Watermark 并发射出去。这种方式的好处是开销小,不管数据流量多大,Watermark 的发射频率是固定的。生产中绝大多数作业都用这种方式,尤其适合 Kafka 这种持续不断的高吞吐数据流。
逐条生成器则是每来一条数据就判断一次,是否要生成一个新的 Watermark。它的优势是 Watermark 的精度高、响应快,如果数据流里有特殊的“标记事件”可以代表一个阶段的结束,用逐条生成很合适。但它的开销也大,实时计算场景里如果每条数据都触发一次 Watermark 判断,反而会成为性能瓶颈,所以实际用得并不多。
我在生产环境里 90% 的情况都用WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(5)),它就是一个周期性的、带固定最大乱序容忍度的生成器。只有特殊业务,比如某个事件代表“批次结束”时,才会自定义PunctuatedWatermark。新手不要一开始就陷入自定义 Watermark 的细节里,用现成的策略理解清楚就够了。
2.3 Watermark 与 Window 触发逻辑:为什么说它是“触发线”
Watermark 和窗口的配合逻辑,是理解整个机制的核心。
Flink 的窗口通常是左闭右开的区间。比如一分钟窗口,代表的是[12:00:00, 12:01:00),也就是说 12:01:00 这一秒的数据不会进入这个窗口。当一条数据到达时,Flink 先看它的事件时间属于哪个窗口,然后直接把它丢进对应窗口的状态里。窗口什么时候触发计算呢?当传入的 Watermark 大于等于窗口结束时间时,就触发。
举例来说,窗口[12:00:00, 12:01:00)的结束时间是 12:01:00。如果当前 Watermark 推进到了 12:01:00,就意味着事件时间小于等于 12:01:00 的数据都到齐了,这个窗口可以触发计算了。
这里有一个很多人容易混淆的点:Watermark 不是用来决定“数据进入哪个窗口”的,它是决定“窗口什么时候触发计算”的。数据只要事件时间落在窗口范围内,不管多晚到达,只要 Watermark 还没越过窗口结束时间,它就能进入窗口;一旦 Watermark 越过窗口结束时间,窗口触发并关闭,之后再到的属于这个窗口的数据就成了迟到数据。
所以 Watermark 推进得快慢,直接决定了窗口触发的时间。Watermark 推进得越慢,窗口触发就越晚,结果产出越滞后,但晚到数据被遗漏的概率越低。Watermark 推进得越快,结果越实时,但乱序数据丢失的可能性也越大。这是一个需要权衡的核心矛盾。
2.4 如何估算乱序容忍度:从 P95 时延到业务容忍度
设置 Watermark 的容忍度时,很多人一拍脑袋就写个 5 秒、30 秒,然后再也不管了。这样做不是不行,但往往要么数据丢得厉害,要么指标迟迟出不来。
我的习惯是,先看上游数据从产生到到达 Kafka 的端到端时延分布。你可以统计最近一天的数据,把每条记录的时间差算出来,看 P50、P95、P99 分别是多少。假设 P95 是 8 秒,P99 是 30 秒,那么如果业务允许偶尔丢 1% 的极端迟到数据,可以把容忍度设为 10 到 15 秒;如果业务要求尽量不丢数,那就设 30 秒甚至更高。
还要考虑业务对结果产出的要求。如果做的是实时大屏,希望指标尽量分钟级更新,那么容忍度设得太大就不合适,因为窗口触发会一直往后拖延。这种情况下,宁可接受少量迟到数据丢失,也要保证实时性。如果是金融交易、对账系统,晚 30 秒出结果没关系,那就可以把容忍度调大,优先保证数据完整性。
我在实际工作中一般会给表加一个“数据到达延迟”的监控指标,这样上线后能看到真实乱序情况,再回来调整参数。不要指望一次就能调准,Watermark 参数是需要根据数据特征迭代优化的。
3. 实战:在 Flink 作业里正确落地 Watermark
3.1 环境准备:Flink 安装部署与基础依赖
先说环境。我自己的习惯是本地开发用 Docker 一键起一个 Flink 集群,生产再用独立集群。如果你还不熟悉 Flink 的安装部署,这里给一个最简思路:Flink 集群需要一个 JobManager 和若干个 TaskManager,任务提交时通过 Web UI 或者命令行把 JAR 包提交进去。
docker run -d --name jobmanager \ -p 8081:8081 \ -p 6123:6123 \ -e FLINK_PROPERTIES="jobmanager.rpc.address: jobmanager" \ flink:1.17.1 jobmanager docker run -d --name taskmanager \ --link jobmanager:jobmanager \ -e FLINK_PROPERTIES="jobmanager.rpc.address: jobmanager" \ flink:1.17.1 taskmanager实际生产里,部署还涉及checkpoint存储、状态后端、高可用配置等,但核心思路一致:把你编译好的作业 JAR 包丢给 Flink,而不是直接跑一个 Java 进程。写完代码后要在 pom 里加上关键依赖,注意版本一定要和集群版本一致,否则经常会出现“连接器类找不到”的诡异问题。
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>1.17.1</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka</artifactId> <version>1.17.1</version> </dependency>版本不一致导致的问题在社区里非常多,最典型的就是ClassNotFoundException和NoSuchMethodError。遇到这类异常,第一反应不是去翻源码,而是检查 Flink 版本和依赖版本是否匹配。
3.2 DataStream API 实现事件时间 Watermark
代码永远是理解技术最直接的路径。下面这段代码是我在生产环境里最常用的写法,用forBoundedOutOfOrderness设定 5 秒乱序容忍度,然后从数据里提取事件时间作为水位线的时间戳。
DataStream<SensorReading> stream = env .addSource(new SensorSource()) .assignTimestampsAndWatermarks( WatermarkStrategy .<SensorReading>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, timestamp) -> event.getTimestamp()) ); DataStream<SensorReading> windowed = stream .keyBy(SensorReading::getId) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .allowedLateness(Time.seconds(30)) .sideOutputLateData(lateTag) .aggregate(new AvgAggregate());forBoundedOutOfOrderness的意思是:我预计乱序程度最多 5 秒,所以 Watermark 会始终保持在当前最大事件时间减去 5 秒的位置。这里有个容易被忽视的点:如果 source 一开始还没有收到数据,Watermark 是-infinity,窗口不会触发。
.withTimestampAssigner是必须的,它告诉 Flink 去数据里哪个字段取事件时间。如果这个字段是字符串类型,需要先解析成long类型的毫秒数,或者Timestamp类型。
Flink 1.14 之后,已经不再需要显式调用env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime),因为事件时间已经是默认语义了。如果你还看到老博客让你设置这个,说明那篇文章是基于老版本的,照抄会浪费你半小时排查问题。
3.3 Table / SQL API 中声明 Watermark 的两种方式
如果你更喜欢写 SQL,Flink 也支持在建表语句里声明 Watermark,这也是目前实时数仓里最常用的方式。
CREATE TABLE clicks ( user_id STRING, url STRING, click_time TIMESTAMP(3), WATERMARK FOR click_time AS click_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'format' = 'csv', 'topic' = 'clicks', 'properties.bootstrap.servers' = 'localhost:9092', 'scan.startup.mode' = 'latest-offset' ); SELECT user_id, COUNT(*) AS cnt, TUMBLE_END(click_time, INTERVAL '1' MINUTE) AS window_end FROM clicks GROUP BY user_id, TUMBLE(click_time, INTERVAL '1' MINUTE);这里WATERMARK FOR click_time AS click_time - INTERVAL '5' SECOND的含义和 Java 里的forBoundedOutOfOrderness(Duration.ofSeconds(5))完全一致。注意两个细节:
第一,事件时间字段必须是TIMESTAMP(3)类型,也就是毫秒精度。如果你从 Kafka 拿到的原始字段是BIGINT类型,得先用TO_TIMESTAMP(FROM_UNIXTIME(ts / 1000))之类的函数转一下,再声明 Watermark。
第二,Watermark 表达式必须打在事件时间字段上,并且两者类型要一致,这是 SQL 里最容易踩的坑。如果你看到 SQL 作业提交时报Invalid Watermark之类的错,优先检查这两点。
3.4 用测试数据验证 Watermark 触发窗口
写完代码,怎么确认 Watermark 真的在按预期推进?我一般会做一套最小化的本地测试,用自定义 Source 发几条固定时间戳的数据,故意制造乱序,然后观察窗口输出。
假设我们的窗口是 1 分钟,容忍度设为 5 秒。按顺序发送这几条数据:
| 发送顺序 | 事件时间 | 说明 |
|---|---|---|
| 1 | 12:00:10 | 正常数据 |
| 2 | 12:00:50 | 正常数据 |
| 3 | 12:01:10 | 已越过 12:00 窗口,但 Watermark 还未到 12:01:00 |
| 4 | 12:00:55 | 迟到数据,但 12:00 窗口尚未触发,仍能进入 |
| 5 | 12:01:20 | 推进 Watermark |
当第 5 条数据的事件时间到达 12:01:20 时,Watermark 会推进到 12:01:15(最大事件时间减 5 秒),已经超过 12:01:00,所以 12:00-12:01 的窗口会触发。这时第 4 条发送的 12:00:55 数据已经成功进入了窗口。
如果你发现窗口迟迟不触发,最直接的排查方式是看 Flink Web UI 的 Watermark 信息。如果 Watermark 一直显示-infinity,那说明数据源的事件时间字段根本没被正确提取,或者数据没有进来。这个检查点能省你至少一个小时的猜测时间。
4. 常见问题排查:Watermark 不前进、数据丢、空闲分区,一次说清
4.1 Watermark 一直不更新或不变
Watermark 停在某个值不动,是最常见的问题。出现概率最高的原因有三个。
第一个原因是 Kafka 分区没有新数据。Watermark 是跟着数据走的,如果某个分区的数据长时间不更新,这个分区对应的 Watermark 也停滞不前。多并行度时,算子取所有输入分区 Watermark 的最小值作为当前 Watermark,只要有一个分区停滞,整体的 Watermark 就会被拖住。
第二个原因是事件时间字段没有正确提取。最常见的是时间戳用了秒而不是毫秒,比如1690000000,结果 Flink 解析出来的时间比真实时间小了 1000 倍,Watermark 自然永远到不了窗口触发线。这种问题通过看 Web UI 上的 Watermark 数值就能发现。
第三个原因是 Watermark 策略设置在了错误的算子上。有些初学者在map之后再assignTimestampsAndWatermarks,但 source 数据已经被加工过,事件时间字段被丢弃或者改写了。解决方法是把assignTimestampsAndWatermarks尽量放在离 source 最近的位置。
4.2 窗口结果迟迟不触发,或该出没出
还有一个高频问题:窗口结果一直不出来,看一眼时间已经超过窗口结束时间很久了。
这种情况十有八九是 Watermark 没有越过窗口的结束时间。比如窗口是[12:00:00, 12:01:00),结束时间是 12:01:00,而 Watermark 因为容忍度设置的关系,一直停留在 12:00:55,那窗口就永远不能触发,因为 Flink 需要等到 Watermark >= 12:01:00。
解决办法是检查 Watermark 当前值与窗口结束时间的差距,然后回看事件的真实时间戳。如果你发现从当前时间看,Watermark 已经落后了很久,那说明数据源里事件时间本身就是老的,或者容忍度设置过大。大批量消费 Kafka 里历史数据但没有按顺序消费时,也会出现这种“当前 Watermark 比真实时钟小很多”的现象,因为事件时间本来就是过去的时间。
4.3 并行子任务 Watermark 对齐机制与热点
Flink 的 Watermark 在多并行度下有“对齐”机制。简单说,一个下游算子会同时接收多个上游分区的数据,它维护的 Watermark 是所有输入分区 Watermark 的最小值。这个设计很安全,因为只要有一个分区还没到,下游就不能确定数据都到齐了。
但这个安全设计也有代价。假设你有 10 个 Kafka 分区,其中 9 个分区每秒都来几百条数据,还剩 1 个分区业务量很低,可能一分钟才来一条。这 1 个分区的 Watermark 就会严重落后,导致下游所有窗口的触发都被拖慢。我遇到过一次生产事故:某交易场景正常指标 5 秒内应该出结果,结果因为一个低流量分区卡住,整个窗口延迟了 2 分钟。
解决思路有两个。一是尽量保证上游数据分布均匀,别让某个分区数据量极端偏低。二是给 Source 设置withIdleness,把长时间没有数据的分区暂时“隔离”,不参与 Watermark 对齐。
4.4 空闲分区导致 Watermark 停滞,withIdleness 怎么用
withIdleness的处理逻辑是:如果某个上游分区超过设定时间没有数据进来,就暂时把它标记为空闲,下游在计算 Watermark 时忽略它。
Java 里这样写:
DataStream<SensorReading> stream = env .addSource(new SensorSource()) .assignTimestampsAndWatermarks( WatermarkStrategy .<SensorReading>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withIdleness(Duration.ofSeconds(30)) .withTimestampAssigner((event, ctx) -> event.getTimestamp()) );SQL 里也可以通过配置项来设置空闲超时。
CREATE TABLE clicks ( ... WATERMARK FOR click_time AS click_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', ... 'scan.watermark.idle-timeout' = '30s' );注意withIdleness的时间不要设得太短。如果分区只是偶尔没有数据,比如每秒数据量本身不高,设得太短会导致分区被频繁标记为空闲,反而让 Watermark 的推进变得不稳定。我一般会根据业务数据最长的间隔来设置,通常在 30 秒到 2 分钟之间。
4.5 迟到数据三重防线:allowedLateness、sideOutputTag、超期清除
即使设了 Watermark,仍然会有数据在窗口触发之后才到达。这时候有三道防线可以处理。
第一道是allowedLateness。它允许窗口在 Watermark 越过结束时间后,再额外保留一段时间,这段时间内每来一条迟到数据,窗口就重新触发一次,输出修正后的结果。比如allowedLateness(Time.seconds(30)),窗口触发后 30 秒内还能接受属于它的迟到数据。
第二道是侧输出流sideOutputLateData(tag)。窗口彻底关闭之后,迟到数据会被丢进侧输出流,不会进入主结果。你可以定期把这部分数据落库,用于数据质量分析和重算。
第三道是状态清理。allowedLateness意味着窗口状态要保留更长时间,这会增加状态后端压力。Flink 会在 Watermark 越过窗口结束时间加allowedLateness之后,彻底清除窗口状态。所以如果数据迟到时间超过allowedLateness,就一定救不回来了。
这三道防线要配合使用。我的建议是主结果走窗口聚合,侧输出流单独接一个下游做延迟数据监控,每天对比两个指标,看看真正无法补救的数据量有多大。如果每天有大量数据进入侧输出流,说明 Watermark 容忍度设置得不够,需要调大。
5. 进阶:Watermark 在真实项目中的联动与调优
5.1 实时数仓场景下 Watermark 的取舍
实时数仓里使用 Watermark 时,取舍的维度更多。拿 Flink CDC Pipeline 来说,它现在常被用来把上游数据库的 Binlog 同步到 Kafka、Iceberg、Hive 等下游。CDC 数据本质上是按事务提交时间产生的,如果在同步链路里还嵌入了窗口聚合逻辑,Watermark 的作用就很关键。
比如你要实时统计“下单后 10 分钟内付款”的用户数,订单事件和支付事件可能来自不同的表,它们到达 Flink 的时间顺序并不保证。如果不设置 Watermark,两股数据在双流 join 时可能一直对不上。设置了统一的 Watermark 之后,双流 join 才能在一个合理的时间范围内把关联数据配对。
但实时数仓里也不能盲目追求“Watermark 越大越好”。因为下游 Hive、Iceberg 通常有分区的提交策略,窗口延后会直接影响数据产出时间,进而影响下游依赖这个指标的数仓任务。我的经验是,先明确业务能接受的结果延迟底线,再反推 Watermark 容忍度,而不是让 Watermark 无限放大。
5.2 自定义 DataSource / Sink 时如何传递 Watermark
热词里经常出现“自定义 DataSource 与 DataSink”的需求。这里单独提醒一下:如果你自己写了一个SourceFunction,并且希望 Watermark 能正常工作,需要在 source 里显式调用collectWithTimestamp或者通过SourceContext.emitWatermark来发射 Watermark。
使用 DataStream API 时,很多企业级 source connector 已经封装好了 Watermark 的生成,你只需要在 source 之后调用assignTimestampsAndWatermarks。但如果你用的是自定义 source,又要用事件时间,常见的做法是:
@Override public void run(SourceContext<MyEvent> ctx) throws Exception { while (running) { MyEvent event = readFromExternal(); ctx.collectWithTimestamp(event, event.getEventTime()); if (hasNewMaxTime()) { ctx.emitWatermark(new Watermark(getCurrentMaxTimestamp() - maxOutOfOrderness)); } } }自定义 sink 时,Watermark 的传递并不像 source 那样重要,因为 sink 通常只需要按自己的连接器逻辑写入数据。但有一点要注意:如果你在 sink 端做了幂等写入或者事务性写入,Watermark 的延迟不能作为拒绝写入的条件,否则会造成数据丢失。
5.3 JDBC/Hive 连接器常见异常与超时参数
把窗口聚合结果写到 MySQL 或者 Hive,是常见的落地方式。热词里的“flink 的 jdbc 连接器异常”“flink sink hive 表数据不入表”,我都遇到过,很多问题和 Watermark 有间接关系。
先说 JDBC 异常。窗口触发通常是一个集中爆发的过程,比如一分钟窗口到时,几千条聚合结果同时要写入数据库,如果连接池不够,就会出现Too many connections、Connection reset这类报错。这时候你需要调整sink.buffer-flush.max-rows、sink.buffer-flush.interval、max-retries等参数,同时把数据库连接池上限调高。别把锅甩给 Watermark,它只是把压力集中到了触发那一刻。
再说 Hive 表不入数据。这个问题最常见的原因是窗口没触发,所以根本没有数据产出到下游。你若在 Flink Web UI 看到窗口一直不触发,那就是 Watermark 的问题。另一个原因是用了StreamingFileSink或FileSink,数据到了但还没有做分区提交。Hive 分区提交依赖 checkpoint,如果你没开 checkpoint,数据永远只停留在临时目录,看起来就像是“数据没写进去”。这时候开启 checkpoint,并配置sink.partition-commit.policy.kind=success-file之类参数,才能真正把临时文件提交成 Hive 分区。
5.4 通过火焰图排查反压与 Watermark 延迟
有时候 Watermark 不推进,不是逻辑写错了,而是作业被反压拖住了。所谓反压,就是下游处理不过来,上游只能停下来等待,结果数据源读不到新数据,Watermark 自然就停住了。
这时候光看 Web UI 可能不够。打开 Flink 的火焰图,你能看到每个算子 CPU 耗时集中在哪个方法上。比如某个KeyBy之后的热 key 导致某个子任务计算量巨大,其他子任务都空闲,这个热点子任务的 Watermark 就会一直落后,拖住整个作业。
我有一次排查线上问题,发现所有事件都集中在一个 userId 上,那个并行子任务的负载是其他子任务的 20 倍。通过火焰图定位到这个热点之后,我给 key 做了加盐处理,把一个热 key 拆成多个子 key,负载才均衡下来。这个问题的表象一直是“Watermark 不准、窗口触发晚”,但根因完全不在 Watermark 逻辑,而在数据分布。
排查这类问题的建议是:先确认 Web UI 上 Watermark 的具体数值和 Last Checkpoint Size,再结合火焰图看 CPU 热点,最后才去怀疑 Watermark 参数设置。顺序反了,很容易在错误的方向上浪费大量时间。
最后再分享几个实际经验
我在生产环境踩过最深的坑,就是把 Watermark 容忍度设得太大,结果每五分钟窗口的结果要等 15 秒才出来,业务方受不了。后来改小容忍度,数据又开始丢。最后的解决办法是分两层处理:主结果用小容忍度保证实时性,侧输出流把被丢弃的迟到数据单独落库,每天晚上做一次离线对账,修正当天的统计结果。
还有一个小技巧,流里每条数据尽量带上两个时间字段,一个是业务发生时间,一个是进入 Kafka 的时间。这样可以非常方便地画出“网络传输延迟分布”,用来验证你的 Watermark 容忍度是否合理。否则你只是知道“乱序”,但不知道乱序有多严重,调参就很盲目。
Watermark 不是一个死板的配置,它本质上是在实时性和准确性之间找一个动态平衡点。希望这篇文章能把原理中的“为什么”讲透,也能让这套思路直接复用到你自己的 Flink 作业里。