Flink集群架构
Flink集群属于Master-Worker的架构模式,其主要包含JobManager、TaskManager以及Client三大组件:
- JobManager: 又称JM,是集群的管理节点,负责整个集群的计算资源、Job的调度与执行、以及CheckPoint的协调。
- TaskManager: 又称TM,集群的工作节点,每个TM一般部署在不同的节点上,JM会将整个Job切分成多个Task,然后分别发送到不同的TM上运行计算;
- Client: 负责接收用户提交的应用程序,然后在本地解析生成一个JobGraph对象(即Job),再通过RPC的方式将JobGraph对象提交到JobManager上,如果提交成功,JM会给Client端返回一个JobClient,然后Client通过JobClient可以和JM进行通信,获取Job的执行状态。
在任务执行中,Flink会对Task的状态进行管理和监控,以便在出现故障或异常情况时进行恢复或重试。
组件 | JobManager架构

Checkpoint Coordinator
- 指的是TM里Checkpoint的协调工作,Checkpoint会通过文件检查点的方式来记录一些状态性的数据,主要和容错相关,而JM会负责TM里面的Checkpoint的协调执行。
JobGraph -> Execution Graph
- Client端会在内部生成一个JobGraph对象给JM,JM收到JobGraph时又会生成一个Execution Graph(物理执行逻辑图),类似与SQL语句和SQL执行计划,然后JM再将Execution Graph 拆分成多个执行单元(Task)发送到TM上。
RPC通信(Actor System)
- JM和TM,以及和Client都是通过Akka来实现RPC通信,而Akka里面最核心的组件就是Actor System,它负责两个进程之间的通信。
Job接收(Job Dispatch)
- Client会提交不同的JobGraph,所以JM会接收不同的JobGraph,然后在拆分成多个Task分发到不同的TM上,这一过程涉及到Job的分发,而JM内部有一个Job Dispatch 组件专门负责干这件事。
ResourceManager(集群资源管理)
- JM内部有一个集群资源管理器,JM会负责整个集群的资源管理,对于不同的部署模式会有对应不同实现的资源管理器,比如 On Standalone、On Kubernetes、On Yarn 等等。
TaskManager的注册与管理
- JM会负责TM的注册,当启动一个TM时,它会主动和JM建立RPC连接,这个过程中JM会将对应的TM的信息保存在本地,之后会对TM进行心跳检测。
组件 | TaskManager架构

JM会将Task发送到TM上,TM(整个进程)内部会有一个Task Scheduling(线程池),当JM提交到Task到TM时,就会从线程池中取一个线程来负责执行任务,而这个执行任务的线程就是Task Slot(任务资源槽),然后是Data Exchange,可以理解为Spark中的shuffle操作,在Flink 中,如果数据进行GroupByKey操作,就会涉及到数据的交互,在TM中就是通过Shuffle Environment来支持Shuffle 操作。Shuffle操作意味着数据要跨节点传输,这种传输时通过RPC的方式,通过Network Manager组件提供的基于Netty实现的网络通信站。
在TM中有一个比较重要的Memory Management,用于管理内存。
TM在启动之后会向JM进行注册,TM要知道当前集群中都有哪些TM以及相关的状态。
组件 | Client架构

用户可以使用Flink API(Java/python/Scala/SQL)编写Flink程序,然后发送给Client, 再由Client提交到集群。客户端在接收应用程序时,会在本地启动一个相应的Client进程,该进程负责解析用户提交的应用程序,解析时会将应用程序的main方法拿出来在自己的进程里执行,执行的主要目的是生成对应的JobGraph对象。
Client中有几个核心概念:
- Context Environment: 表示上下文环境,即Client在第一步创建Context Environment,然后在其中执行main方法。
- Job Submit: 当生成JobGraph对象时,会将它和依赖的一些Jar包一块提交到JM上。这个过程也是RPC通信。
对象 | JobGraph架构

JobGraph本质上是用户编写的应用程序的一个DAG。无论用户采用哪种API编写应用程序,最终都会解析编译成一个Jar包或SQL脚本,再调用flink run命令来执行。
其核心逻辑是,通过反射的方式调用应用程序的main方法,对应Application Code的执行,然后调用应用程序的Execute(Client端)方法,将应用程序转成StreamGraph(Dataflow),接着讲StreamGraph转换成JobGraph,此时会对每一个算子进行拆解、指定相应的并行度,最后调用Submit将生成的JobGraph提交到JM中,而DAG里的算子执行,则依赖JM将任务调度到TM中。
Flink集群的运行模式
Per-Job模式在1.17中被弃用
- 每一个Job对应一个集群,每个Job独占一个JM
Session Mode

主要运行流程:
- 下载应用依赖的Jar包,以及安装Client
- 在Client中执行Main方法生成JobGraph对象
- 将JobGraph和依赖一起提交到集群运行
- 等待Job运行结果
主要问题是,在第三步中,上传资源过程很吃带宽,并且Job每次运行相关依赖都需要每次上传。
主要实现逻辑:
所有的Job都在一个Runtime中运行。JM和多个TM组成一个集群,所有的Job都会发送到同一个JM中,其特点就是多个Job共享一个JM,对于Client来说,会执行一些方法生成相应的JobGraph对象,然后连通依赖的Jar包一起传到JM中。此模式JM的生命周期不受Job影响,不管提交多少个Job,JM都会一直处于运行状态。
优缺点:
- 优点:资源充分共享,利用率高
- 缺点:资源隔离差,对于非Native类型部署时,TM的扩展差,Slot的计算伸缩性差。
Application Mode
Session 和Per-Job的痛点:首先生成 JobGraph 是需要消耗 CPU 资源的,任务多的话会导致客户端压力增大。而且生成 Jobgraph 是同步的,任务一多的话也可能会造成阻塞,加上 JobGraph 和依赖的 Jar 包的提交也需要时间,如果 Jar 包比较大的话,会非常依赖带宽,更何况这些依赖每次都要上传。而对于流平台而言重要的就是实时,不能陷入等待。因此后来就有人提出,能不能把生成 JobGraph 这一步从客户端移到 JM 上,这样可以释放本地客户端的压力。最终在客户端里,只需要负责命令的提交、下发,以及等待 Job 的运行结果即可。
为了解决Session 和Per-Job模式存在的问题而诞生的新模式。

其模式的主要实现思路是将生成JobGraph过程移到了JM中执行,JM首先去分布式文件系统拉取存储的依赖包,再根据应用程序生成JobGraph对象,然后执行、调度。客户端不需要再将Jar包提交到JM中,从而大大避免了网络传输的消耗,以及Client的负载,与此同时也是实现了资源隔离,所有的Job共享一个JM,但JM内部又在Application层实现了资源隔离。
- Application Mode的隔离属于容器层面隔离,提高资源的利用率。
- Per-Job Mode的隔离属于虚拟机层面的隔离,每个Job对应的JM可以看出单独的虚拟机,资源浪费明显。
Flink集群资源管理器
Native 集群部署: 当以 Session 模式启动集群时,只启动 JM 不会启动 TM。只有在客户端提交 Job 之后,JM 才会跟资源管理器进行交互、申请资源,动态启动 TM 满足计算要求。
YARN 和K8S都支持Native部署模式。
目前主要支持的资源管理器有:
- Local
- Standalone
- YARN
- Kubernetes
Flink on Standalone

在此模式中,TM必须要事先进行注册,这就意味着能够计算的资源在一开始就已经确定,然后JM 和TM运行的时候也会有相应的进程,如果JM和TM在同一个节点,那就是单机Standalone模式,如果在多个节点,则是多机Standalone模式
注意:如果资源管理器是Standalone,那Flink的运行只支持Session
Flink on YARN

Resource Manager:
- 处理Client请求。Client若提交任务,需要经过Resource Manager, 他是整个资源的管理者,包括集群的CPU、内存、磁盘等
- 监控Node Manager
- 启动或监控Application Master
- 资源分配和调度
Node Manager:
- 管理单个节点上的资源。是当前节点的资源管理者,同时也需要和Resource Manager汇报。
- 处理来自Resource Manager的命令。
- 处理来自Application Master的命令。
Application Master:
- 某个任务的管理者。当任务在Node Manager上运行时,由Application Master负责管理,因为每个任务对应一个AM.
- 负责数据切分。
- 为应用程序申请资源并分配给内部的任务
- 任务的监控与容错。
Container:
- 是YARN中资源的抽象,它封装了节点上的多维度资源,如内存、磁盘、CPU、网络等,Container是为AM服务的,任务在运行时,需要的内存、CPU等资源都被虚拟化到了Container中。
Flink的三种运行模式,YARN都是支持的。
Pyflink架构和运行原理
Pyflink整体架构

Python API 是建立在Java API之上的。在Client,会起一个Python VM,然后启动一个Java VM,两个VM进行socket 通信。

- Python 和Java VM 里面都用Py4J各自启动一个Gateway,然后Gateway会维护一些对象。
- 在 Python 这边创建一个 table 对象的时候,它也会在相应的 Java 这边创建一个相同 table 对象。如果创建一个 TableEnvironment 对象的时候,在 Java 这边也会创建一个 TableEnvironment 对象,当调用 table 对象上的方法,那么也会映射到 Java 这边
- 基于这一套架构,可以得出一个结论:如果用 Python Table API 写出了一个作业,这个作业没有 Python UDF 的时候,那么这个作业的性能跟用 Java 写出来的作业性能是一样的。因为它底层的架构都是同一套 Java 的架构。
Pyflink UDF架构

主要流程如下:
- 在open方法中进行Java Operator 和Python Operator环境的初始化。
- 环境初始化好后,会进行数据处理。当Java Operator 收到数据后,先把数据放到一个input buffer缓冲区中,达到一定阈值后,才会flash到Python这边,Python这边处理完后,也会先将数据放到一个结果的缓冲区中,当达到一定阈值,比如达到一定的记录的行数,或一定的时间位置,才会把结果flash到java这边。
- state 访问的链路。
- logging 访问的链路。
- metrics 汇报的链路。