一套集群 — JobManager 做决定,TaskManager 来执行
一句话总结
Flink 集群由一个负责决策的进程(JobManager)和多个负责运行的进程(TaskManager)组成。一行 SQL 会变成作业图,JobManager 再把它分配到 TaskManager 的槽位上。什么东西在哪里运行、为什么挂掉,全都会留在 REST API 里。
为什么需要它
第一次负责流处理时,通常从“写了 SQL 就有结果”起步。可是在生产环境里收到的问题并不是 SQL。“作业为什么停了”“并行度调高了为什么没变化”“TaskManager 内存给了 4GB,为什么堆只有 1GB”“为什么没有重启”。这些问题全都关乎引擎的结构。
批处理脚本挂了重新跑一遍就行。流处理作业会挂上好几周,期间还持有状态。所以必须知道“哪个进程负责什么”,出故障时才能决定该看哪里。本模块先把这张地图铺开。
工作原理
JobManager 由三个部件组成(官方文档中的 Flink Architecture)。
| 部件 | 职责 |
|---|---|
| Dispatcher | 开放 REST 接口(默认 8081)和 Web UI,作业提交进来时创建一个 JobMaster |
| ResourceManager | 管理槽位。在独立(standalone)部署中,它只能分配现有 TaskManager 的槽位,不能启动新的 TaskManager |
| JobMaster | 负责单个作业的执行。每个作业对应一个 |
TaskManager 负责实际执行运算并收发数据。资源分配的最小单位是任务槽位,槽位数就是可以同时运行的任务数。有两点需要注意。第一,槽位并不划分 CPU——按文档的说法,目前的槽位只划分托管内存(managed memory)。第二,多个算子可以被连成算子链,在同一个线程中运行。这就是 Source → 过滤 → 转换会成为同一个任务的原因。线程之间的传递和缓冲消失了,吞吐量随之提高。
并行度是指一个算子要复制成几份来运行。并行度为 2 的聚合会拆成两个子任务,每个各占一个槽位。批处理与流处理在这里出现分歧。批处理作业的默认调度器是自适应批调度器,它会按输入数据的大小重新确定未指定并行度的算子的并行度。标准是每个任务平均 16MB(execution.batch.adaptive.auto-parallelism.avg-data-volume-per-task)。对很小的文件设置 parallelism.default = 2,并行度仍然是 1,这不是 bug,而是这种设计。关闭自动决定的开关是 execution.batch.adaptive.auto-parallelism.enabled。
内存是本模块最常被误解的部分。taskmanager.memory.process.size 是 JVM 进程整体的大小。其中要先扣除 JVM 元空间和 JVM 开销,剩下的 Flink 内存再分为框架堆、任务堆、托管内存和网络内存。
프로세스 전체 = Flink 메모리 + JVM 메타스페이스 + JVM 오버헤드
Flink 메모리 = 프레임워크 힙 + 태스크 힙 + 프레임워크 오프힙 + 태스크 오프힙 + 관리 메모리 + 네트워크 메모리
这个 Pod 的上限是 2Gi,所以 TaskManager 设为 1024MiB。这样一来,任务堆只剩 200MiB 出头。“给了内存却发生 OOM”的报告,大多是把这张表看反了。实际数值会以字节为单位出现在 /taskmanagers 响应的 memoryConfiguration 中。
失败与重启。 作业挂掉后,JobMaster 会查看重启策略。根据文档(Task Failure Recovery),不开启检查点时默认不重启,开启后默认是 exponential-delay。所以没有检查点的作业,异常记录的最上面打出来的不是原因,而是 Recovery is suppressed by NoRestartBackoffTimeStrategy。真正的原因在下面 Caused by: 链的末尾。
curl -s localhost:8081/overview # TaskManager 수 · 슬롯 수 · 버전
curl -s localhost:8081/jobs/overview # 잡 목록과 상태(FINISHED · FAILED · RUNNING)
curl -s localhost:8081/jobs/<jid> # 정점(vertex)별 병렬도
curl -s localhost:8081/jobs/<jid>/exceptions # 실패 기록(exceptionHistory)
在现场相遇的样子
最常见的场景是“并行度调高了却没有变快”。原因通常是三者之一:槽位不够,作业在等待资源;是批处理作业,自适应调度器按数据大小压低了并行度;Source 的分片(文件数、分区数)比并行度少,部分子任务闲着。这三种情况都可以通过查看 /jobs/<jid> 中每个顶点的并行度和子任务数来区分。如果只看仪表板上的吞吐量曲线,三者看起来一模一样。
第二种是改了配置却没有生效。numberOfTaskSlots 或内存大小这类进程启动时读取的值,只改文件是不会有任何反应的,必须重新启动 TaskManager。反过来,SET 'parallelism.default' 这类在提交作业时读取的值,下一个作业就立刻生效。不清楚属于哪一种,就只会积累“重启之后就好了”的经验之谈。
第三种是故障报告。告警只会以“作业 FAILED”一行送来。如果把异常记录最上面的那句话原样粘贴进去,所有人都会误以为重启策略是原因。找出链末尾的 Caused by 并写下来,这个习惯决定了报告的质量。
下一项实验要做什么
启动集群,用 REST 查看 TaskManager 和槽位,并把内存预算换算成数字。运行一个以批处理方式聚合订单文件的作业,在作业列表中找到它,再用并行度 2 重新运行同一个作业,确认顶点的并行度。把槽位改为 4 个并重新启动后,故意制造一个把字符串转换成整数时挂掉的作业,把真正的原因和重启策略整理成报告。