高校服务地方专题网站建设,网络管理系统的基本组成和功能,查看自己网站访问量,外贸网店怎么开Flink 的网络协议栈是组成 flink-runtime 模块的核心组件之一#xff0c;是每个 Flink 作业的核心。它连接所有 TaskManager 的各个子任务(Subtask)#xff0c;因此#xff0c;对于 Flink 作业的性能包括吞吐与延迟都至关重要。与 TaskManager 和 JobManager 之间通过基于 A…Flink 的网络协议栈是组成 flink-runtime 模块的核心组件之一是每个 Flink 作业的核心。它连接所有 TaskManager 的各个子任务(Subtask)因此对于 Flink 作业的性能包括吞吐与延迟都至关重要。与 TaskManager 和 JobManager 之间通过基于 Akka 的 RPC 通信的控制通道不同TaskManager 之间的网络协议栈依赖于更加底层的 Netty API。
本文将首先介绍 Flink 暴露给流算子(Stream operator)的高层抽象然后详细介绍 Flink 网络协议栈的物理实现和各种优化、优化的效果以及 Flink 在吞吐量和延迟之间的权衡。
1.逻辑视图
Flink 的网络协议栈为彼此通信的子任务提供以下逻辑视图例如在 A 通过 keyBy() 操作进行数据 Shuffle
这一过程建立在以下三种基本概念的基础上
▼ 子任务输出类型ResultPartitionType Pipelined有限的或无限的一旦产生数据就可以持续向下游发送有限数据流或无限数据流。 Blocking仅在生成完整结果后向下游发送数据。
▼ 调度策略 同时调度所有任务(Eager)同时部署作业的所有子任务用于流作业。 上游产生第一条记录部署下游(Lazy)一旦任何生产者生成任何输出就立即部署下游任务。 上游产生完整数据部署下游当任何或所有生产者生成完整数据后部署下游任务。
▼ 数据传输 高吞吐Flink 不是一个一个地发送每条记录而是将若干记录缓冲到其网络缓冲区中并一次性发送它们。这降低了每条记录的发送成本因此提高了吞吐量。 低延迟当网络缓冲区超过一定的时间未被填满时会触发超时发送通过减小超时时间可以通过牺牲一定的吞吐来获取更低的延迟。
我们将在下面深入 Flink 网络协议栈的物理实现时看到关于吞吐延迟的优化。对于这一部分让我们详细说明输出类型与调度策略。首先需要知道的是子任务的输出类型和调度策略是紧密关联的只有两者的一些特定组合才是有效的。
Pipelined 结果是流式输出需要目标 Subtask 正在运行以便接收数据。因此需要在上游 Task 产生数据之前或者产生第一条数据的时候调度下游目标 Task 运行。批处理作业生成有界结果数据而流式处理作业产生无限结果数据。
批处理作业也可能以阻塞方式产生结果具体取决于所使用的算子和连接模式。在这种情况下必须等待上游 Task 先生成完整的结果然后才能调度下游的接收 Task 运行。这能够提高批处理作业的效率并且占用更少的资源。
下表总结了 Task 输出类型以及调度策略的有效组合
注释 [1]目前 Flink 未使用 [2]批处理 / 流计算统一完成后可能适用于流式作业
此外对于具有多个输入的子任务调度以两种方式启动当所有或者任何上游任务产生第一条数据或者产生完整数据时调度任务运行。要调整批处理作业中的输出类型和调度策略可以参考 ExecutionConfig#setExecutionMode()——尤其是 ExecutionMode以及 ExecutionConfig#setDefaultInputDependencyConstraint()。
2.物理数据传输
为了理解物理数据连接请回想一下在 Flink 中不同的任务可以通过 Slotsharing group 共享相同 Slot。TaskManager 还可以提供多个 Slot以允许将同一任务的多个子任务调度到同一个 TaskManager 上。
对于下图所示的示例我们假设 2 个并发为 4 的任务部署在 2 个 TaskManager 上每个 TaskManager 有两个 Slot。TaskManager 1 执行子任务 A.1A.2B.1 和 B.2TaskManager 2 执行子任务 A.3A.4B.3 和 B.4。在 A 和 B 之间是 Shuffle 连接类型比如来自于 A 的 keyBy() 操作在每个 TaskManager 上会有 2x4 个逻辑连接其中一些是本地的另一些是远程的
不同任务远程之间的每个网络连接将在 Flink 的网络堆栈中获得自己的 TCP 通道。但是如果同一任务的不同子任务被调度到同一个 TaskManager则它们与同一个 TaskManager 的网络连接将多路复用并共享同一个 TCP 信道以减少资源使用。在我们的例子中这适用于 A.1→B.3A.1→B.4以及 A.2→B.3 和 A.2→B.4
每个子任务的输出结果称为 ResultPartition每个 ResultPartition 被分成多个单独的 ResultSubpartition- 每个逻辑通道一个。Flink 的网络协议栈在这一点的处理上不再处理单个记录而是将一组序列化的记录填充到网络缓冲区中进行处理。每个子任务本地缓冲区中最多可用 Buffer 数目为每个发送方和接收方各一个:
#channels * buffers-per-channel floating-buffers-per-gate
单个 TaskManager 上的网络层 Buffer 总数通常不需要配置。有关如何在需要时进行配置的详细信息请参阅配置网络缓冲区的文档。
▼ 造成反压1
每当子任务的数据发送缓冲区耗尽时——数据驻留在 Subpartition 的缓冲区队列中或位于更底层的基于 Netty 的网络堆栈内生产者就会被阻塞无法继续发送数据而受到反压。接收端以类似的方式工作Netty 收到任何数据都需要通过网络 Buffer 传递给 Flink。如果相应子任务的网络缓冲区中没有足够可用的网络 BufferFlink 将停止从该通道读取直到 Buffer 可用。这将反压该多路复用上的所有发送子任务因此也限制了其他接收子任务。下图说明了过载的子任务 B.4它会导致多路复用的反压也会导致子任务 B.3 无法接受和处理数据即使是 B.3 还有足够的处理能力。
为了防止这种情况发生Flink 1.5 引入了自己的流量控制机制。
3.Credit-based 流量控制
Credit-based 流量控制可确保发送端已经发送的任何数据接收端都具有足够的能力(Buffer)来接收。新的流量控制机制基于网络缓冲区的可用性作为 Flink 之前机制的自然延伸。每个远程输入通道(RemoteInputChannel)现在都有自己的一组独占缓冲区(Exclusive buffer)而不是只有一个共享的本地缓冲池(LocalBufferPool)。与之前不同本地缓冲池中的缓冲区称为流动缓冲区(Floating buffer)因为它们会在输出通道间流动并且可用于每个输入通道。
数据接收方会将自身的可用 Buffer 作为 Credit 告知数据发送方1 buffer 1 credit。每个 Subpartition 会跟踪下游接收端的 Credit也就是可用于接收数据的 Buffer 数目。只有在相应的通道(Channel)有 Credit 的时候 Flink 才会向更底层的网络协议栈发送数据(以 Buffer 为粒度)并且每发送一个 Buffer 的数据相应的通道上的 Credit 会减 1。除了发送数据本身外数据发送端还会发送相应 Subpartition 中有多少正在排队发送的 Buffer 数称之为 Backlog给下游。数据接收端会利用这一信息(Backlog)去申请合适数量的 Floating buffer 用于接收发送端的数据这可以加快发送端堆积数据的处理。接收端会首先申请和 Backlog 数量相等的 Buffer但可能无法申请到全部甚至一个都申请不到这时接收端会利用已经申请到的 Buffer 进行数据接收并监听是否有新的 Buffer 可用。
Credit-based 的流控使用 Buffers-per-channel 来指定每个 Channel 有多少独占的 Buffer使用 Floating-buffers-per-gate 来指定共享的本地缓冲池(Local buffer pool)大小可选3通过共享本地缓冲池Credit-based 流控可以使用的 Buffer 数目可以达到与原来非 Credit-based 流控同样的大小。这两个参数的默认值是被精心选取的以保证新的 Credit-based 流控在网络健康延迟正常的情况下至少可以达到与原策略相同的吞吐。可以根据实际的网络 RRT (round-trip-time)和带宽对这两个参数进行调整。
注释3如果没有足够的 Buffer 可用则每个缓冲池将获得全局可用 Buffer 的相同份额±1。
▼ 造成反压2
与没有流量控制的接收端反压机制不同Credit 提供了更直接的控制如果接收端的处理速度跟不上最终它的 Credit 会减少成 0此时发送端就不会在向网络中发送数据数据会被序列化到 Buffer 中并缓存在发送端。由于反压只发生在逻辑链路上因此没必要阻断从多路复用的 TCP 连接中读取数据也就不会影响其他的接收者接收和处理数据。
▼ Credit-based 的优势与问题
由于通过 Credit-based 流控机制多路复用中的一个信道不会由于反压阻塞其他逻辑信道因此整体资源利用率会增加。此外通过完全控制正在发送的数据量我们还能够加快 Checkpoint alignment如果没有流量控制通道需要一段时间才能填满网络协议栈的内部缓冲区并表明接收端不再读取数据了。在这段时间里大量的 Buffer 不会被处理。任何 Checkpoint barrier触发 Checkpoint 的消息都必须在这些数据 Buffer 后排队因此必须等到所有这些数据都被处理后才能够触发 Checkpoint“Barrier 不会在数据之前被处理”。
但是来自接收方的附加通告消息向发送端通知 Credit可能会产生一些额外的开销尤其是在使用 SSL 加密信道的场景中。此外单个输入通道( Input channel)不能使用缓冲池中的所有 Buffer因为存在无法共享的 Exclusive buffer。新的流控协议也有可能无法做到立即发送尽可能多的数据如果生成数据的速度快于接收端反馈 Credit 的速度这时则可能增长发送数据的时间。虽然这可能会影响作业的性能但由于其所有优点通常新的流量控制会表现得更好。可能会通过增加单个通道的独占 Buffer 数量这会增大内存开销。然而与先前实现相比总体内存使用可能仍然会降低因为底层的网络协议栈不再需要缓存大量数据因为我们总是可以立即将其传输到 Flink一定会有相应的 Buffer 接收数据。
在使用新的 Credit-based 流量控制时可能还会注意到另一件事由于我们在发送方和接收方之间缓冲较少的数据反压可能会更早的到来。然而这是我们所期望的因为缓存更多数据并没有真正获得任何好处。如果要缓存更多的数据并且保留 Credit-based 流量控制可以考虑通过增加单个输入共享 Buffer 的数量。 注意如果需要关闭 Credit-based 流量控制可以将这个配置添加到 flink-conf.yaml 中taskmanager.network.credit-model:false。但是此参数已过时最终将与非 Credit-based 流控制代码一起删除。 4.序列号与反序列化
下图从上面的扩展了更高级别的视图其中包含网络协议栈及其周围组件的更多详细信息从发送算子发送记录(Record)到接收算子获取它。
在生成 Record 并将其传递出去之后例如通过 Collector#collect()它被传递给 RecordWriterRecordWriter 会将 Java 对象序列化为字节序列最终存储在 Buffer 中按照上面所描述的在网络协议栈中进行处理。RecordWriter 首先使用 SpanningRecordSerializer 将 Record 序列化为灵活的堆上字节数组。然后它尝试将这些字节写入目标网络 Channel 的 Buffer 中。我们将在下面的章节回到这一部分。
在接收方底层网络协议栈Netty将接收到的 Buffer 写入相应的输入通道(Channel)。流任务的线程最终从这些队列中读取并尝试在 RecordReader 的帮助下通过 SpillingAdaptiveSpanningRecordDeserializer 将累积的字节反序列化为 Java 对象。与序列化器类似这个反序列化器还必须处理特殊情况例如跨越多个网络 Buffer 的 Record或者因为记录本身比网络缓冲区大默认情况下为32KB通过 taskmanager.memory.segment-size 设置或者因为序列化 Record 时目标 Buffer 中已经没有足够的剩余空间保存序列化后的字节数据在这种情况下Flink 将使用这些字节空间并继续将其余字节写入新的网络 Buffer 中。
4.1 将网络 Buffer 写入 Netty
在上图中Credit-based 流控制机制实际上位于“Netty Server”和“Netty Client”组件内部RecordWriter 写入的 Buffer 始终以空状态无数据添加到 Subpartition 中然后逐渐向其中填写序列化后的记录。但是 Netty 在什么时候真正的获取并发送这些 Buffer 呢显然不能是 Buffer 中只要有数据就发送因为跨线程写线程与发送线程的数据交换与同步会造成大量的额外开销并且会造成缓存本身失去意义如果是这样的话不如直接将将序列化后的字节发到网络上而不必引入中间的 Buffer。
在 Flink 中有三种情况可以使 Netty 服务端使用发送网络 Buffer
写入 Record 时 Buffer 变满或者Buffer 超时未被发送或发送特殊消息例如 Checkpoint barrier。
▼ 在 Buffer 满后发送
RecordWriter 将 Record 序列化到本地的序列化缓冲区中并将这些序列化后的字节逐渐写入位于相应 Result subpartition 队列中的一个或多个网络 Buffer中。虽然单个 RecordWriter 可以处理多个 Subpartition但每个 Subpartition 只会有一个 RecordWriter 向其写入数据。另一方面Netty 服务端线程会从多个 Result subpartition 中读取并像上面所说的那样将数据写入适当的多路复用信道。这是一个典型的生产者 - 消费者模式网络缓冲区位于生产者与消费者之间如下图所示。在1序列化和2将数据写入 Buffer 之后RecordWriter 会相应地更新缓冲区的写入索引。一旦 Buffer 完全填满RecordWriter 会3为当前 Record 剩余的字节或者下一个 Record 从其本地缓冲池中获取新的 Buffer并将新的 Buffer 添加到相应 Subpartition 的队列中。这将4通知 Netty服务端线程有新的数据可发送如果 Netty 还不知道有可用的数据的话4。每当 Netty 有能力处理这些通知时它将5从队列中获取可用 Buffer 并通过适当的 TCP 通道发送它。
注释4如果队列中有更多已完成的 Buffer我们可以假设 Netty 已经收到通知。
▼ 在 Buffer 超时后发送
为了支持低延迟应用我们不能只等到 Buffer 满了才向下游发送数据。因为可能存在这种情况某种通信信道没有太多数据等到 Buffer 满了在发送会不必要地增加这些少量 Record 的处理延迟。因此Flink 提供了一个定期 Flush 线程(the output flusher)每隔一段时间会将任何缓存的数据全部写出。可以通过 StreamExecutionEnvironment#setBufferTimeout 配置 Flush 的间隔并作为延迟5的上限对于低吞吐量通道。下图显示了它与其他组件的交互方式RecordWriter 如前所述序列化数据并写入网络 Buffer但同时如果 Netty 还不知道有数据可以发送Output flusher 会3,4通知 Netty 服务端线程数据可读类似与上面的“buffer已满”的场景。当 Netty 处理此通知5时它将消费获取并发送Buffer 中的可用数据并更新 Buffer 的读取索引。Buffer 会保留在队列中——从 Netty 服务端对此 Buffer 的任何进一步操作将在下次从读取索引继续读取。
注释5严格来说Output flusher 不提供任何保证——它只向 Netty 发送通知而 Netty 线程会按照能力与意愿进行处理。这也意味着如果存在反压则 Output flusher 是无效的。
▼ 特殊消息后发送
一些特殊的消息如果通过 RecordWriter 发送也会触发立即 Flush 缓存的数据。其中最重要的消息包括 Checkpoint barrier 以及 end-of-partition 事件这些事件应该尽快被发送而不应该等待 Buffer 被填满或者 Output flusher 的下一次 Flush。
▼ 进一步的讨论
与小于 1.5 版本的 Flink 不同请注意a网络 Buffer 现在会被直接放在 Subpartition 的队列中b网络 Buffer 不会在 Flush 之后被关闭。这给我们带来了一些好处
同步开销较少Output flusher 和 RecordWriter 是相互独立的在高负荷情况下Netty 是瓶颈直接的网络瓶颈或反压我们仍然可以在未完成的 Buffer 中填充数据Netty 通知显著减少
但是在低负载情况下可能会出现 CPU 使用率和 TCP 数据包速率的增加。这是因为Flink 将使用任何可用的 CPU 计算能力来尝试维持所需的延迟。一旦负载增加Flink 将通过填充更多的 Buffer 进行自我调整。由于同步开销减少高负载场景不会受到影响甚至可以实现更高的吞吐。
4.2 BufferBuilder 和 BufferConsumer
更深入地了解 Flink 中是如何实现生产者 - 消费者机制需要仔细查看 Flink 1.5 中引入的 BufferBuilder 和 BufferConsumer 类。虽然读取是以 Buffer 为粒度但写入它是按 Record 进行的因此是 Flink 中所有网络通信的核心路径。因此我们需要在任务线程Task thread和 Netty 线程之间实现轻量级连接这意味着尽量小的同步开销。你可以通过查看源代码获取更加详细的信息。
5. 延迟与吞吐
引入网络 Buffer 的目是获得更高的资源利用率和更高的吞吐代价是让 Record 在 Buffer 中等待一段时间。虽然可以通过 Buffer 超时给出此等待时间的上限但可能很想知道有关这两个维度延迟和吞吐之间权衡的更多信息显然无法两者同时兼得。下图显示了不同的 Buffer 超时时间下的吞吐超时时间从 0 开始每个 Record 直接 Flush到 100 毫秒默认值测试在具有 100 个节点每个节点 8 个 Slot 的群集上运行每个节点运行没有业务逻辑的 Task 因此只用于测试网络协议栈的能力。为了进行比较我们还测试了低延迟改进如上所述之前的 Flink 1.4 版本。
如图使用 Flink 1.5即使是非常低的 Buffer 超时例如1ms对于低延迟场景也提供高达超时默认参数100ms75 的最大吞吐但会缓存更少的数据。
6.结论
了解 Result partition批处理和流式计算的不同网络连接以及调度类型Credit-Based 流量控制以及 Flink 网络协议栈内部的工作机理有助于更好的理解网络协议栈相关的参数以及作业的行为。后续我们会推出更多 Flink 网络栈的相关内容并深入更多细节包括运维相关的监控指标(Metrics)进一步的网络调优策略以及需要避免的常见错误等。
原文链接 本文为云栖社区原创内容未经允许不得转载。