关于Flink中的Latency track
目的
为了监控Flink数据端到端的数据延迟
解决方案
- 定时在source处定时的发送一个特殊的event,类似于watermark的处理方式,这里叫做LatencyMark。
- latencyMarker仅在source处产生,latencyMarker对象包含的是source operator的区分标志 subIndex以及vertexID(根据subIndex以及vertexID区分是否是同一个Marker)以及携带了发送的时间
- LatencyMark在各个operator间传递,在每个operator处将会比较LatencyMark和它当前的系统时间来决定延迟的大小,并存入LatencyGauge,每一个operator都会维护这样一个Metric(因此LatencyMark的实现就是基于TM和JM集群的机器系统时间是进行过同步的)
- 当Operator有多个Output的时候,他会随机选择一个来发送,这确保了每一个Marker在整个流中只会出现一次,repartition也不会导致LatencyMark的数量暴增。
在sink operator处会维护source的最近128个latencyMarker,通过一个LatencyGauge来展示
具体实现
默认latencyTrackingInterval是2000,也就是2s发送一个LatencyMarker。在StreamSource中判断,如果开启了latency track.那么就会定期发送LatencyMarker。
在StreamSource(Operator)中的定时发送
1 | latencyMarkTimer = processingTimeService.scheduleAtFixedRate( |
在Operator中处理:首先进行LatencyMarker处理再发送,
在AbstractStreamOperator类中定义的latencyMarker处理
1 | public void reportLatency(LatencyMarker marker, boolean isSink) { |
然后只要不是sink类型的operator,就会往后继续传递LatencyMarker。随机选择一个Channel来发送,这里就是为了保证一个latencyMarker在整个流中只会出现一次。这里和watermark的机制有点不一样,waterMark是遍历全部的channel来发送。
1 | public void emitLatencyMarker(LatencyMarker latencyMarker) { |
LatencyMarker不会参与窗口时间的计算,应该说是不参与任何operator的计算,因此他只能用来衡量数据在整个DAG中流通的速度不能衡量operator计算的时间,这个只能通过单测来进行计算StreamInputProcessor#processInput 这里进行对进入的element进行处理,对于watermark和LatencyMarker类型会先处理发送掉,不会经由后面的windowOperator或其他operator来处理
1 | if (result.isFullRecord()) { |
sink中只进行report不再进行forward了(StreamSink.java)
1 | protected void reportOrForwardLatencyMarker(LatencyMarker maker) { |
总结
- LatencyMarker能够较好的监控因网络抖动或数据反压引起的延迟,可以提前预警反压情况
- 在正常情况(没有反压)下数据在DAG图中的流动延迟大概0.5s左右,所以说Flink的确是一个很快的引擎:)