How do Flink session mode and application mode differ when you submit a job?
answer
- one cluster many jobs, or one cluster one job
- where does main() actually execute
- the client stays thin in one of them
- isolation versus submission latency
- the YARN-only middle option was removed
basics
~10 sIn Flink session mode a long-lived cluster hosts many jobs, and a CLI submission runs main() on the client. In application mode each application gets its own cluster and main() runs on the JobManager.
solid answer
~50 s**Session mode** starts a cluster first and submits jobs into it. From the CLI, the client executes the application's `main()`, builds the JobGraph and ships it with the jars to the Dispatcher (a jar run through the REST API runs `main()` on the JobManager instead); the cluster outlives the job. Submission is fast and resources are shared, but there is no isolation: one job's failure, memory leak or noisy resource use hits its neighbours. **Application mode** (on Flink 2.3, `flink run -t kubernetes-application` or `-t yarn-application`) spins up a cluster dedicated to one application and runs `main()` *on the JobManager*, so jar download and graph building happen in the cluster; the cluster shuts down when the application finishes. It is the recommended shape for production long-running jobs. Flink 2.0 removed both the old per-job mode, which built the graph on the client but got a dedicated cluster, and the separate `run-application` CLI action.
code
bash · 8 lines# session: start once, submit many; a CLI submission builds the graph on the client
./bin/flink run -m jobmanager:8081 -p 8 orders.jar
# application: cluster per application, main() runs on the JobManager
./bin/flink run -t kubernetes-application \
-Dkubernetes.cluster-id=orders-etl \
-Dkubernetes.container.image.ref=registry/orders:1.4 \
local:///opt/flink/usrlib/orders.jargo deeper
Know that a Flink job is submitted with the flink run CLI against a cluster, and that the cluster may be long-lived and shared or dedicated to your job.
Explain where main() executes in each mode and what that means for the client machine, and state the isolation difference between a shared session cluster and a per-application cluster.
Be ready to justify a choice for a specific workload and to walk through the operational consequences: failure blast radius, per-job upgrades, resource accounting, and how each mode maps onto Kubernetes objects.
Own the platform default. Decide whether teams get a shared session cluster, application-mode deployments, or both, and what that means for cost attribution, upgrade cadence and how many JobManagers your organisation runs.
## Two orthogonal choices When deploying Flink you pick two things independently: **who provides the resources** (standalone, YARN, native Kubernetes) and **what the cluster's lifecycle is bound to** (a session, or one application). The deployment mode is the second choice, and it is what determines isolation and where your `main()` method executes. ## Session mode A session cluster is started ahead of time and lives independently of any job: a JobManager, some TaskManagers, an idle Web UI. Jobs are then submitted into it with `flink run -t remote` (or `-t yarn-session` / `-t kubernetes-session` after starting the session). What happens on submit: 1. The **client** runs the application's `main()` method — for a CLI or SQL Client submission; a jar submitted through the REST API or the Web UI has its `main()` run on the JobManager instead, with the cluster still shared. `env.execute()` does not run anything locally — it builds the JobGraph and hands it to the deployment code. 2. The client uploads the user jars and the JobGraph to the Dispatcher over REST. 3. The Dispatcher spawns a JobMaster, which requests slots from the already-registered TaskManagers. Strengths: submission latency is low (no cluster to start), resources are pooled across many short jobs, and one Web UI shows everything. Weaknesses are all about sharing: every job depends on the same JobManager, so its failure or overload affects all of them; TaskManagers hold subtasks from different jobs, so a TaskManager crash restarts several jobs; user-code class loading conflicts and jar bloat accumulate in one place; and resource accounting per job is muddy. For CLI submissions the client machine also has to be able to run `main()`, which for a large SQL or Table API application can mean real CPU, memory and jar downloads on a laptop or a CI runner. ## Application mode Application mode inverts the flow. `flink run -t kubernetes-application -Dkubernetes.cluster-id=orders-etl ...` asks the resource provider to start a JobManager *whose entrypoint is your application*. The `main()` method executes on the JobManager, inside the cluster. The client's only job is to hand over the request and exit. On Flink 2.3 there is no separate `run-application` action any more — 2.0 removed it — and the ordinary `run` action switches to application mode whenever the target ends in `-application`. That buys three things: - **Isolation.** One application, one cluster, one JobManager, its own TaskManagers. A crash, a class-loading conflict or a memory problem is contained, and the cluster's resource usage is exactly that application's usage. - **A thin client.** Jar downloading and graph building happen in the cluster on the JobManager, so the submitting machine does not need the dependencies or the memory. This matters most for big Table API/SQL programs whose planning is expensive. - **A bound lifecycle.** The cluster terminates when the application does, which makes it a natural fit for a Kubernetes Deployment per job or an operator-managed `FlinkDeployment`. The cost is startup latency — every submission provisions a cluster — and, if you run hundreds of tiny jobs, one JobManager's overhead per job. ## Per-job mode, and why it is gone Per-job mode (YARN only) also gave each job a dedicated cluster, but the client still executed `main()` and built the JobGraph before the cluster existed. It combined the client-side cost of session mode with the startup cost of application mode, and application mode supersedes it. It was deprecated before 2.0 and removed in Flink 2.0, whose release notes point users to application mode; `-t yarn-per-job` is no longer an accepted target, so it is something you meet only on old platforms. For SQL jobs, 2.0 also taught the SQL Gateway to run in application mode as the replacement. ## Choosing A long-running streaming job in production wants application mode: it will run for months, it needs its own failure domain, and a per-application cluster maps cleanly onto Kubernetes objects and onto per-team cost accounting. A pool of short ad-hoc SQL queries, an interactive notebook, or a CI environment wants session mode: the jobs are seconds-to-minutes long and paying cluster startup per query would dominate. Many platforms run both — a shared session cluster for exploration, application-mode deployments for anything on a schedule or a pager. ## Operational consequences In session mode, upgrading Flink means draining and restarting a cluster many teams depend on. In application mode, each application upgrades on its own schedule, and rolling out a new Flink version becomes a per-job rollout with a savepoint. In session mode the Web UI is the shared place to find a job; in application mode each job has its own UI endpoint, which is why platforms in this shape almost always add a central catalogue or use the Flink Kubernetes Operator to keep track.
- Why does application mode need the jar path to be reachable from inside the cluster?Because `main()` runs on the JobManager, not on the client. The entrypoint must load the application jar itself, so it is normally baked into the container image and referenced with a `local://` path, or fetched from a distributed filesystem the JobManager can read. There is no session-style jar upload to the Dispatcher; native Kubernetes can opt in to `kubernetes.artifacts.local-upload-enabled`, which copies a client-local jar to DFS before deployment.
- Which mode would you pick for dozens of short ad-hoc SQL queries, and why?Session mode. Each query lives for seconds to minutes, so paying cluster provisioning per query would dominate the runtime, and pooling TaskManagers across queries keeps utilization sane. The isolation you give up matters little when nothing is long-running or on a pager; you accept that one bad query can disturb its neighbours.
- What does application mode change about upgrading the Flink version?It turns a fleet-wide event into a per-job rollout. A session cluster's upgrade forces every tenant to drain at once; with one cluster per application you take a savepoint, redeploy that application on the new image, and restore — job by job, on each team's own schedule, with an easy rollback to the previous image and the same savepoint.
saying these in an interview costs you the question
- Saying main() runs on the client in application mode
- Treating session mode as isolated because jobs have separate JobMasters
- Recommending per-job mode on a current Flink version
- Confusing the deployment mode with the resource provider
- Claiming application mode requires Kubernetes