You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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相关问题,解答如下:

你通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.24 09:45:01