Flink中latency metrics的含义及是否可有效评估应用延迟?
Flink Latency Metrics 说明与压测实践方案
你当前开发的Flink数据管道代码如下:
SingleOutputStreamOperator<String> stream = ... DataStream<String> branch2 = stream .getSideOutput(outputTag2) .keyBy(MetricObject::getRootAssetId) .window(TumblingEventTimeWindows.of(Time.seconds(180))) .trigger(ContinuousEventTimeTrigger.of(Time.seconds(15))) .aggregate(new CountDistinctAggregate(),new CountDistinctProcess()) .name("windowed-count-distinct") .uid("windowed-count-distinct") .map(AggregationObject::toString) .name("get-toString");
针对你关注的latency metrics相关问题,解答如下:
1. Flink 内置 Latency Metrics 的具体含义
你通过env.getConfig().setLatencyTrackingInterval(1000)开启的延迟追踪,核心逻辑是:
- 每隔指定时间(这里是1000ms),Flink会从Source端生成特殊的延迟标记(Latency Marker),标记携带了Source端的生成时间戳
- 延迟标记不会参与任何业务计算逻辑,只会跟着数据流的传输路径,经过各个算子的网络队列、输入队列流转
- 每个算子收到延迟标记时,会计算当前时间和标记生成时间的差值,作为该算子的latency指标上报
这个指标本质上统计的是数据在Flink集群内部传输、队列排队的耗时,不包含任何业务逻辑的计算耗时。
2. 能否有效评估你的应用延迟
不能直接用来评估你这个场景的业务端到端延迟,原因如下:
- 你的管道核心是180s滚动窗口+15s连续触发的聚合逻辑,业务数据需要在窗口中等待事件时间推进、积攒到触发条件才会输出,这个等待耗时是你业务延迟的核心组成部分,但延迟标记会直接跳过窗口计算逻辑流转,完全不会统计这部分耗时,最终得到的latency指标会远低于实际业务延迟
- 它的适用场景是判断管道的底层传输/排队瓶颈:当出现反压、算子队列堆积时,latency指标会立刻上升,这个判断是准确的
3. 压测场景下的使用方案
你计划用不同速率压测来监控吞吐量、延迟、反压节点,可以搭配两类指标使用:
3.1 内置Latency Metrics的用法
- 反压判断:逐步提升发送速率时,如果某个算子的latency指标突然陡增,说明该算子处理能力达到上限,上游已经开始出现反压
- 吞吐量瓶颈判断:配合监控算子的
numRecordsOutPerSecond指标,如果该指标不再随发送速率提升而增长,同时latency指标持续上升,说明集群吞吐量已经达到上限
3.2 补充自定义业务延迟指标
要获取真实的业务端到端延迟,你需要新增自定义指标:
- 在Source端给每条原始数据追加一个
sourceIngestTimestamp字段,记录数据进入Flink的时间 - 在窗口的
CountDistinctProcess处理逻辑中,每次触发输出聚合结果时,统计当前窗口内所有数据的sourceIngestTimestamp最大值,用当前系统时间减去这个最大值,得到本次窗口输出结果的端到端延迟,将该值作为自定义Gauge指标上报
这个自定义指标才是你业务层面真实的延迟数据,适合用来判断不同压测速率下业务延迟的变化情况。
内容的提问来源于stack exchange,提问作者Canelupo
相关产品推荐
相关产品推荐

