flink的网络传输在1.5版本进行了重构,本文就这个feature来对flink网络传输进行系统的源码分析
发送端
首先我们先来看数据发送端的主要流程如下:
数据发送链路
RecordWriter => ResultPartition => ResultSubPartition => ResultSubpartitionView => BufferAvailabilityListener => PartitionRequestQueue
解释一下: 用户程序中调用output.collect(),首先会通过RecordWriter进行数据或者event的序列化。并且将其从堆内内存拷贝至堆外内存,然后添加至相应的ResultPartition中。ResultPartition根据数据selectChannel发送给下游的哪个subIndex,BufferConsumer就会被添加到相应的subpartition所维护的一个双端队列中。在某些条件下需要通过,在服务启动最开始注册上来的ResultSubpartitionView去通知消费端来进行消费buffer,view做的事情就是通过调用BufferAvailableListener的具体实现来进行通知事件通知。最终在netty端,通过PartitionRequestQueue进行最终的buffer发送。
上面讲述了大体的流程,下面我们来结合代码来进行细节分析,下面代码可能会结合1.4和1.7两个版本来进行讲解
RecordWriter
RecordSerializer
在1.4版本中序列化器是和下游的并发度一一绑定的,这样会导致一个问题,比如发送下游是hash的分区模式的话,在上游的每一个并发度就会存储5MB的序列化后的缓存数据,当下游的并发较大的时候就会占据比较大的内存,带来一定的gc问题。序列化器会负责做这样几件事情:
- 数据的序列化
- 在写数据的时候并不会校验缓存块的大小
- 写的时候同时用一个4字节的bytebuffer记录数据的大小,有多少个字节
- 数据序列化结果的拷贝,对拷贝结果的判断
- 数据拷贝了一部分,memorysegment已经满了
- 拷贝了完整记录
- 拷贝了完整记录,并且segment满了
- 缓存清理
- …
重点拷贝过程
1 | private boolean copyFromSerializerToTargetChannel(int targetChannel) throws IOException, InterruptedException { |
1 | @Override |
1 | public int append(ByteBuffer source) { |
BufferBuilder,BufferConsumer,PositionMarker
在上面copy代码中看到其实拷贝的时候是依赖buffer的,如果没有申请到BufferBuiler,是会一直blocking的,那么这个bufferbuilder是什么呢?
1 | private BufferBuilder requestNewBufferBuilder(int targetChannel) throws IOException, InterruptedException { |
在向BufferProvider,一般是localBufferPool申请完得到一个memorysegment后,将其封装成一个bufferbuilder,每一个bufferbuilder会对应
一个bufferconsumer和positionMarker,positionMarker会标记生产端的数据写到多少个字节了,这个在消费端的时候也会用到这个position,
由于是多线程使用所以position的值需要被标记成volatile来保证数据的可见性,每次消费端拉取数据的时候,对于没有写完的buffer同样可以进行消费,
消费前更新一个buffer的position真实位置,这里用到了一个小技巧,由于数据在生产的时候需要频繁的更新position,如果是volatile的,
虽然比较轻量,频繁更新也是比较大的开销,因此加入了一个cachedPosition,在写数据的时候只需要更新builder中的cachedPosition,生产端每次
完成一批的书写才会commit给volatile position,以此来减少缓存刷新。
从一个正在写的bufferbuiler中构建一个可消费的slice
1 | public Buffer build() { |
PartitionRequestQueue
在将bufferConsumer添加到subpartition的队列之后,同时会在partitionRequestQueue中维护一个availableReader的队列,这个队列表示可以往下
下游发送的buffer数据,这样通过一个while true循环持续的将队列中的数据往下游发送,当然这个availableReader队列的维护既考量了上游subpartition
有没有buffer的因素,也考量了下游要发送的receiver端的credit的情况,如果没有credit也是无法进入这个待发送队列的。
消费端
数据接收链路
CreditBasedPartitionRequestClientHandler => RemoteInputChannel => SingleInputGate => BarrierHandler => StreamInputProcessor => StreamOperator
首先会通过netty client进行数据的接收,然后从localbufferpool申请内存接收数据,然后根据backlog的信息去决定是不是要给上游分发credit,以及数据处理的流程
这里主要分析下credit的判断逻辑
1 | /** |
整理流程图

netty内存的优化
以下是message encode的时候的一段代码
1 | // only allocate header buffer - we will combine it with the data buffer below |
可以看到这里和以前版本不一样的地方就是不需要再去申请一块netty内存做一次拷贝,因为这里将buffer对象的实现直接改成了继承netty的ByteBuf类,
所以减少了一次netty申请directBuffer以及从堆外拷贝到netty directBuffer的开销。在buffer处理完由netty回收时会放回localBufferPool中
1 | @Override |
和flink1.4相比有了哪些改进
https://docs.google.com/document/d/1chTOuOqe0sBsjldA_r-wXYeSIhU2zRGpUaTaik7QZ84
https://issues.apache.org/jira/browse/FLINK-7282?subTaskView=all