In a Flink cluster, what do the JobManager and the TaskManager each do?
answer
- two process kinds, one coordinates
- master hands out slots, workers run threads
- records never pass through the master
- dispatcher, a master per job, slot bookkeeper
basics
~20 sFlink's JobManager is the coordinator: it accepts submissions, turns a job into a schedule, allocates task slots and triggers checkpoints. TaskManagers are the worker processes that run the operator subtasks and exchange records directly with each other.
solid answer
~50 sA Flink cluster is one **JobManager** process plus one or more **TaskManager** processes. The JobManager is the master: its *Dispatcher* exposes the REST endpoint and Web UI and starts one *JobMaster* per submitted job; the *JobMaster* turns the JobGraph into an ExecutionGraph, schedules subtasks into slots, and coordinates checkpoints; Flink's own *ResourceManager* tracks task slots and, on YARN or Kubernetes, asks the resource provider for more TaskManager containers when slots are missing. A TaskManager is a JVM worker: it offers `taskmanager.numberOfTaskSlots` slots, runs each subtask as a thread, holds the network buffers used to shuffle records, and holds local state. Records never travel through the JobManager — TaskManagers exchange data directly. If a TaskManager dies, the job restarts from the last checkpoint; if the JobManager dies without HA configured, its jobs go down with it.
code
bash · 3 lines./bin/start-cluster.sh
./bin/flink run -m localhost:8081 -p 4 target/job.jar
./bin/flink list -m localhost:8081go deeper
Be able to name the two processes and say which coordinates and which executes. Knowing that the JobManager exposes the Web UI on port 8081 and that TaskManagers provide the slots is enough at this level.
Explain the JobManager's internal parts — Dispatcher, one JobMaster per job, Flink's own ResourceManager — and how a subtask gets from a submitted jar to a running thread in a slot.
Be ready to reason about failure: what a TaskManager loss costs in reprocessing, what high availability actually recovers, and how the same architecture maps onto standalone, YARN and native Kubernetes deployments.
Own the blast-radius argument. The JobManager is the shared, HA-critical component, so the decision of how many jobs share one is a reliability and tenancy decision, not a packing decision.
## The two process roles Every Flink deployment, whatever the resource provider, is built from exactly two kinds of long-lived process. The **JobManager** is the coordinator (there is one active JobManager per cluster, possibly with standbys). The **TaskManagers** are the workers; you run as many as you need capacity for. Nothing else is required — the Web UI, the REST API and the scheduler all live inside the JobManager process. ## What lives inside the JobManager The JobManager is not one monolithic component; it hosts three: - **Dispatcher** — the REST endpoint and Web UI. It receives job submissions and the job's jars, spawns a JobMaster for each job, and keeps the archive of finished jobs. - **JobMaster** — one instance *per job*. It converts the submitted JobGraph into the parallel ExecutionGraph, requests slots for every subtask, deploys those subtasks to TaskManagers, tracks their state, triggers and confirms checkpoints, and drives the restart strategy when something fails. - **ResourceManager** — Flink's own slot bookkeeper, not to be confused with YARN's ResourceManager. It knows every registered TaskManager and every free slot. In a standalone cluster it can only hand out the slots that already exist; in a native Kubernetes or YARN deployment it actively *requests new TaskManager pods or containers* when a job asks for more slots than exist, and releases idle ones after a timeout. ## What a TaskManager does A TaskManager is a single JVM. It advertises a fixed number of **task slots** (`taskmanager.numberOfTaskSlots`, default 1) and registers itself with the JobManager's ResourceManager. Once subtasks are deployed into its slots, the TaskManager: - runs each subtask as a **thread** inside the JVM (chained operators run in the same thread); - owns the **network buffer pools** used to serialize records and ship them to other TaskManagers — a shuffle is TaskManager-to-TaskManager traffic, never via the JobManager; - owns the **local state**, including the working directory of the embedded RocksDB instance when that state backend is used; - sends periodic **heartbeats** to the JobManager and metrics to the metrics reporter. ## How they find each other On startup a TaskManager registers with the ResourceManager over RPC. The JobMaster asks the ResourceManager for slots, the ResourceManager offers slots from registered TaskManagers, and the JobMaster deploys subtasks straight to the TaskManager. If a job needs more slots than are registered and the deployment cannot create more TaskManagers, the job stays in a scheduling state and eventually fails with a `NoResourceAvailableException` (for example "Could not acquire the minimum required resources.") rather than running at reduced parallelism. That is the default scheduler; the opt-in Adaptive Scheduler is the one that can start a job at whatever parallelism the available slots allow. ## Failure behaviour - **TaskManager loss**: the JobMaster notices the missed heartbeats, cancels the job's remaining subtasks and restarts the job (per its restart strategy) from the last completed checkpoint, once replacement slots are available. Work done since that checkpoint is redone. - **JobManager loss**: without high availability, everything it coordinated goes down. With HA configured (ZooKeeper or the Kubernetes HA services), a standby JobManager is elected leader, reads the job metadata from the HA storage, and recovers running jobs from their last checkpoint. That asymmetry is why the JobManager is the piece people put HA around and the piece that determines the blast radius of a shared cluster. ## How it maps onto real deployments The roles do not change between resource providers — only who starts the processes does: - **Standalone**: you (or a script, or a Kubernetes Deployment) start the JobManager and TaskManager processes yourself. - **Native Kubernetes**: the JobManager runs as a pod and Flink's ResourceManager creates TaskManager pods on demand through the Kubernetes API. - **YARN**: the JobManager runs inside the YARN ApplicationMaster container, and Flink's ResourceManager asks YARN's ResourceManager for TaskManager containers. ## Things people get wrong The most common misconception is that data flows through the JobManager — it does not; it carries control-plane traffic (deployment, checkpoint triggers, heartbeats, metrics) only, which is why a single JobManager can coordinate a job pushing millions of records per second. The second is confusing Flink's ResourceManager with YARN's: they are different components that talk to each other. The third is assuming a TaskManager runs one job — in a session cluster its slots can hold subtasks from several different jobs at once, which is exactly why one bad job can hurt its neighbours.
- If the JobManager restarts under high availability, does the running job keep processing?No — the job is not running while there is no leader. A standby JobManager wins the leader election, reads the job's metadata from the HA storage (ZooKeeper or Kubernetes ConfigMaps plus a durable path), and restarts the job from its last completed checkpoint. Processing resumes with a gap, and everything since that checkpoint is reprocessed.
- Why can one Flink JobManager coordinate a job moving millions of records per second?Because it is a control-plane component only. It deploys subtasks, triggers checkpoints, receives heartbeats and collects metrics; the records themselves flow TaskManager to TaskManager over the network stack. JobManager load scales with the number of subtasks, jobs and checkpoints, not with throughput — which is why very large ExecutionGraphs and very frequent checkpoints are what actually strain it.
- What is the difference between Flink's ResourceManager and YARN's ResourceManager?They are separate components with the same name. Flink's ResourceManager lives inside the JobManager and tracks Flink task slots. YARN's ResourceManager is the cluster-wide scheduler that hands out containers. In a YARN deployment Flink's ResourceManager is a *client* of YARN's: it asks YARN for containers and then starts a TaskManager inside each one.
The JobManager is the shift supervisor with the roster and the clipboard; the TaskManagers are the workers at the benches. The supervisor decides who does what and calls the breaks, but the parts move directly from bench to bench.
saying these in an interview costs you the question
- Saying records are routed through the JobManager
- Calling the TaskManager a single-job process
- Confusing Flink's ResourceManager with YARN's ResourceManager
- Assuming a JobManager crash is transparent without HA configured
- Believing the default scheduler runs a job at reduced parallelism when slots are short