[toc]
本文翻译自flink官网的一篇博文,详细介绍Flink架构中的网络栈,下面就让我们来一睹为快吧。
简介
Flink网络栈是flink中的核心组件,是flink-runtime模块的一部分。它连接了所有TaskManager中独立的工作单元(subtasks)。这是流入数据流经的地方,因此你所观察到的吞吐量和延迟都和他息息相关,可以说Flink的网络栈决定了Flink框架本身性能的好坏。和TaskManager,Jobmanager之间通信所使用的akka rpc框架不同的是,flink网络栈采用了更底层的网络api,使用的是Netty框架。
Logical View

它抽象了以下三个概念的不同设置:
- Subtask output type (ResultPartitionType):
- pipelined (bounded or unbounded): 上游一产生数据,就一条一条的往下游发送,作为有界或无界的数据流
- blocking: 只有当上游的全部结果就绪之后才向下游发送数据
- Scheduling type:
- all at once (eager): 同时部署所有的subtask(流式应用采用这种模式)
- next stage on first output (lazy): 当上游的生产者开始有输出结果的时候,才开始不是下游的subtask,是一种lazy模式
- next stage on complete output: 当上游的数据全部就绪之后才开始下游subtask的部署。
- Transport:
- high throughput: 不采用一条一条发送数据的模式,Flink缓存一批数据到network buffer中,攒批发送。这个减少了网络开销的单条边际成本,带来了高吞吐
- low latency via buffer timeout: 通过减少发送未攒满一个buffer的timeout时间,牺牲一定的吞吐带来更低的延迟
我们将在下面这部分讨论吞吐和延迟的优化。这一部分将查看网络栈的物理层,对于这部分,我们将详细阐述output type以及调度模式。首先,subtask的输出类型和调度类型是紧密交织在一起的,两者的特定组合才有效。Pipelined result partition是流式的输出,流式输出需要将数据发送到一个正在工作的subtask,因此目标task就需要在上游结果产出下发之前deploy完成或者在任务启动最初完成deploy。 Batch作业产出有限的结果,而stream作业产出无限的结果。
Batch作业也可以以阻塞的方式产出结果,具体取决于operator和conector的使用模式。在这种方式下,下游算子需要再上游结果完全ready之后才进行部署,在这种方式下资源使用效率会更高.
下表总结了有效的组合方式:
| Output Type | Scheduling Type | Applies to… |
|---|---|---|
| pipelined, unbounded | all at once | Streaming jobs |
| next stage on first output | n/a¹ | |
| pipelined, bounded | all at once | n/a² |
| next stage on first output | Batch jobs | |
| blocking | next stage on complete output | Batch jobs |
【1】: 当前没有再flink中使用
【2】: 在Batch/Streaming unification完成后可能在流式作业中使用
另外对于有多个输入的subtask的batch作业,调度开始有两种模式,在所有input产出数据或者任意一个输入产出数据的时候。
Physical Transport
为了理解真实的吴礼数据的连接,请回想一下,在Flink中,不同的task可能会共享一个同一个slot,通过slot sharing group机制,TaskManager也可以提供多个slot来允许一个task的多个subtask跑在一个TaskManager上。
举个下图中的例子,我们假想一个有四个并发的任务,部署在两个分别有2个slot的TaskManager上。TaskManager 1 运行subtask A.1,A.2,B.1 和 B.2 而TaskManager 2 运行subtaskA.3,A.4,B.3,B.4。假设A和B之间的shuffle方式是keyBy(),这样在每一个TaskManager上都有2x4个逻辑连接,有些走local的,有些是通过网络的,如下图所示。

不同task之间的每个(远程)网络连接,都将在flink网络栈中获得自己的TCP通道,但是,如果同一个任务的不同子任务被调度到同一个TaskManager上,他们和另一个TaskManager上的TCP连接将会共享(多路复用),在我们的例子中A.1 -> B.3, A.1 -> B.4 以及A.2 -> B.3,A.2 -> B.4将会复用一个tcp连接(这里有点疑问,按我的理解应该是每个TM之间都会共享channel才对)。

每个subtask的输出被称作ResultPartition, 每一个又被细分为ResultSubPartition,一个逻辑channel会有一个。在这个阶段,Flink已经不再单独处理每条记录了,而是将一组序列化完的数据打包拷贝到network buffer中,每一个subtask中local buffer pool(发送端和接收端各有一个pool)中所能获取的最多的buffer数量是通过以下的配置决定的
1 | channels * buffers-per-channel + floating-buffers-per-gate |
通常一个TaskManager上总的buffer数量不需要配置,在需要是可以查看相关的配置项Configuring the Network Buffers
Inflicting Backpressure (1)
当一个subtask的发送端的buffer用尽之后 - buffer可能用于result subpartition的buffer队列,也可能正用于低阶的Netty的网络栈中尚未回收。在这种情况下producer就被block住,无法进一步发送数据。
消费端也是一样的方式,从底层netty网络栈传输来的buffer需要通过network buffer才能被正确消费,如果在消费端的network buffer用尽了,Flink将停止从channel中读入数据,直到有新的network buffer用来做数据的转化。这种情况下将会反压上游所有通过这个tcp channel的多路复用发送数据的生产端,这样也就限制了其他下游的消费能力。在下图阐述了一个有性能瓶颈的B.4 subtask将会引起backpressure,最终将会引起B.3无法消费和处理新的数据,尽管B.3还有足够的network buffer。

为了彻底解决这个问题,Flink1.5引入了流控机制。
Credit-based Flow Control
基于流控的网络传输能够确保正在传输的数据在接受端都有可以接受的buffer(传输的数据都是得到确认之后才可以向下游发送的)。新的传输模式仍然基于Flink原有的network buffer的能力,在其上做了一些扩展。除了仅仅有一个共享的本地buffer pool,每一个remote inputchannel现在会有其自己独占的一批buffer池,相对地,共享池中的buffer就被称为发floating buffer因为他们对于每一个inputchannel都是可以获取的。
接收端现在将会以Credits的形式通知发送端告知它能够接收多少buffer(1 buffer = 1 credit)。 每一个result subpartition将持续记录相对应channel的credits。只有当下游的通道有credits的时候,才会通过netty将上游的buffer发送出去,每发送一个buffer,就会减去一个credits,除了发送buffer数据,同时还会携带当前的backlog信息(指的是当前这个subparititon还有多少buffer在等待发送),接收端拿到这个信息之后就会去申请相应的合适数量的floating buffer来处理subpartition中排队等待处理的buffer。接收端一般会申请和backlog一样多的buffer,但这个并不总是能申请到这么多,也可能当前根本没有buffer可以申请。接收端将会通过监听buffer池,等待buffer使用完回收的时候就可以开始新的buffer申请

流控模式下使用buffers-per-channel来指定每一个channel的独占buffer的大小,使用floating-buffers-per-gate来指定local buffer pool的大小这两个参数的默认值理论上可以达到和非流控模式下的最大吞吐量,可能你需要根据你的网络带宽或者rt要求来调整相关的参数
Inflicting Backpressure (2)
和没有流控的模式相比,credits能够提供更加直接有效的控制: 如果消费端无法快速处理造成消费端的buffer用尽,这样该recevier的credits就会降为0,发送端就不会再发送数据到这个partition。反压只发生在了这个通道上,并不会影响tcp通道的数据传输,因此其他接收端的消费能力并不会受到影响。
What do we Gain? Where is the Catch?
通过这样的优化整体的资源使用率应该会得到上升,并且通过对正在传输的数据量的直接控制,也带来了checkpoint对齐时间的优化。但是receiver端发送消息也带来了额外的开销(在这之前接收端是只要负责收数据就好了,没有发消息的动作),尤其是在开启了SSL加密的情况下。并且一个input channel无法使用所有的buffer pool,因为独占的buffer无法共享。而且如果你产出数据比你发布credits慢,就会导致无法尽可能快的下发数据,虽然这些情况会带来性能上的损失,但是通常还是建议开启流控的模式,因为带来的诸多好处。
另一件你可能会注意到的是因为我们在发送端和接收端之间会缓存更少的数据可能会更早的碰到反压,可以通过调节上述的exclusive和floating buffer的大小来缓解
Writing Records into Network Buffers and Reading them again
下图对上面的图做了一些扩展,添加了一些周边组件使用的一些细节

在创建一个record将其传递给下游,比如通过Collector#collect(),它将被传递到RecordWriter中,recordwriter将这个java对象序列化成二进制数组,最终被拷贝至networkbuffer中像上述几节中的处理方式进行处理。RecordWriter首先将数据通过SpanningRecordWriter序列化到对上的一个byte数组中。然后将会将这些byte数组写到相应目标channel的的network buffer中,我们再最后一小节再回来讨论这块细节。
在接收端,netty会将接收到的buffer写入相应的input channel,流式任务的task最后将会读取这些队列中的buffer,将其通过RecordReader反序列化出来,和序列化过程类似,反序列也需要处理一些特殊情况,比如一条记录跨过了多个network buffer,可能是因为一条记录比较大,大过了单个networkbuffer大小(32K),也有可能数据被加入networkbuffer的时候已经没有足够的容纳空间了。
Flushing Buffers to Netty
在上图中,流控的机制实际上就是在NettyServer和NettyClient组件中实现的,RecordWriter正在写的buffer总是最开始是空的没有数据的状态被加入到result subpatition中,然后再渐渐的被序列化的记录填满,但是netty是什么时候真正处理这些buffer呢?
显然他不能一有数据就去请求,因为这会带来巨大的开销,涉及到跨线程通信和同步,并且也会将整个buffer废弃掉。
在Flink中有3个情况可以将buffer变为NettyServer可以消费的状态:
- buffer写满了
- buffer timeout时间条件满足了
- 一个特殊的event发送了,比如checkpoint barrier,为了保证一致性就需要发送
Flush after Buffer Full
在序列化完成后,RecordWriter会将这些bytes写出到合适的subpartition的network buffer队列中去,尽管一个RecordWriter可以处理多个subpartition,但是每一个subpartition都只会有一个writer向其写数据。NettyServer会从多个subpartition读取buffer,通过一个channel发送出去。这是经典的生产者和消费者模式,正如下图所示。在(1)序列化和(2)将数据写出到buffer,RecordWriter会更新buffer的writer下标,一旦buffer已经满了,RecordWriter会从他的local buffer pool中申请一块新的buffer用来写当前记录的剩余的bytes,或者写下一条记录,并且将新的buffer添加搭配subpartition的queue中。 (4)通知NettyServer数据已经ready,当netty能够发送数据时,就会(5)从queue中获取buffer,通过TCP channel将数据发送给下游。

Flush after Buffer Timeout
为了能够支持低延迟的处理场景,我们不能只依赖buffer变满的下发信号。可能会由于上游某些通道数据不多带了发送延迟。因此还有周期性的发送机制:OutputFlusher。下图展示了这个flusher机制是如何和其他组件协同工作的。
RecordWriter进行序列化并将数据写入network buffer。但是并行的output flusher(3,4)会通知NettyServer进行数据的消费。当NettyServer收到通知之后,他将会消费buffer中可以消费的数据,并且更新buffer的reader index。buffer仍然会保留在队列中,下一次对该buffer的读取操作会从上次的reader index开始。严格的说flusher并没有保障数据的发送,他仅仅是通知NettyServer可以进行数据发送,如果正处于反压状态,flusher并没有什么实际效果

Flush after special event
一些特殊消息也会触发数据的flush,最重要的就是checkpoint barrier和end-of-stream
Future remarks
和Flink1.5之前相比,network buffers现在被放置在了subpartition的队列中,而且每次flush我们并不会关闭buffer,这给我们带来一些好处:
- 减少了同步带来的消耗(output flusher和RecordWriter是相互独立的)
- 在负载高,Netty是性能瓶颈的情况下我们仍然可以在不完整的buffer中累积数据
- 减少了netty的通知信号
但是你可能也注意到在低负载的场景下,cpu使用率和TCP发包率会变高,这是因为flink会利用cpu来达到想要的低延迟,在高负载情况下性能可能会更好一些,因为去除了一些同步的消耗。
Buffer Builder & Buffer Consumer
这一节可以参考我之前博客中的相关分析BufferBuiler/BufferConsumer
Latency vs. Throughput
Network buffers是被使用来以期望达到更高的资源利用率,并且让数据再buffer中等待一段时间来攒批发送来达到更高的吞吐。尽管这个等待时长可以通过设置buffer timeout时间。你可能会好奇想找出latency和thoughout之间的平衡关系。显然,你是不可能同时拥有这两者的。下图显示了不同的time out时间设置,0ms-100ms所带来的吞吐量的提升,这些测试是跑在100个就节点,每个节点8个slots,没有业务逻辑只有纯粹的网络栈开销。

如你所见,Flink1.5+,即使是非常低的缓冲区超时(例如1ms)也提供高达默认超时的75%的最大吞吐量。
参考: https://flink.apache.org/2019/06/05/flink-network-stack.html