Skip to content
WSLDeep Dive Published Updated 7 min readViews unavailable

Apache Flink in WSL: Local Session Clusters and Job Lifecycle

Run a local Flink session cluster in WSL, submit a bounded sample job, inspect TaskManager capacity, and separate local state from production recovery.

Apache Flink’s local installation gives developers a way to bring up a small session cluster on one Linux machine and exercise job submission, the JobManager UI, TaskManager registration, and SQL or DataStream examples. Running it inside WSL is convenient for application development, but it does not create independent process failure domains or a durable streaming service. All local processes share one WSL VM, Windows host, and backing storage.

Use the official documentation that matches the Flink distribution you install. The stable first-steps guide describes a local install and gives the release’s Java requirements; Java support can change, so check that page rather than assuming a JDK from another project will work. Keep the archive, configuration, logs, and job artifacts under the distro’s Linux filesystem to avoid introducing Windows-mounted filesystem behavior into a Linux workload.

Start in foreground and prove process ownership

Download a Flink binary distribution from the Apache project and verify its release and checksum through the official release channels. Extract it into a versioned directory. Before starting anything, inspect Java and available WSL resources:

java -version
free -h
nproc
df -h "$HOME"

The local setup uses the distribution scripts. Start the cluster from the extracted Flink directory and inspect the process and log output:

./bin/start-cluster.sh
ss -ltnp | grep ':8081'

The stable guide describes the web interface at localhost port 8081 and shows a single TaskManager for the local cluster. Treat actual output and the installed version’s configuration as authoritative. Open the UI and confirm the JobManager sees the TaskManager and its available task slots. A listening port alone does not prove the TaskManager registered or can execute a job.

Submit one small example job from the same distribution, then inspect its state in the UI and read the logs. The official local guide shows a bundled TopSpeedWindowing.jar example. Keep the first run bounded and use a job with predictable completion. Afterward, stop the local cluster with ./bin/stop-cluster.sh and verify its processes exit. Avoid launching repeated clusters on the same ports or with the same state directory.

Know what session mode means

A session cluster starts the cluster first and accepts jobs afterward. Multiple jobs can be submitted to the same session, which makes it convenient for interactive development. This differs from application mode, where the cluster lifecycle is associated with an application. Do not confuse local script-based cluster startup with a deployment controller that restarts failed processes; the standalone documentation places process restart and resource allocation responsibility on the operator.

The JobManager coordinates job execution and exposes the web interface. TaskManagers provide task slots and execute subtasks. A job can fail locally because the TaskManager has too few slots, Java cannot reserve memory, a connector or filesystem plugin is absent, or a job artifact cannot be loaded. Read the JobManager and TaskManager logs together; the UI is a useful state view but not a full diagnosis.

For SQL testing, start the SQL Client included in the distribution after the cluster is running. Use a bounded local source or a generated table and confirm the resulting job in the UI. Do not use an unbounded source in an experiment that expects the job to finish. A streaming job that remains RUNNING may be correct, but define a time-bounded cancellation step so it does not consume laptop resources indefinitely.

Configure memory and local data deliberately

Flink’s JobManager and TaskManager use JVM processes with separate memory and resource configurations. Start from the release defaults and collect process RSS, logs, job duration, and WSL memory pressure before tuning. A local VHDX’s remaining disk capacity is also relevant for temporary spill and checkpoint data. Do not set JVM memory to the full WSL cap; leave room for Linux, Windows, and other workloads.

Keep local checkpoints, state, and temporary paths on the distro filesystem. A local filesystem checkpoint can help exercise checkpoint configuration, but it is not an independent recovery copy. Losing the WSL virtual disk or host can remove both the running cluster and local checkpoint directory. For meaningful durability tests, configure the same external storage and failure domains that the target deployment uses, then test restore and rescaling there.

Separate a checkpoint from a savepoint in the test plan. Checkpoints support Flink’s configured state recovery behavior during execution; savepoints are an operator-triggered mechanism used for controlled stop, upgrade, or migration workflows. A local run can demonstrate that the job reaches a checkpoint and can be restarted under the same local conditions, but it cannot prove that a checkpoint survives loss of the host that stores it. Keep the exact job artifact, configuration, and state metadata together for any restore exercise, and test compatibility with the target release before discarding the original environment. Avoid deleting a state directory while a job is still using it.

Flink’s configuration uses specific options for ports, REST address, RPC address, memory, and temporary directories. Consult the configuration reference for the selected release before changing a file. Avoid changing memory, parallelism, and slot counts simultaneously; alter one variable, rerun the same job, and compare a measured outcome. The default local setup is a starting point, not a sizing recommendation for production.

WSL networking and host access

The Flink UI is an HTTP service. Keep it bound to the local development boundary, especially if the server has no authentication configured for this isolated lab. Do not publish its port to a LAN to bypass a Windows localhost problem. Verify reachability in stages: check the Linux listener, open it from a WSL browser or curl, then use the Windows client path expected by your active WSL networking mode.

Microsoft documents NAT and mirrored networking separately. The Windows host’s ability to reach a service through localhost does not mean a remote LAN client can reach it, and a port forward does not change Flink’s internal JobManager/TaskManager communication. When components run in separate containers or machines, use the network addresses appropriate to that topology; a localhost URI inside one process is not a universal address for peers.

Failure handling and cleanup

If startup fails, check Java, port conflicts, writable directories, and the latest JobManager log. If the UI shows no TaskManager, inspect TaskManager startup and its connection address. If a job is accepted but fails, review the exception at the failing operator, classpath and plugin availability, input path, and resource state. If the service runs but the Windows browser cannot connect, diagnose WSL forwarding instead of rebinding every listener to 0.0.0.0.

Use the UI’s checkpoint view to distinguish a job that is running from one that is making recoverable progress. For a short local test, it is enough to verify that the configured checkpoint completes and that the job can restore from the intended state path after a controlled restart. Keep checkpoint history and job logs together with the test notes. If the job has no state or checkpointing is disabled, say so explicitly; a stateless example cannot validate recovery behavior.

Before an upgrade, stop the cluster and preserve job source, configuration, and any state that matters. Do not copy active state files as a backup. For a disposable local lab, prefer reproducible configuration and fixtures over preserving an opaque state directory. When validating savepoints or checkpoints, use the exact workflow supported by the Flink release and retain the metadata required to restore them.

Acceptance criteria

Accept the local session cluster when the Flink and Java versions are recorded, the cluster starts with a JobManager and TaskManager, a bounded example job completes or reaches its intended streaming state, logs explain the observed outcome, and shutdown leaves no daemons running. Document that local state and service availability share the WSL VM’s failure domain.

This setup is useful for job-development feedback and basic API compatibility. It is not a production deployment, independent checkpoint store, multi-node scheduling test, or guarantee of fault recovery.

Related:

Sources:

Comments