Fair enough, and I'm happy to leave the specifics for the FLIP.

If someone points storageDir at /tmp or an emptyDir, they presumably know
what they're signing up for. My concern is narrower than that: because a
MiniCluster loses everything at once rather than just a TaskManager, it
seems worth leaning harder on recommending durable storage here than we
would for a FlinkDeployment. This could be tackled by a proper disclaimer
in the docs maybe.

One other thought. I suspect people will pick this up for testing and early
adoption before they use it in production, at least to begin with. If
that's right, the migration path out to a full FlinkDeployment is probably
the part worth getting right first.

On Mon, Aug 3, 2026 at 2:29 AM Dale Lane <[email protected]>
wrote:

> I completely agree that enabling HA for this would be sensible (my example
> CR gives away that most of my PoC testing was with it on) but I'm not sure
> about making mandatory though.
>
> I can see the argument for it (don't let users do silly things that will
> get them into trouble).
>
> Are you thinking that we'd make the storageDir config mandatory as well to
> force the user to enable the HA persistent storage? I guess that wouldn't
> stop them using ephemeral storage like /tmp for that or even an emptydir
> volume. But at least that'd be a conscious decision rather than unknowingly
> doing something risky
>
> Kind regards
>
> D
> --
> dalelane.co.uk
>
> Sent with Proton Mail secure email.
>
> On Sunday, 2 August 2026 at 10:41, P Sinha <[email protected]>
> wrote:
>
> > 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>
> > > <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
> > >
> > >
> >
>


-- 
Best Regards,
Purushottam Sinha
शुभकामनाएं
पुरूषोत्तम सिन्हा

Reply via email to