flink网络传输的前世今生

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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
private boolean copyFromSerializerToTargetChannel(int targetChannel) throws IOException, InterruptedException {
// We should reset the initial position of the intermediate serialization buffer before
// copying, so the serialization results can be copied to multiple target buffers.
// 这一步reset是为了在数据发送如果是broadcast这种一份数据需要发送多个下游通道的时候,就可以只序列化一次,后续数据发送的时候只需要将bytebuffer
// 的position值值置到0就可以了。
serializer.reset();

boolean pruneTriggered = false;
BufferBuilder bufferBuilder = getBufferBuilder(targetChannel);
SerializationResult result = serializer.copyToBufferBuilder(bufferBuilder);
// buffer没写满说明数据肯定已经写完了,直接进行下面的逻辑
while (result.isFullBuffer()) {
// buffer写满了,首先将bufferBuilder标记为写完了,就是将positionMarker置为相反数
numBytesOut.inc(bufferBuilder.finish());
numBuffersOut.inc();

// If this was a full record, we are done. Not breaking out of the loop at this point
// will lead to another buffer request before breaking out (that would not be a
// problem per se, but it can lead to stalls in the pipeline).
// buffer写满,并且记录也写满了,那么发送到这个channel就完成了,否则就需要继续申请bufferBuilder继续拷贝
if (result.isFullRecord()) {
pruneTriggered = true;
bufferBuilders[targetChannel] = Optional.empty();
break;
}

bufferBuilder = requestNewBufferBuilder(targetChannel);
result = serializer.copyToBufferBuilder(bufferBuilder);
}
}
1
2
3
4
5
6
7
8
@Override
public SerializationResult copyToBufferBuilder(BufferBuilder targetBuffer) {
targetBuffer.append(lengthBuffer);
targetBuffer.append(dataBuffer);
targetBuffer.commit();

return getSerializationResult(targetBuffer);
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
public int append(ByteBuffer source) {
checkState(!isFinished());

int needed = source.remaining();
int available = getMaxCapacity() - positionMarker.getCached();
// segment不一定足够大,可能存不下这批buffer, 堆外内存拷贝的时候需要提前计算好可以拷贝的量,否则会有异常
int toCopy = Math.min(needed, available);

// 将source buffer中的数据/堆内存,put至memorySegment中,利用Unsafe进行数据拷贝
memorySegment.put(positionMarker.getCached(), source, toCopy);
// 设置新的position
positionMarker.move(toCopy);
return toCopy;
}

BufferBuilder,BufferConsumer,PositionMarker

在上面copy代码中看到其实拷贝的时候是依赖buffer的,如果没有申请到BufferBuiler,是会一直blocking的,那么这个bufferbuilder是什么呢?

1
2
3
4
5
6
7
8
9
private BufferBuilder requestNewBufferBuilder(int targetChannel) throws IOException, InterruptedException {
checkState(!bufferBuilders[targetChannel].isPresent() || bufferBuilders[targetChannel].get().isFinished());

BufferBuilder bufferBuilder = targetPartition.getBufferProvider().requestBufferBuilderBlocking();
bufferBuilders[targetChannel] = Optional.of(bufferBuilder);
// 一个bufferbuilder对应一个bufferconsumer
targetPartition.addBufferConsumer(bufferBuilder.createBufferConsumer(), targetChannel);
return bufferBuilder;
}

在向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
2
3
4
5
6
7
8
9
10
public Buffer build() {
// 获取最近builder,commit到的position
writerPosition.update();
int cachedWriterPosition = writerPosition.getCached();
// slice 切分只读区块
Buffer slice = buffer.readOnlySlice(currentReaderPosition, cachedWriterPosition - currentReaderPosition);
currentReaderPosition = cachedWriterPosition;
// 增加引用计数
return slice.retainBuffer();
}

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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
/**
* Receives the backlog from the producer's buffer response. If the number of available
* buffers is less than backlog + initialCredit, it will request floating buffers from the buffer
* pool, and then notify unannounced credits to the producer.
*
* @param backlog The number of unsent buffers in the producer's sub partition.
*/
void onSenderBacklog(int backlog) throws IOException {
int numRequestedBuffers = 0;

synchronized (bufferQueue) {
// Similar to notifyBufferAvailable(), make sure that we never add a buffer
// after releaseAllResources() released all buffers (see above for details).
if (isReleased.get()) {
return;
}

numRequiredBuffers = backlog + initialCredit;
// 检查当前input通道的buffer是否做够上游produce所需要的buffer,如果不够就去bufferpool申请
while (bufferQueue.getAvailableBufferSize() < numRequiredBuffers && !isWaitingForFloatingBuffers) {
Buffer buffer = inputGate.getBufferPool().requestBuffer();
if (buffer != null) {
// 申请到buffer之后先占据住
bufferQueue.addFloatingBuffer(buffer);
numRequestedBuffers++;
// 没有足够的buffer,那么注册回调等buffer回收
} else if (inputGate.getBufferProvider().addBufferListener(this)) {
// If the channel has not got enough buffers, register it as listener to wait for more floating buffers.
isWaitingForFloatingBuffers = true;
break;
}
}
}

// 如果生产端有buffer需求,并且之前的unannouncedCredit为0那么就需要通知上游有buffer了
if (numRequestedBuffers > 0 && unannouncedCredit.getAndAdd(numRequestedBuffers) == 0) {
notifyCreditAvailable();
}
}

整理流程图

flink-network

netty内存的优化

以下是message encode的时候的一段代码

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// only allocate header buffer - we will combine it with the data buffer below
headerBuf = allocateBuffer(allocator, ID, messageHeaderLength, buffer.readableBytes(), false);

receiverId.writeTo(headerBuf);
headerBuf.writeInt(sequenceNumber);
headerBuf.writeInt(backlog);
headerBuf.writeBoolean(isBuffer);
headerBuf.writeInt(buffer.readableBytes());

CompositeByteBuf composityBuf = allocator.compositeDirectBuffer();
composityBuf.addComponent(headerBuf);
composityBuf.addComponent(buffer);
// update writer index since we have data written to the components:
composityBuf.writerIndex(headerBuf.writerIndex() + buffer.writerIndex());
return composityBuf;

可以看到这里和以前版本不一样的地方就是不需要再去申请一块netty内存做一次拷贝,因为这里将buffer对象的实现直接改成了继承netty的ByteBuf类,
所以减少了一次netty申请directBuffer以及从堆外拷贝到netty directBuffer的开销。在buffer处理完由netty回收时会放回localBufferPool

1
2
3
4
@Override
protected void deallocate() {
recycler.recycle(memorySegment); // 在网络传输完内存释放的时候直接将segment回收到localbufferpool中
}

和flink1.4相比有了哪些改进

https://docs.google.com/document/d/1chTOuOqe0sBsjldA_r-wXYeSIhU2zRGpUaTaik7QZ84

https://issues.apache.org/jira/browse/FLINK-7282?subTaskView=all

谢谢支持