Flink 的作业提交流程,就像一座自动化工厂的运转全景: 从计划处(CLI)下达任务,到厂长(JobManager)组织生产,后勤(ResourceManager)分配资源,车间(TaskManager)开动机器,工位(Slot)全速生产——最终让数据在这座工厂中实时流动、加工、输出。

一、从 WordCount 说起

上一篇文章我们介绍了《从论文到生产:Flink 的诞生、思想与架构哲学》,现在让咱们带着“思想地图”一起来研究下 Flink 源码。编程语言的入门示例是 “Hello World”,而在大数据计算领域,对应的经典例子则是 WordCount。在学习 Flink 的过程中,官网通常会让我们执行以下命令来提交这个最简单的 demo 并在本地试运行:

$ ./bin/start-cluster.sh    //启动集群
$ ./bin/flink run examples/streaming/WordCount.jar  //提交demo
$ ./bin/stop-cluster.sh     //停止集群

同时当我们想要在分布式集群 Yarn/K8s 跑的时候执行的又是诸如这些命令:

// 运行在Yarn,官网参考命令
$ ./bin/flink run-application -t yarn-application ./xx/my-flink-job.jar

// 运行在K8s,官网参考命令
$ ./bin/flink run-application \\
    --target kubernetes-application \\
    -Dkubernetes.cluster-id=my-first-application-cluster \\
    -Dkubernetes.container.image.ref=custom-image-name \\
    local:///opt/flink/usrlib/my-flink-job.jar

Flink源码工程中,我们可以找到 WordCount 的代码:flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/wordcount/WordCount.java

主要核心逻辑像这样:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<String> text = env.readTextFile("input.txt");
text.flatMap(...).keyBy(...).sum(...).print();
env.execute("WordCount");

学到这里的时候大家会不会有疑问:执行这些命令后 Flink 内部到底发生了什么?是谁接管了作业?又是如何分发到集群中执行的?带着这些问题,我们一起来看看 Flink 内部到底做了哪些事。

二、运行模式与背景

在 Flink 中,作业可以提交到不同类型的集群(如 Standalone、YARN、K8s),并支持多种运行模式:Session、Per-Job(1.15弃用,2.0移除)、Application。不同模式主要影响的是集群生命周期资源管理方式,而整体流程的核心逻辑一致,本文将以 Application on K8s 为例,追踪从执行命令到env.execute() 再到作业执行的完整流程。

Application on K8s模式 下,Flink 将作业与集群捆绑在一起:一个 Job 有一个专用集群实例,JobManager 与 TaskManager 由 Kubernetes 动态创建,这种模式因资源隔离好、生命周期绑定方便,所以在生产中应用广泛。

注:本文源码解析基于 Flink 1.20 版本

三、宏观全景:一次作业是如何提交的?

我们先从一张图来整体看看一次作业是如何提交的

从图中可以看到,整个提交流程可分为 10 个步骤,和K8s相关的大家忽略,主要关注 Flink 的整体运行逻辑即可:

  1. 用户执行命令并解析后上传相关所需资源(用的 demo 包含在镜像,所以无需上传)

  2. Flink KubeClient 请求集群创建 Flink JM deployment 并启动 JM 的 pod

  3. 从 K8s 专用 Entrypoint 类进入,创建 JobManager 关键角色:Dispatcher(创建会触发执行用户代码生成 StreamGraph 和 JobGraph)、Resource Manager 和 JobMaster

  4. JobMaster 创建会生成 ExecutionGraph,创建 SlotPool 并计算所需资源,向 RM 注册并请求 Slot

  5. RM 发现 slot 资源不够会向 K8s 申请新的资源

  6. 在新申请的 Pod 中通过 KubernetesTaskExecutorRunner 启动 TaskExecutor(即 TaskManager)

  7. Task Executor 将自身的 Slot 向 RM 进行注册

  8. RM 对已注册的 Slot 进行分配

  9. Task Executor 向 JM 提供自身 Slot

  10. JM 将 Task 提交部署到 Task Executor 运行

四、微观拆解:提交到执行全流程详解

1. CLI 入口:run-application 命令背后发生了什么?

当我们执行 ./bin/flink run-application --target kubernetes-application -D...... local:///opt/flink/usrlib/my-job.jar 命令时,首先进入的是 Flink 的命令行客户端 CliFrontend。可以把它理解成 Flink 的“前台接待员”——负责解析参数、识别命令、启动正确的执行流程。整个调用链大致如下:

  1. CliFrontend 入口:flink脚本会进入 CliFrontend.main() 方法,解析命令行参数。
  2. 命令分发:识别到 run-application 命令,调用 runApplication() 方法。
  3. 参数处理:解析命令行参数,创建 ProgramOptions 和 Configuration 对象。
  4. 部署器创建:实例化 ApplicationClusterDeployer,用于部署应用集群。
  5. 资源上传:通过 artifactUploader.uploadAll() 判断按需上传本地资源(JAR文件等)到远程存储

此时:Flink 已经把用户命令转译成一个可执行的“集群部署请求”。那这一阶段的核心目标就是“准备运行环境”——CLI 把所有参数和依赖打包成 Flink 可识别的任务描述。

2. 集群创建:K8s 如何拉起 JobManager?

接下来,ApplicationClusterDeployer 会调用KubernetesClusterDescriptor.deployApplicationCluster() ,进入集群创建阶段:

  1. 验证配置:确保 deployment.target 设置为 kubernetes-application
  2. 应用配置:将应用配置合并到Flink配置中。
  3. 集群创建:调用 deployClusterInternal() 创建集群:
    • 设置执行模式和入口点类
    • 创建 KubernetesJobManagerParameters 和 FlinkPod 模板
    • 构建 KubernetesJobManagerSpecification
    • 调用 client.createJobManagerComponent() 创建Deployment、ConfigMap和Service
    private ClusterClientProvider<String> deployClusterInternal(...){
        // ...
		    final KubernetesJobManagerSpecification kubernetesJobManagerSpec =
                KubernetesJobManagerFactory.buildKubernetesJobManagerSpecification(
                  podTemplate, kubernetesJobManagerParameters);

        client.createJobManagerComponent(kubernetesJobManagerSpec);
        // ...
    }

执行到这里,Kubernetes 会启动一个JobManager Pod。容器启动命令会指定入口类为KubernetesApplicationClusterEntrypoint

此时:JobManager Pod 已在集群中启动,等待执行用户程序。

3. JobManager 启动:三大核心角色登场

JobManager 是 Flink 集群的“大脑”,它启动后会创建三个关键组件:Dispatcher、ResourceManager和JobMaster,启动大致流程如下:

  1. JM Pod启动:Kubernetes根据Deployment创建JM Pod,执行 KubernetesApplicationClusterEntrypoint.main()
  2. ClusterEntrypoint初始化:初始化配置、日志、安全等。
  3. 组件工厂创建:调用 createDispatcherResourceManagerComponentFactory() 创建 DefaultDispatcherResourceManagerComponentFactory
  4. Dispatcher创建
    • 使用 ApplicationDispatcherLeaderProcessFactoryFactory 创建调度器工厂
    • ApplicationDispatcherBootstrap中,通过runApplicationEntryPoint()调用ClientUtils.executeProgram执行用户代码
    • 用户代码最后都必须有env.execute(),执行后会生成StreamGraph
    • 然后在 EmbeddedExecutor 执行器中调用PipelineExecutorUtils.getJobGraph方法,该方法使用FlinkPipelineTranslationUtil.getJobGraphStreamGraph转换为JobGraphStreamGraphTranslator通过StreamingJobGraphGenerator完成具体的转换逻辑
    • JobGraph包含了实际执行所需的所有信息,如作业配置、并行度、资源需求等
  5. ResourceManager创建:创建KubernetesResourceManager,使用 KubernetesResourceManagerDriver 管理资源。
  6. JobMaster创建:当作业提交后,由Dispatcher创建JobMaster。

此时:Flink 集群的中枢已经准备完毕。Dispatcher 拿到 JobGraph,ResourceManager 准备调度资源,JobMaster 即将开始调度任务。

4.JM 触发生成 ExecutionGraph 并请求资源

JobMaster 是真正掌控作业执行的核心,在 JobMaster 构造函数中创建 SchedulerBase 实例时,会调用 createAndRestoreExecutionGraph 方法触发 JobGraph 转化为 ExecutionGraph —— Flink 的物理执行计划,核心调用链如下:

SchedulerBase.createAndRestoreExecutionGraph()
 → ExecutionGraphFactory.createAndRestoreExecutionGraph()
 → DefaultExecutionGraphBuilder.buildGraph()
  • 核心转换逻辑:在DefaultExecutionGraphBuilder
// DefaultExecutionGraphBuilder.java
ExecutionGraph executionGraph = DefaultExecutionGraphBuilder.buildGraph(jobGraph, ...);

ExecutionGraph描述了任务的实际执行单元(ExecutionVertex)和依赖关系。

构建完成后,JobMaster 会向 ResourceManager 请求 Slot 资源。

此时:JobMaster 已经画好了调度图纸,等待 ResourceManager 分配资源。

5.RM 申请资源启动 TaskManager

ResourceManager 是 Flink 的资源协调者,在 Kubernetes 模式下,它通过 KubernetesResourceManagerDriver 与 API Server 交互,请求新的 TM Pod:

  1. 资源请求处理:当ResourceManager需要资源时,调用 requestResource()
  2. Pod创建
    • 使用 KubernetesTaskManagerFactory 构建TaskManager Pod
    • 设置Pod名称格式为 {clusterId}-taskmanager-{attemptId}-{podIndex}
    • 调用 client.createTaskManagerPod() 创建Pod
  3. 资源跟踪:维护 requestResourceFutures 跟踪资源请求状态。

此时:Kubernetes 开始调度 TaskManager Pod;Flink 正在等待节点启动并注册。

6. Task Executor 启动

KubernetesTaskExecutorRunner 负责启动Task Executor:

  1. 入口点执行:TaskManager Pod 启动后,执行/docker-entrypoint.sh taskmanager,最终调用KubernetesTaskExecutorRunner.main()方法。
  2. 配置加载:从环境变量获取并设置 TASK_MANAGER_NODE_ID
  3. TaskManager启动:调用 TaskManagerRunner.runTaskManagerProcessSecurely() 启动TaskExecutor。
public static void main(String[] args) {
    EnvironmentInformation.logEnvironmentInfo(LOG, "Kubernetes TaskExecutor runner", args);
    SignalHandler.register(LOG);
    JvmShutdownSafeguard.installAsShutdownHook(LOG);

    runTaskManagerSecurely(args);
}

private static void runTaskManagerSecurely(String[] args) {
    Configuration configuration = null;

    try {
        configuration = TaskManagerRunner.loadConfiguration(args);
        final String nodeId = System.getenv().get(Constants.ENV_FLINK_POD_NODE_ID);
        // ...
        configuration.set(TaskManagerOptionsInternal.TASK_MANAGER_NODE_ID, nodeId);
    } catch (FlinkParseException fpe) {
        // ...
    }

    TaskManagerRunner.runTaskManagerProcessSecurely(checkNotNull(configuration));
}

此时:TaskExecutor 已经启动,准备向集群注册自己。

7. Task Executor 注册到 RM

TaskExecutor 启动后,会主动与 ResourceManager 建立连接并注册。

注册过程相当于 “我是谁、我有几个 Slot、请给我任务”。

  1. 连接RM:TaskExecutor启动后,在connectToResourceManager()方法中连接RM
  2. 注册TM:通过远程调用最终会通过RM的registerTaskExecutor()方法把TM进行注册

注册信息包含TaskExecutor的资源ID、可用slot数量等

// ResourceManager接收TaskExecutor注册请求
public CompletableFuture<RegistrationResponse> registerTaskExecutor(
        final TaskExecutorRegistration taskExecutorRegistration, final Time timeout) {
    // 建立与TaskExecutor的RPC连接
    CompletableFuture<TaskExecutorGateway> taskExecutorGatewayFuture = 
            getRpcService().connect(
                    taskExecutorRegistration.getTaskExecutorAddress(), 
                    TaskExecutorGateway.class);
    // ...
    
    return taskExecutorGatewayFuture.handleAsync((gateway, throwable) -> {
        if (throwable != null) {
            return new RegistrationResponse.Failure(throwable);
        } else {
            return registerTaskExecutorInternal(gateway, taskExecutorRegistration);
        } // ...
    }, getMainThreadExecutor());
}

此时:资源池(Slot)已准备就绪。

8. RM 对已注册的 Slot 进行分配

  1. 需求处理:ResourceManager 通过 SlotManager 接收并处理 JobMaste r的资源需求
  2. 资源跟踪:更新作业的资源需求信息
  3. 分配触发:调用 checkResourceRequirementsWithDelay()触发 slot 分配逻辑
  4. 分配策略:根据作业需求和可用资源进行 slot 匹配和分配
// FineGrainedSlotManager处理作业资源需求
public void processResourceRequirements(ResourceRequirements resourceRequirements) {
    checkInit();
    if (resourceRequirements.getResourceRequirements().isEmpty()
            && resourceTracker.isRequirementEmpty(resourceRequirements.getJobId())) {
        return; // 跳过空需求
    }
    // ...
    // 更新资源需求并触发分配检查
    resourceTracker.notifyResourceRequirements(
            resourceRequirements.getJobId(), resourceRequirements.getResourceRequirements());
    checkResourceRequirementsWithDelay();
}

此时:Flink 内部的资源调度链闭合,准备进入任务部署阶段。

9. Task Executor 向 JM 提供自身 Slot

TaskExecutor 收到分配指令后,会向 JobMaster 提交自身可用的 Slot 信息:

// TaskExecutor向JobMaster提供slot
private void internalOfferSlotsToJobManager(JobTable.Connection jobManagerConnection) {
    final JobID jobId = jobManagerConnection.getJobId();

    if (taskSlotTable.hasAllocatedSlots(jobId)) {
        final JobMasterGateway jobMasterGateway = jobManagerConnection.getJobManagerGateway();
        final JobMasterId jobMasterId = jobManagerConnection.getJobMasterId();

        // 收集已分配的slot
        final Iterator<TaskSlot<Task>> reservedSlotsIterator =
                taskSlotTable.getAllocatedSlots(jobId);
        final Collection<SlotOffer> reservedSlots =
                CollectionUtil.newHashSetWithExpectedSize(2);

        while (reservedSlotsIterator.hasNext()) {
            SlotOffer offer = reservedSlotsIterator.next().generateSlotOffer();
            reservedSlots.add(offer);
        }

        // 发送slot offer
        final UUID slotOfferId = UUID.randomUUID();
        currentSlotOfferPerJob.put(jobId, slotOfferId);

        CompletableFuture<Collection<SlotOffer>> acceptedSlotsFuture =
                jobMasterGateway.offerSlots(
                        getResourceID(),
                        reservedSlots,
                        Time.fromDuration(taskManagerConfiguration.getRpcTimeout()));

		        // 处理接受结果 ...
    }
}

JobMaster 接收后,确认 Slot 有效并将其绑定到对应的任务。

Slot 是 Flink 中“Task运行的容器”,这一过程相当于 TaskExecutor 告诉 JobMaster:‘我有空位,快派任务过来。’

10. JM 将 Task 提交部署到 Task Executor 运行

最后,JobMaster 根据 ExecutionGraph 将具体Task下发到对应的 TaskExecutor:

// TaskExecutor接收任务提交请求
public CompletableFuture<Acknowledge> submitTask(
        TaskDeploymentDescriptor tdd, JobMasterId jobMasterId, Time timeout) {
    final JobID jobId = tdd.getJobId();

    // 验证JobManager连接、JobMaster ID、slot状态 ...

    // 加载任务数据
    try {
        tdd.loadBigData(
                taskExecutorBlobService.getPermanentBlobService(),
                jobInformationCache,
                taskInformationCache,
                shuffleDescriptorsCache);
    } catch (IOException | ClassNotFoundException e) {
        throw new TaskSubmissionException(
                "Could not re-integrate offloaded TaskDeploymentDescriptor data.", e);
    }
    // ...

    // 创建并启动Task实例...
    if (taskAdded) {
        task.startTaskThread();
        // ...
    } else // ...
}

TaskExecutor 收到任务后会创建 Task 实例,加载用户代码并启动执行线程。

至此,一个 Flink 作业在 Kubernetes 上的部署过程正式完成。

这一刻,WordCount 的算子真正开始在 TaskManager 中运行。

五、源码结构小结—化繁为简

看完前面复杂的源码与组件调用,我们不妨退一步,从更高的视角重新审视整个 Flink 提交流程。

无论是 Local、YARN 还是 Kubernetes 模式,核心主线其实都围绕这五大步骤展开:

CLI → 集群资源管理API → JobManager → ResourceManager → TaskManager
解析命令   →  创建集群   →   构建作业图   →   分配资源   →   部署任务

其中:

1.资源提供者既可以是 K8s,也可以是 Yarn 等,而源码实现只是对应不同资源提供者,客户端和交互方式有所差异,核心流程基本一致。

2.Session和Per-Job两种模式下 StreamGraph/JobGraph 在Client 生成,Application 则是都在Cluster端的 JM 中生成。

Flink 的作业提交流程,就像一座自动化工厂的运转全景: 从计划处(CLI)下达任务,到厂长(JobManager)组织生产,后勤(ResourceManager)分配资源,车间(TaskManager)开动机器,工位(Slot)全速生产——最终让数据在这座工厂中实时流动、加工、输出。

六、延伸思考

至此,我们从 flink run-application 命令一路追踪到 TaskExecutor 真正执行算子任务。

这一过程揭示了 Flink 在 Kubernetes 上的完整生命周期闭环:提交 → 启动 → 调度 → 执行。

在分布式集群中,Flink 的 JM、TM 都运行在不同的机器上,那他们是如何互相通知、调用和传输数据的呢?


下一篇,我们就来聊聊 Flink 内部组件的通信机制——厂长(JobManager)和车间主任(TaskManager)的“无线对讲机”

更多推荐