Flink 架构入门:从 JobManager、TaskManager 到并行度的核心机制拆解

发布时间:2026/10/7 14:32:39
Flink 架构入门:从 JobManager、TaskManager 到并行度的核心机制拆解 1. 从一次本地作业提交说起JobManager、TaskManager 与并行度到底怎么配合刚接触 Flink 的时候最容易懵的不是算子怎么写而是「我提交了一个作业它到底跑在哪、谁在管、为什么并行度设错了就报错」。这篇就围绕 Flink 架构里最核心的三个概念——JobManager、TaskManager、并行度parallelism——用一个本地 Standalone 集群把整条链路走一遍。你不需要先懂 YARN 或 K8s只要一台能跑 Java 的机器就能把作业提交、任务调度、Slot 分配这几件事看清楚。Flink 运行时是一个典型的 Master-Slave 结构。JobManager 是 Master负责整个集群的资源管理、任务调度、Checkpoint 协调和故障恢复TaskManager 是 Slave负责真正执行任务并管理自己节点上的资源。Client 负责把作业提交给 JobManager提交方式可以是命令行也可以是 Web UI。提交之后Client 可以断开detached 模式也可以保持连接接收任务报告attached 模式。这里有个关键点并行度不是「越多越快」的开关而是决定一个算子被拆成多少个子任务subtask去并行执行的参数。一个算子的 subtask 数量就是它的并行度。整个流程序的并行度通常取所有算子中最大的那个。而 TaskManager 能承载多少 subtask取决于它有多少个 Task Slot。Slot 是 TaskManager 资源的固定子集比如一个 TaskManager 有 3 个 slot它就把管理的内存大致分成三份。注意slot 目前只隔离受管理的内存不隔离 CPU。所以三者关系可以这样理解JobManager 拿到作业后根据并行度算出需要多少个 subtask再向 ResourceManager 申请对应数量的 SlotTaskManager 提供 Slotsubtask 被分配到 Slot 里执行。如果并行度大于集群总 Slot 数作业就会因为资源不足而失败。这也是为什么很多人第一次跑 Flink 会看到「No enough slots」之类的报错。下面我会先讲清楚 JobManager 内部的三个组件再给出可复制的flink-conf.yaml配置然后实际提交一个作业最后验证 TaskManager 注册和 Slot 使用情况。整个过程都在本地 Standalone 模式下完成适合刚入门的开发者跟着做一遍。2. JobManager 内部拆解与 Standalone 集群前置准备JobManager 本身是一个 JVM 进程但它内部并不是铁板一块而是由三个组件协作ResourceManager、Dispatcher、JobMaster。理解这三个组件才能明白作业提交后到底发生了什么。ResourceManager 负责管理 TaskManager 的 Slot。在 Standalone 部署下它只能分配已有 TaskManager 的空闲 Slot不能自己启动新的 TaskManager。当 JobMaster 申请 Slot 时ResourceManager 会把有空闲 Slot 的 TaskManager 分配出去。如果 Slot 不够它就没法满足请求作业会等待或失败。它还负责终止空闲的 TaskManager 来释放资源。Dispatcher 提供 REST 接口接收作业提交并为每个提交的作业启动一个新的 JobMaster。它还启动了 Flink WebUI用来展示作业执行信息。Dispatcher 在某些提交方式下不是必需的但在 Standalone 会话集群里它是核心入口。JobMaster 负责管理单个 JobGraph 的执行。一个集群里可以同时跑多个作业每个作业有自己的 JobMaster。JobMaster 会把 JobGraph 转换成 ExecutionGraph然后向 ResourceManager 申请 Slot最后把 subtask 部署到 TaskManager 上。一个集群里只能有一个 active 的 JobManager。如果是 HA 集群其他 JobManager 处于 standby 状态。这一点在本地 Standalone 里通常只有一个所以不用太担心。TaskManager 也是 JVM 进程负责当前节点上的任务运行和资源管理。它的资源通过 Task Slot 划分。至少需要一个 TaskManager。TaskManager 中 Slot 的最大数量就是它能并发处理任务的上限。多个算子可以在一个 Slot 中执行这就是 Slot 共享机制。默认情况下Flink 允许同一作业的不同 subtask 共享 Slot只要它们来自同一个作业。好处是集群所需 Slot 数等于作业中最高并行度而不是所有算子并行度之和同时资源利用率更高。在开始配置之前你需要准备一台机器安装 JDK 8 或 11Flink 1.10 附近版本对 JDK 8 支持较好新版本建议 JDK 11。下载 Flink 二进制包解压到本地目录比如/opt/flink。确认JAVA_HOME已设置。Standalone 模式下你只需要启动 JobManager 和 TaskManager 两个进程。启动脚本在bin/目录下start-cluster.sh会同时启动一个 JobManager 和一个 TaskManager。默认配置下TaskManager 的 Slot 数是 1这个值在flink-conf.yaml里由taskmanager.numberOfTaskSlots控制。这里先给出一份可复制的flink-conf.yaml关键配置片段。路径通常是conf/flink-conf.yaml。你可以直接修改这几项# JobManager 的 RPC 地址本地 Standalone 用 localhost jobmanager.rpc.address: localhost # JobManager 的 RPC 端口 jobmanager.rpc.port: 6123 # JobManager 管理的内存本地测试给 1024m 足够 jobmanager.memory.process.size: 1024m # TaskManager 管理的内存 taskmanager.memory.process.size: 2048m # 每个 TaskManager 提供的 Slot 数量这里设为 4 taskmanager.numberOfTaskSlots: 4 # 作业默认并行度设为 2 parallelism.default: 2注意taskmanager.numberOfTaskSlots和parallelism.default是两个不同层面的东西。前者是静态的表示 TaskManager 能提供多少并发执行能力后者是动态的表示作业运行时实际使用的并发能力。假设你有 3 个 TaskManager每个 3 个 Slot总共 9 个 Slot。如果parallelism.default1那 9 个 Slot 只用了 1 个剩下 8 个空闲。所以设置合适的并行度才能提高效率。配置改完后用bin/start-cluster.sh启动集群。启动后可以访问http://localhost:8081看到 WebUI。如果 WebUI 打不开先检查 JobManager 日志log/flink-*-standalone-*.log。3. 可复制配置并行度、Slot 与作业提交参数这一节把配置和提交命令完整串起来。你可以在本地直接复制执行。首先确认conf/flink-conf.yaml里已经设置了taskmanager.numberOfTaskSlots: 4和parallelism.default: 2。然后启动集群cd /opt/flink bin/start-cluster.sh启动后用jps应该能看到StandaloneSessionClusterEntrypointJobManager和TaskManagerRunnerTaskManager两个进程。接下来准备一个简单的作业。Flink 自带示例 jar 包在examples/streaming/目录下比如TopSpeedWindowing.jar。我们用它来提交bin/flink run -p 4 examples/streaming/TopSpeedWindowing.jar这里的-p 4表示把作业并行度设为 4。因为我们的 TaskManager 有 4 个 Slot所以刚好能跑满。如果你不指定-p就会用parallelism.default的值也就是 2。提交后命令行会输出 JobID 和运行状态。你可以用bin/flink list查看正在运行的作业bin/flink list如果想看已完成的作业用bin/flink list -a。现在重点来了并行度和 Slot 的关系。假设你把-p设成 5而总 Slot 只有 4作业会怎样它会一直处于RESTARTING或直接失败日志里会出现类似「No enough slots」的提示。这就是因为 JobMaster 申请不到足够的 Slot。另外Slot 共享机制会让事情变得不那么直观。默认开启 Slot 共享时一个 Slot 可以容纳整个作业管道。也就是说即使作业有多个算子只要最高并行度是 44 个 Slot 就能跑完。如果你关闭 Slot 共享在代码里设置slotSharingGroup或配置cluster.evenly-spread-out-slots等那每个算子都会单独占 Slot所需 Slot 数会变成所有算子并行度之和。所以入门阶段建议保持默认的 Slot 共享。再给一个代码层面的并行度设置示例。如果你写 DataStream 程序可以在环境里设置StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(4);也可以在单个算子后面跟.setParallelism(2)来覆盖全局并行度。算子的并行度可以不同但整个作业的并行度取最大值。提交作业时还可以用-d参数让 Client 断开连接detached 模式bin/flink run -d -p 4 examples/streaming/TopSpeedWindowing.jar这样提交完命令就返回作业在后台跑。适合长期运行的流作业。4. 验证请求与成功结果TaskManager 注册与 Slot 使用情况作业提交后怎么确认 TaskManager 注册成功、Slot 被正确使用有三个途径WebUI、命令行、日志。先看 WebUI。打开http://localhost:8081首页会显示 TaskManager 数量、Slot 总数、可用 Slot 数。点击左侧「Task Managers」能看到每个 TaskManager 的 ID、地址、Slot 数、可用 Slot 数。如果这里显示 1 个 TaskManager、4 个 Slot说明注册成功。再看作业详情。点击正在运行的作业进入「Overview」能看到作业的并行度、Task 数量、状态。如果状态是RUNNING说明调度成功。点击「Exceptions」可以看有没有报错。命令行方式bin/flink list输出会列出作业 ID 和状态。如果想看更详细的 Slot 信息可以查看 JobManager 日志tail -f log/flink-*-standalone-*.log日志里会打印 TaskManager 注册、Slot 分配、subtask 部署等信息。比如你会看到类似「Registered TaskManager at ...」「Allocated slot ...」的记录。还有一个实用命令是bin/flink info可以查看作业的执行计划bin/flink info examples/streaming/TopSpeedWindowing.jar它会输出算子和并行度信息帮你确认并行度设置是否符合预期。如果你提交的作业一直处于RESTARTING先检查 Slot 是否够用。在 WebUI 的 Task Managers 页面看「Free Slots」是不是 0。如果是 0说明 Slot 被占满了要么降低并行度要么增加 TaskManager 或 Slot 数。成功的结果应该是作业状态RUNNINGTaskManager 页面显示 Slot 总数和可用数符合预期日志里没有资源不足的报错。5. 本篇常见错排查401、local proxy failed、reading choices、OAuth 等真实报错这一节整理几个入门阶段容易遇到的报错以及对应的排查思路。注意这些报错不一定都出现在 Standalone 本地场景但你在接入外部服务或使用 API 时可能会碰到。报错一401 Unauthorized如果你在调用外部 API 或模型服务时看到 401通常是 Key 无效或没带上。检查请求头里的 Authorization 字段是否正确Key 是否过期。如果你用的是 TaoToken 这类服务确认 Base URL 和 Key 是否匹配。Base URL 一般是https://taotoken.net/apiKey 在控制台生成。报错二local proxy failed这个报错通常出现在网络请求被本地代理拦截时。检查环境变量HTTP_PROXY、HTTPS_PROXY是否设置成了不可用的地址。在本地开发时可以临时取消这些环境变量再试。另外确认你的请求地址是可达的不要被本地防火墙拦截。报错三reading choices这个报错常见于调用模型对话接口时返回体里没有choices字段。原因可能是请求格式不对或者模型 ID 写错了。检查你的请求 JSON 里model字段是否和平台提供的 Model ID 一致。如果你用的是 TaoToken 的模型对话功能可以在控制台确认可用的模型列表。报错四OAuth 相关错误如果你在配置 Claude Code 或类似工具时看到 OAuth 报错通常是认证流程没走完或 token 过期。检查配置文件里的认证信息是否完整。对于 Claude Code 接入需要配置 Base URL、Key、Model ID 三件套。Base URL 用https://taotoken.net/apiKey 用控制台生成的Model ID 按文档填写。报错五No enough slots这是 Flink 本身的报错。原因就是并行度大于总 Slot 数。解决办法降低-p值或者增加taskmanager.numberOfTaskSlots或者多启动几个 TaskManager。在 Standalone 下你可以手动再启动一个 TaskManagerbin/taskmanager.sh start每启动一个就多一份 Slot。报错六TaskManager 注册不上如果 WebUI 里 TaskManager 数量一直是 0检查jobmanager.rpc.address是否配置正确。本地 Standalone 用localhost。如果 JobManager 和 TaskManager 不在同一台机器要改成 JobManager 的实际 IP。另外检查端口 6123 是否被占用。排查时日志是最好的朋友。JobManager 日志在log/flink-*-standalone-*.logTaskManager 日志在log/flink-*-taskexecutor-*.log。先看有没有ERROR或Exception再顺着堆栈找原因。6. 继续深入从 Standalone 到生产部署的路径把本地 Standalone 跑通之后你对 JobManager、TaskManager、并行度、Slot 的关系应该有了直观感受。接下来可以往两个方向深入。第一个方向是资源管理平台。Standalone 下 ResourceManager 只能分配已有 TaskManager 的 Slot不能动态拉起新 TaskManager。生产环境通常用 YARN 或 KubernetesResourceManager 可以向资源平台申请容器来启动 TaskManager实现弹性伸缩。流程上Client 提交作业后ApplicationMaster 会启动 JobManager然后 JobManager 向 ResourceManager 申请资源启动 TaskManagerTaskManager 启动后向 JobManager 发送心跳并等待分配任务。这个链路和本地 Standalone 的逻辑是一致的只是资源来源不同。第二个方向是执行图。你写的 Stream API 代码先生成 StreamGraph经过优化比如算子链生成 JobGraph提交给 JobManager 后再转换成 ExecutionGraph最后部署到 TaskManager 上形成物理执行图。理解这四层图能帮你定位「为什么我的算子没有链在一起」「为什么并行度和我设的不一样」这类问题。算子链的条件是并行度相同且是 one-to-one 操作两个条件缺一不可。如果你在接入模型服务或 API 时遇到认证、Base URL、Model ID 配置问题可以到 TaoToken 的接入文档看具体步骤https://taotoken.net/doc?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewrite 。需要生成 Key 的话控制台在 https://taotoken.net/console?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewrite API Keys 管理页在 https://taotoken.net/api-keys?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewrite 。想先验证模型对话是否通可以用 https://taotoken.net/models?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewrite 。如果是长期做编码或 Agent 开发Coding Plan 在 https://taotoken.net/coding-plan?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewrite 。Claude Code 接入参考 https://taotoken.net/claude-code?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewrite 。最后留一个实用技巧本地调试时把taskmanager.numberOfTaskSlots设成 4parallelism.default设成 2这样既有余量观察 Slot 使用又不会因为并行度太高导致资源紧张。提交作业时用-p显式指定并行度方便对照 WebUI 看 Slot 分配。每次改完配置重启集群再验证避免旧配置干扰。