Hi Dale,
Thanks for the writeup. Yes, regarding your second question, I do see a benefit here. One thing for now. Should HA be required rather than optional for a stateful FlinkMiniCluster? In a normal deployment you can survive without HA because the common failure is TaskManager loss, which the surviving JobManager handles. In a single pod, there are no TaskManager-only failures, every fatal failure is a JobManager failure. So making HA mandatory should make more sense, I believe. On Tue, Jul 28, 2026 at 1:38 AM Dale Lane <[email protected]> wrote: > **TLDR - I'm starting a discussion about adding support for the Flink > Kubernetes Operator to run Flink's MiniCluster in a single pod for > low-throughput jobs that require isolation. I've created a PoC to > demonstrate feasibility, and would like to gauge initial reaction and > gather feedback ahead of writing up a FLIP** > > > ## Motivation > > The Flink Kubernetes Operator supports two approaches for running a job: a > `FlinkDeployment` (a separate JobManager Deployment plus a TaskManager > Deployment/pod group, provisioned either natively or in standalone mode) or > a `FlinkSessionJob` submitted into an existing Flink deployment. Both > approaches run Flink as a distributed cluster with > independently-schedulable JobManager and TaskManager pods. > > This provides scalability and high availability, at the cost of a fixed > baseline cost of at least one JobManager pod and one or more TaskManager > pods. > > A smaller, lighter-weight alternative would be useful for small or > intermittent jobs, where the resource cost of two separately-scheduled pods > is disproportionate to the job itself. A single-pod, self-contained Flink > job that starts fast and needs no multi-pod coordination could be a good > fit for low-throughput jobs that aren't suitable for session clusters > because they need isolation. > > Apache Flink has a mechanism suited to this: MiniCluster. MiniCluster is > an in-JVM cluster that runs a real JobManager and one or more TaskManagers > merged into a single Java process. It instantiates the same component > classes as a full distributed deployment, meaning that job graphs, > checkpoints, save points, state backends, connectors, metrics, and more are > all fully compatible with a normal distributed Flink cluster. > > Users who could use this to run very small Flink deployments in Kubernetes > currently need to deploy this themselves as a bare MiniCluster in a > hand-rolled pod, losing the benefits of the operator's lifecycle > management, status reporting, and savepoint / checkpoint tooling. > > I'm proposing to add first-class support for running a job in MiniCluster > through the operator - allowing users to get the same declarative lifecycle > management, status reporting, and snapshot tooling they already have for > Flink jobs in Kubernetes today, but now for a single-pod topology. > > I do have a working proof-of-concept, validated end-to-end in Kubernetes, > that has been helpful in confirming the mechanism is viable inside the > operator's existing architecture. > > > ## Proposed Change > > ### Parity with existing Flink deployments > > My goal would be to make running a Flink job in a MiniCluster feel as > similar as possible to running a full distributed Flink job. I've been > validating this in my proof-of-concept; confirming that a metrics scrape, > web dashboard, state stored in a mounted PVC, and more all work end-to-end > against a real cluster. > > #### Savepoints and checkpoints > > It will be possible to trigger a savepoint or checkpoint against a > MiniCluster job using the exact same `FlinkStateSnapshot` mechanism already > used for `FlinkDeployment` / `FlinkSessionJob` jobs, with no separate > snapshot workflow for MiniClusters. > > (This helps to enable the migration workflow described below.) > > #### Metrics > > A `metrics.reporter.*` key in `flinkConfiguration`, plus whatever pod > annotations a scraper needs (e.g. Prometheus annotations on `podTemplate`), > will produce a scrapeable metrics endpoint in MiniCluster pods in the same > way it does today, with no MiniCluster-specific metrics configuration > needed. > > #### Persistent volumes > > Volume and volumeMount definitions will be supportable for the single > MiniCluster pod exactly as they are to `FlinkDeployment`'s > `jobManager`/`taskManager` pods, with `state.checkpoints.dir`, > `state.savepoints.dir`, and RocksDB's local-directory configuration > continuing to work as ordinary `flinkConfiguration` keys. > > #### Flink MiniCluster launcher > > A MiniCluster-specific launcher will be needed to construct the > MiniCluster, start it, and run the user's job. This could be provided as > example code (in a similar way to how `examples/flink-sql-runner-example` > is today). Or this could be a bootstrap class added to core Flink, so that > it is available out of the box in a regular Flink image. > > Whichever approach is taken, the goal is that a user's Flink job should > not need to have anything MiniCluster-aware or MiniCluster-specific in it. > > If the user has written a regular Flink application, > StreamExecutionEnvironment.getExecutionEnvironment would give them the > in-process MiniCluster in the same way it does with a real Flink cluster > today. > > If they've written Flink SQL, they can use the existing, unmodified > `sql-runner.jar` from `flink-sql-runner-example` in the same way that they > would for a FlinkDeployment today. > > (This working assumption would also help enable the migration workflow > described below.) > > ### Migration: promoting a MiniCluster job to a full FlinkDeployment > > Because job graphs, savepoints, and checkpoints are format-compatible > between MiniCluster and a full distributed cluster, I want to support a > "start small on MiniCluster, promote to a full distributed > `FlinkDeployment` later" workflow: > > 1. Create a MiniCluster job > 2. Create a FlinkStateSnapshot that has a jobReference pointing at the > MiniCluster job to trigger a savepoint > 3. Use the savepoint path to create a full "normal" FlinkDeployment with > an initialSavepointPath for that savepoint, with > `job.jarURI`/`entryClass`/`args` pointing at the **same unmodified job > jar** used in the MiniCluster job > > > > ## New or Changed Public Interfaces > > There are different ways that we could surface this new Flink topology in > the Operator API. The two obvious options are to create a new kind specific > to MiniClusters, or to modify the existing FlinkDeployment kind to support > a third MiniCluster mode. > > My preference is to add a new kind: `FlinkMiniCluster` - and this is the > approach I've taken with my proof-of-concept. > > The new CRD, `FlinkMiniCluster`, would live alongside the existing > `FlinkDeployment`, `FlinkSessionJob`, `FlinkBlueGreenDeployment`, and > `FlinkStateSnapshot` kinds. > > It represents a single job running inside a single pod, with JobManager > and TaskManager merged into one process — no native/standalone distinction, > no separately-scheduled TaskManager pods. > > Example: > > ```yaml > apiVersion: flink.apache.org/v1beta1 > kind: FlinkMiniCluster > metadata: > name: minicluster-job > spec: > image: my-registry/my-minicluster-driver:latest > flinkVersion: v2_2 > serviceAccount: flink > flinkConfiguration: > high-availability.type: kubernetes > high-availability.storageDir: file:///flink-data/high-availability > state.savepoints.dir: file:///flink-data/savepoints > state.checkpoints.dir: file:///flink-data/checkpoints > execution.checkpointing.interval: "30s" > podTemplate: > spec: > containers: > - name: flink-main-container > volumeMounts: > - mountPath: /flink-data > name: flink-state-volume > volumes: > - name: flink-state-volume > persistentVolumeClaim: > claimName: job-state > resources: > requests: > memory: "512Mi" > cpu: "0.5" > limits: > memory: "1Gi" > cpu: "1" > job: > jarURI: local:///opt/flink/sql-runner.jar > args: [ "/opt/flink/example-job.sql" ] > parallelism: 1 > upgradeMode: savepoint > status: > clusterInfo: {} > jobManagerDeploymentStatus: READY > jobStatus: > checkpointInfo: > lastPeriodicCheckpointTimestamp: 0 > jobId: 123bae75a4bb3fbaebd07052f2f193bd > jobName: insert-into_default_catalog.default_database.output > savepointInfo: > lastPeriodicSavepointTimestamp: 0 > savepointHistory: [] > startTime: "1784833854278" > state: RUNNING > updateTime: "1784833882038" > lifecycleState: STABLE > observedGeneration: 1 > reconciliationStatus: > lastReconciledSpec: '{"spec":... > lastStableSpec: '{"spec":... > reconciliationTimestamp: 1784833851180 > state: DEPLOYED > ``` > > The approach I've taken for the shape of this resource: > > - **Reuse status entirely.** > `jobStatus` (jobId, state,`upgradeSavepointPath`), `lifecycleState`, > `reconciliationStatus`, and `error` are the same fields that > `FlinkDeploymentStatus` and `FlinkSessionJobStatus` already expose. Any > tooling that polls or displays status for Flink resources, and any > `kubectl` workflows that operate on FlinkDeployment, won't need to do > anything to special-case for this new kind. > - **Reuse spec wherever the underlying concept is the same.** > `flinkConfiguration`, `image`, `imagePullPolicy`, `flinkVersion`, > `logConfiguration`, `ingress`, `serviceAccount`, the `job:` block > (`jarURI`, `parallelism`, `entryClass`, `args`, `state`, `upgradeMode`, > `initialSavepointPath`, `savepointRedeployNonce`, `autoscalerResetNonce`) > are all the same fields already used by `FlinkDeployment`. A job definition > is portable between `FlinkMiniCluster` and `FlinkDeployment` by > copy-pasting the spec and config. > - **Exclude what doesn't apply.** > There is one `podTemplate` + one `resources` block instead of separate > `jobManager`/`taskManager` blocks, since there's one pod, not two pod > groups. There is no `mode` field (native / standalone). > > One other key difference is that `spec.job` would be required, so there is > no job-less "session" variant at this stage. > > The idea here is that the pod boots the cluster and runs its one job. > > (I'm treating a job-less session MiniCluster that waits for external job > submission as explicitly out of scope here. It could be a candidate for a > future, separate proposal.) > > This nets out to: > > - **New CRD**: `FlinkMiniCluster` (`flink.apache.org/v1beta1` > <http://flink.apache.org/v1beta1>), a new kind alongside the four > existing operator kinds. New spec/status types, new RBAC rules > (`flinkminiclusters`, `flinkminiclusters/finalizers`, > `flinkminiclusters/status`) in the Helm chart, and a new admission-webhook > dispatch branch (validating/mutating). > - `FlinkStateSnapshot`: `jobReference.kind` gains a new accepted value, > `FlinkMiniCluster`, additive to the existing `FlinkDeployment`/ > `FlinkSessionJob` values. > > **No changes** to any existing kind's spec or status shape > (`FlinkDeployment`, `FlinkSessionJob`, `FlinkBlueGreenDeployment`, or > existing `FlinkStateSnapshot` fields) > > > ## Next steps > > ### Performance measurements > > I haven't yet done specific benchmarking or tuning to see how small an > Operator-managed FlinkMiniCluster could be, however Robert Metzger's talk > on MiniCluster > https://speakerdeck.com/rmetzger/tiny-flink-minimizing-the-memory-footprint-of-apache-flink > described examples of a ~250mb memory footprint for Flink while still being > able to support a throughput of 100mb/s. > > I'm confident that there is a lot of potential here, but more concrete > numbers need to follow if we go with this. > > ### Feedback > > I'm looking for feedback on the idea of promoting MiniCluster to something > that can be managed by the Kubernetes Operator before tidying this up as a > full FLIP, and sharing the prototype for a deeper discussion about > implementation and possible approaches. > > What do you think? > > Do you see a benefit for such an Operator feature? > > Do you see any challenges or obstacles that would be worth exploring with > the prototype? > > > > Thanks in advance for your time! > > Dale > -- > dalelane.co.uk > >
