游乐游手机版
首页/AI教程/文章详情

Apache Flink网络协议栈原理解析与深度理解

时间:2026-08-15 14:57
Flink网络协议栈连接TaskManager子任务,控制通道基于AkkaRPC,数据通道依赖Netty。逻辑视图包含输出类型与调度策略,通过缓冲和超时机制平衡吞吐与延迟。物理传输复用TCP连接,Credit-based流控防止反压阻塞其他信道,序列化通过RecordWriter与Reader完成。

Flink 的网络协议栈,可以说是整个 flink-runtime 模块的心脏,也是每个 Flink 作业真正跑起来的核心纽带。它连接着所有 TaskManager 的各个子任务(Subtask),因此对吞吐量和延迟这两个关键指标有着直接的影响。控制通道那边,TaskManager 和 JobManager 之间用的是基于 Akka 的 RPC 通信;而 TaskManager 这边,网络协议栈则依赖更底层的 Netty API,走的是另一条路。

这篇文章会先聊聊 Flink 给流算子(Stream operator)提供的高层抽象,然后深入到物理实现和优化细节,最后看看优化效果,以及 Flink 在吞吐和延迟之间是怎么做取舍的。

1. 逻辑视图

Flink 的网络协议栈给彼此通信的子任务提供了这样一个逻辑视图——比如 A 通过 keyBy() 操作进行数据 Shuffle 时,看起来是这样的:

\

这背后,其实依赖三个基本概念:

▼ 子任务输出类型(ResultPartitionType):
Pipelined(有限的或无限的):一旦产生数据就可以持续向下游发送有限或无限的数据流。
Blocking:只有在生成完整结果后,才向下游发送数据。

▼ 调度策略:
同时调度所有任务(Eager):作业的所有子任务一起部署(流作业常用)。
上游产生第一条记录时部署下游(Lazy):一旦有任何生产者生成了输出,就立刻部署下游任务。
上游产生完整数据后部署下游:等所有或部分生产者生成完整数据后,才部署下游任务。

▼ 数据传输:
高吞吐: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 有 2 个 Slot。TaskManager 1 执行子任务 A.1、A.2、B.1、B.2;TaskManager 2 执行子任务 A.3、A.4、B.3、B.4。A 和 B 之间是 Shuffle 连接,比如来自 A 的 keyBy() 操作。在每个 TaskManager 上,会有 2×4 个逻辑连接,有些是本地的,有些是远程的:

\

不同任务(远程)之间的每个网络连接,在 Flink 的网络堆栈里都会获得自己的 TCP 通道。但是,如果同一个任务的不同子任务被调度到同一个 TaskManager,它们与同一个 TaskManager 的网络连接就会多路复用,共享同一个 TCP 信道,这样能减少资源使用。在我们的例子中,这适用于 A.1→B.3、A.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。如果相应子任务的网络缓冲区里没有足够可用的 Buffer,Flink 就会停止从该通道读取,直到 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 数目)。只有对应通道有 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 来指定共享的本地缓冲池大小(可选 3)。通过共享本地缓冲池,Credit-based 流控能使用的 Buffer 总数可以达到和原来非 Credit-based 流控同样的大小。这两个参数的默认值经过精心挑选,保证在网络健康、延迟正常的情况下,至少能达到与原策略相同的吞吐。实际使用时,可以根据网络的 RTT(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 加密信道的场景下。此外,单个输入通道不能使用缓冲池中的所有 Buffer,因为存在无法共享的 Exclusive buffer。新流控协议也有可能无法做到立即发送尽可能多的数据(如果生成数据的速度快于接收端反馈 Credit 的速度),这时可能增长发送数据的时间。虽然这可能会影响作业性能,但考虑到所有优点,新的流量控制通常表现更好。可能有人会通过增加单个通道的独占 Buffer 数量来缓解,但这会增大内存开销。不过,和之前相比,总体内存使用可能仍然会降低,因为底层的网络协议栈不再需要缓存大量数据——我们总能立即把数据传给 Flink(一定有相应的 Buffer 来接收)。

使用新的 Credit-based 流量控制时,可能还会注意到另一件事:由于我们在发送方和接收方之间缓冲的数据变少了,反压可能会更早到来。但这是预期的效果,因为缓存更多数据并没有真正获得任何好处。如果真想缓存更多数据同时保留 Credit-based 流控,可以考虑增加单个输入共享 Buffer 的数量。

\

\

4. 序列化与反序列化

下图从上面的高层视图进一步扩展,展示了网络协议栈及其周围组件的更多细节,从发送算子发送 Record 到接收算子获取它:

\

当 Record 生成并传递出去之后——比如通过 Collector#collect()——它被传递给 RecordWriter。RecordWriter 会把 Ja va 对象序列化成字节序列,最终存储在 Buffer 中,再按照上面描述的方式在网络协议栈中处理。RecordWriter 首先使用 SpanningRecordSerializer 将 Record 序列化为灵活的堆上字节数组。然后,它尝试把这些字节写入目标网络 Channel 的 Buffer 中。我们会在下面回到这一步。

在接收方,底层网络协议栈(Netty)把接收到的 Buffer 写入相应的输入通道。流任务的线程最终从这些队列中读取,并尝试在 RecordReader 的帮助下,通过 SpillingAdaptiveSpanningRecordDeserializer 将累积的字节反序列化为 Ja va 对象。和序列化器类似,这个反序列化器也必须处理特殊情况,比如跨越多个网络 Buffer 的 Record,或者 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 序列化到本地的序列化缓冲区中,再把序列化后的字节逐渐写入对应 ResultSubpartition 队列中的一个或多个网络 Buffer。虽然单个 RecordWriter 可以处理多个 Subpartition,但每个 Subpartition 只会有一个 RecordWriter 向它写入数据。另一方面,Netty 服务端线程会从多个 ResultSubpartition 中读取,然后像上面说的那样把数据写入适当的多路复用信道。这是一个典型的生产者 - 消费者模式,网络缓冲区位于两者之间,如下图所示。在(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。

▼ 进一步的讨论

和 Flink 1.5 之前的版本不同,请注意:(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 中引入的 BufferBuilderBufferConsumer 类。虽然读取是以 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)(对于低延迟场景)也提供了高达超时默认参数(100ms)75% 的最大吞吐,但缓存的数据更少。

6. 结论

理解 Result partition、批处理和流式计算的不同网络连接以及调度类型、Credit-Based 流量控制,以及 Flink 网络协议栈内部的工作机理,有助于更好地理解网络协议栈相关的参数以及作业的行为。后续我们会推出更多关于 Flink 网络栈的内容,并深入更多细节,包括运维相关的监控指标(Metrics)、进一步的网络调优策略,以及需要避免的常见错误等。

来源:https://developer.aliyun.com/article/706462
上一篇AI回答结果解析优化方法:品牌别名库与消歧规则实践 下一篇AI搜索内容策略入门到行业实战完整指南
本站内容用于信息整理与展示,如有侵权或内容问题请及时联系处理。

相关推荐

补充同频道和同主题内容,方便继续浏览更多相关内容。

同类最新

继续查看同栏目最近更新的文章。

更多
CAD零基础入门教程:坐标输入、图层管理与基础绘图命令
AI教程 · 2026-09-01

CAD零基础入门教程:坐标输入、图层管理与基础绘图命令

本文面向CAD零基础学习者,系统讲解坐标输入、图层管理与基础绘图命令的核心用法。通过分步实操与常见问题排查,帮助新手建立精确绘图习惯,掌握规范出图的基础能力。

CAD从入门到项目交付:绘图、标注、图块与实战工作流
AI教程 · 2026-09-01

CAD从入门到项目交付:绘图、标注、图块与实战工作流

掌握CAD的核心在于建立“画得准、标得清、复用快、交付稳”的工作流。本文提供从环境设置、高频命令组合、标注规范、图块标准化到项目分阶段交付的完整路径,帮助初学者避免常见返工陷阱,独立完成可检查、可复用、可打印的工程图纸。

Claude Code 登录指南:个人、Teams 与企业账号区分与授权步骤
AI教程 · 2026-09-01

Claude Code 登录指南:个人、Teams 与企业账号区分与授权步骤

本文详细解析 Claude Code 登录前的账号类型区分方法,涵盖个人订阅、Teams 席位与企业 Enterprise 席位的授权路径差异。提供终端登录命令、环境变量排查及常见异常处理步骤,帮助用户快速完成正确授权并避免登录路径混淆。

Claude Code 文件修改前的权限模式配置与命令审批指南
AI教程 · 2026-09-01

Claude Code 文件修改前的权限模式配置与命令审批指南

本文详细介绍Claude Code在修改文件前的权限模式配置方法,包括defaultMode可选值、permissions allow与deny规则设置、多层级配置文件管理以及 status验证技巧,帮助开发者安全高效地使用AI编程助手。

Claude Code接入VS Code后先测扩展和终端命令
AI教程 · 2026-09-01

Claude Code接入VS Code后先测扩展和终端命令

在VS Code中接入Claude Code后,建议优先验证扩展面板与集成终端两条入口。本文提供标准检查顺序、关键命令与常见故障排查路径,帮助你快速确认环境就绪,避免后续开发受阻。