Flink on Karmada: Building Resilient Data Pipelines on Multi-Cluster K8s - Michas Szacillo & Wang Li
Michas Szacillo, Wang Li
KubeCon + CloudNativeCon Europe 2025 · Session
Overview
This talk, presented by Michas Szacillo and Wang Li from Bloomberg, delves into the critical challenges of operating large-scale, stateful streaming applications like Apache Flink in a multi-cluster Kubernetes environment. It highlights the journey undertaken by Bloomberg's streaming platform team, in close collaboration with the Karmada community, to achieve automated, stateful failover for Flink jobs across multiple Kubernetes clusters. The core of their solution involves leveraging Karmada, an open-source Kubernetes management system, to orchestrate applications and implementing custom enhancements for Flink-specific state preservation during failover events.

Key moments
- 0:00 Introduction to Flink on Karmada talk and agenda
- 2:00 Bloomberg's large-scale Flink platform and overview
- 3:10 Flink's fault tolerance using state snapshots and checkpoints
- 4:00 Flink job deployment and management on Kubernetes
- 6:00 Challenges managing Flink applications at multicluster scale
- 8:00 Envisioning an ideal unified multicluster control plane
Flink on Karmada: Building Resilient Data Pipelines on Multi-Cluster K8s
Speakers: Michas Szacillo, Senior Software Engineer and Tech Lead at Bloomberg; Wang Li, Software Engineer at Bloomberg
Conference: KubeCon EU
YouTube: https://www.youtube.com/watch?v=mqXZ2T-jWuU
Overview
This talk, presented by Michas Szacillo and Wang Li from Bloomberg, delves into the critical challenges of operating large-scale, stateful streaming applications like Apache Flink in a multi-cluster Kubernetes environment. It highlights the journey undertaken by Bloomberg's streaming platform team, in close collaboration with the Karmada community, to achieve automated, stateful failover for Flink jobs across multiple Kubernetes clusters. The core of their solution involves leveraging Karmada, an open-source Kubernetes management system, to orchestrate applications and implementing custom enhancements for Flink-specific state preservation during failover events.
The speakers systematically outline the inherent limitations of Flink's native high-availability mechanisms when faced with complete cluster failures, emphasizing the operational overhead and potential data loss associated with manual recovery processes. They then introduce Karmada as a foundational technology to overcome these hurdles, detailing its features for multi-cluster management, intelligent scheduling, and automated failover. A significant portion of the discussion focuses on the technical intricacies of integrating Flink's internal state management with Karmada's failover capabilities, culminating in a novel approach that uses Kverno webhooks and a custom state preservation enhancement within Karmada to ensure seamless, stateful recovery.
For organizations running critical data pipelines, this talk provides invaluable insights into building robust, resilient streaming platforms that can withstand infrastructure failures beyond the scope of a single Kubernetes cluster. It demonstrates a practical, production-grade solution for achieving true business continuity for stateful applications, significantly reducing manual intervention, minimizing downtime, and enabling more efficient platform operations and maintenance. The collaboration between an enterprise user and an open-source community to drive specific feature development further underscores the talk's relevance to both platform engineers and open-source contributors.
Background
▶ Watch: Introduction to Flink on Karmada talk and agenda (0:00)
Bloomberg operates a massive streaming platform, hosting approximately 1,000 unique Apache Flink jobs across multiple Kubernetes clusters and tiers. These jobs are instrumental in processing critical financial data for various use cases, including ETL, real-time analytics, and event processing, directly supporting Bloomberg's core financial products and providing real-time market insights. Given the scale and criticality, ensuring the reliability and efficiency of this system is paramount.
Apache Flink is a renowned open-source data streaming framework celebrated for its low latency, high scalability, and ability to process vast amounts of data with exactly-once guarantees. A cornerstone of Flink's reliability is its built-in fault tolerance and state management system. Flink jobs are typically long-running, and to ensure continuity in the face of failures, Flink periodically takes state snapshots, known as checkpoints, to persistent storage. Additionally, savepoints can be manually triggered. In the event of a job failure, Flink can automatically restore from the latest checkpoint, minimizing data loss and downtime.
For deployment on Kubernetes, Flink offers native support, enabling scalable and manageable job orchestration. The Flink operator automates the entire job lifecycle, from deployment and upgrades to auto-recovery. Users submit a YAML definition, which the Kubernetes API server processes, and the Flink operator then manages the Flink job within the cluster. Flink's high availability (HA) settings allow it to recover from internal cluster issues such as hardware failures, pod crashes, or transient network glitches. This HA mechanism relies on metadata, often stored in Kubernetes ConfigMaps, pointing to the latest state. However, this recovery model assumes that recovery will occur within the same cluster.
As Bloomberg's streaming platform expanded, managing Flink applications at scale introduced significant challenges. Control planes were inherently tied to single clusters, leading to operational complexity. Users had to manage multiple kubeconfigs, track resource availability across different clusters, and manually coordinate deployments. More critically, while Flink's HA handles intermittent in-cluster failures, it does not address scenarios involving a partial or total cluster failure. If a Kubernetes cluster becomes unavailable, or if critical metadata like the Flink ConfigMap is deleted, the Flink operator cannot reconcile the job's state, requiring manual intervention. This translated into a cumbersome, error-prone manual process for cross-cluster failover, often necessitating users to redeploy jobs to a healthy cluster, understand their last known state, or even run duplicate pipelines in tandem. Furthermore, routine maintenance windows for Kubernetes clusters, which occur quarterly, became expensive and required extensive coordination to migrate users off the affected clusters. These limitations highlighted the urgent need for a more robust, automated, and unified control plane capable of managing Flink jobs across a federation of Kubernetes clusters.
Key Findings
▶ Watch: Flink's fault tolerance using state snapshots and checkpoints (3:10)
The primary challenge identified by Bloomberg was the lack of an automated, stateful failover mechanism for Apache Flink applications across multiple Kubernetes clusters. While Flink provided robust in-cluster fault tolerance, it fell short when an entire cluster experienced significant issues or went offline. This led to a manual, time-consuming, and error-prone recovery process for critical financial data pipelines.
The key findings and solutions explored by the Bloomberg team, in collaboration with the Karmada community, can be summarized as follows:
- Need for a Unified Multi-Cluster Control Plane: The existing operational model, where control planes were tied to individual clusters, created significant overhead. A unified control plane was essential to streamline user experience, simplify deployments, enable intelligent resource scheduling, and centralize health monitoring across a fleet of Kubernetes clusters.
- Introduction of Karmada: The Karmada project was identified as a promising open-source solution. Karmada is a Kubernetes management system designed to manage cloud-native applications across multiple clusters. Its features, including managing groups of clusters from one place, defining application scheduling policies (resource-aware scheduling, cluster affinity), unified authentication, and crucially, support for automated cross-cluster failover, directly addressed Bloomberg's requirements.
- Initial Karmada Integration Limitations: While Karmada offered automated failover, its default configuration lacked awareness of Flink's internal state. When a Flink job was rescheduled to a new cluster by Karmada, it would start from scratch, losing its processing state. This was unacceptable for long-running streaming jobs that require seamless resumption from their last known state.
- Requirement for Flink-Specific State Awareness: To enable graceful recovery, Karmada needed to understand Flink's internal job state machine and, more importantly, preserve the job ID—a critical piece of information that links a running Flink job to its previous state (checkpoints/savepoints). Without preserving the job ID, the new Flink deployment in a healthy cluster could not locate and restore from the latest state snapshot.
- Development of State Preservation Enhancement: This limitation spurred a collaboration with the Karmada community, leading to the development of a state preservation enhancement. This new feature, integrated into Karmada's existing failover API and configured via the propagation policy, allows users to specify a JSON path to extract specific data from the application's status (e.g., Flink's job ID) and inject it as a label into the new deployment during failover. This ensures the necessary metadata is carried over.
- Leveraging Kverno for State Injection: To bridge the gap between the preserved metadata (as a label) and Flink's consumption mechanism, Kverno webhooks were introduced. A mutating webhook intercepts the new Flink deployment request, reads the injected job ID label, and dynamically sets the
initialSavepointPathin the Flink deployment specification. This crucial step enables the Flink operator to start the job from its previous state in the new cluster.
These findings collectively describe the evolution from a manual, single-cluster recovery model to an automated, multi-cluster, state-aware failover system, significantly enhancing the resiliency and operational efficiency of Bloomberg's critical Flink data pipelines.
Technical Deep Dive
▶ Watch: Flink job deployment and management on Kubernetes (4:00)
The journey to achieve stateful Flink failover on Karmada involved several intricate technical steps, building upon Karmada's native capabilities and extending them for Flink-specific requirements.
At its core, Karmada acts as a centralized control plane for managing applications across a federation of Kubernetes clusters. It provides a unified authentication endpoint and an API server where users can apply resources as if they were interacting with a single cluster. Karmada then uses propagation policies to determine how these resources should be distributed and scheduled to member clusters. These policies can incorporate advanced features like resource-aware scheduling and cluster affinity rules, allowing platform owners to define intelligent placement strategies. For instance, Bloomberg configured a propagation policy for Flink deployments to ensure they are scheduled to only a single cluster (maxOneCluster), with all replicas co-located.
Karmada's automated cross-cluster failover mechanism relies on two main health tracking components:
- Cluster Health: Karmada continuously monitors the health of member clusters by calling their existing Kubernetes health endpoints. If a cluster is deemed unhealthy for a configured grace period, Karmada triggers a taint-based eviction, rescheduling applications to a healthier cluster.
- Application Health: For more granular control, Karmada offers a framework for resource interpretation. Users can define how their custom resources (like Flink deployments, which are Custom Resource Definitions or CRDs) should be evaluated for health. This is crucial because Karmada, by default, has no inherent understanding of Flink's internal state.
The initial integration faced a significant hurdle: when Karmada detected an unhealthy Flink job and rescheduled it, the new deployment would start from scratch. This was because Karmada lacked awareness of Flink's state and, critically, the job ID required to restore from a previous checkpoint or savepoint. To address this, a deep understanding of Flink's internal job state machine was necessary. The speakers provided a simplified model:
- Healthy States:
Running,Failed,Finished,Cancelled,Suspended. These represent a stable state, whether active or terminated. - Ephemeral States (Conditionally Healthy/Unhealthy):
Reconciling,Initializing,Created. If a job is stuck in one of these states without an associated user error (e.g., bad image path), it's considered unhealthy. However, if it's a normal transition (e.g.,Reconcilingas the job manager starts), it's healthy. - Short-Lived Transition States (Healthy):
Restarting,Failing,Cancelling. These are temporary states that quickly lead to a terminal state and are treated as healthy during their brief duration.
The health interpreter for Flink deployments was refined to consider not only the job state but also the error field within the status. This allowed for distinguishing between transient issues, user configuration errors (e.g., bad container images, malformed YAML), and actual runtime problems (application bugs, bad upgrades) that might warrant a failover.
The core technical innovation for state preservation was the introduction of a state preservation enhancement within Karmada's failover API, configured directly in the propagation policy. This enhancement consists of two fields:
jsonPath: An expression (e.g.,jobStatus.jobID) that identifies the specific piece of data to extract from the application's status.aliasLabelName: The name of the label under which this extracted data will be injected into the rescheduled resource.
This mechanism ensures that when Karmada detects a failure and prepares to delete the resource from the unhealthy cluster and schedule it to a new one, it first extracts the Flink job ID from the existing job's status. This job ID is then injected as a label (e.g., flink.apache.org/job-id: <extracted-job-id>) into the new Flink deployment manifest that Karmada creates for the healthy target cluster.
The final piece of the puzzle involved making the Flink operator in the member cluster aware of this preserved state. This was achieved using Kverno webhooks. Kverno is a policy management tool for Kubernetes that allows declarative definition of policies for validating and mutating resources. A mutating webhook was implemented to intercept the Flink deployment request as it arrived at the healthy member cluster's API server. This webhook would:
- Inspect the incoming Flink deployment manifest for the specific label injected by Karmada (e.g.,
flink.apache.org/job-id). - If the label is present, extract the job ID.
- Mutate the Flink deployment spec to set the
initialSavepointPathfield, pointing it to the latest checkpoint or savepoint associated with the extracted job ID. This effectively tells the Flink operator to restore the job from its previous state rather than starting fresh.
This multi-faceted approach, combining Karmada's orchestration capabilities with Flink-specific health interpretation, state preservation, and Kverno-driven mutation, provided a comprehensive solution for automated, stateful cross-cluster failover.
Demo / Proof of Concept
▶ Watch: Challenges managing Flink applications at multicluster scale (6:00)
While the talk did not feature a live coding demonstration, the speakers presented a detailed architectural diagram illustrating the finalized flow for stateful failover, effectively serving as a proof of concept for their implementation. This diagram elucidates how various components—Karmada, Flink applications, Kverno webhooks, and Kubernetes clusters—interact to achieve seamless recovery.
The demonstrated workflow begins with a Flink job successfully running on a healthy cluster, say Cluster B. At some point, the application experiences an issue and transitions into a reconciling state. This could indicate a job manager crash or an inability for the Flink operator to retrieve status, suggesting a potential problem beyond transient in-cluster failures.
Here’s a step-by-step breakdown of the demonstrated failover process:
- Issue Detection: Karmada, continuously monitoring the application's health using its custom Flink health interpreter, detects that the Flink job has transitioned to an
unhealthystate (e.g.,reconcilingwithout an associated error). - Toleration Period: Karmada respects the
tolerationSecondsdefined in the resource's propagation policy. This grace period allows for minor, self-correcting issues or initial application startup times. If the application remains unhealthy beyond this threshold, failover is triggered. - State Preservation: Upon triggering failover, Karmada first extracts the critical job ID from the
jobStatus.jobIDfield of the Flink application's status in Cluster B. This extraction is configured via thejsonPathin the state preservation enhancement within the propagation policy. - Resource Deletion: Karmada then deletes the Flink deployment resource from the unhealthy Cluster B.
- Rescheduling and Label Injection: Karmada intelligently schedules the Flink deployment to a new, healthy cluster (e.g., Cluster A), considering resource availability and other policy constraints. Crucially, during this rescheduling, Karmada injects the previously extracted job ID as a label (e.g.,
flink.apache.org/job-id: <extracted-job-id>) into the new Flink deployment manifest. This is where thealiasLabelNamefrom the state preservation enhancement comes into play. - Webhook Interception and Mutation: As the new Flink deployment request arrives at the API server of Cluster A, it is intercepted by the Kverno mutating webhook.
- Initial Savepoint Path Injection: The Kverno webhook reads the
flink.apache.org/job-idlabel from the deployment manifest. Using this job ID, it then dynamically modifies the Flink deployment'sspecto include theinitialSavepointPathparameter, pointing to the latest checkpoint or savepoint associated with that specific job ID. - Stateful Recovery: The Flink operator in Cluster A receives the mutated deployment request. With the
initialSavepointPathexplicitly set, the Flink operator knows to restore the job from its last known state, ensuring that data processing resumes seamlessly from where it left off, rather than starting anew.
This detailed sequence demonstrates a robust, automated mechanism that effectively preserves Flink's operational state across cluster boundaries, transforming a manual, error-prone recovery into an intelligent, hands-off process.
Defensive Implications
▶ Watch: Envisioning an ideal unified multicluster control plane (8:00)
The implementation of Flink on Karmada with stateful failover carries significant defensive implications for organizations operating critical data pipelines:
- Enhanced Data Pipeline Resiliency: The most direct benefit is a dramatic increase in the resiliency of stateful data pipelines. By automating cross-cluster failover, the system can withstand not just individual component failures (handled by Flink's native HA) but also partial or total failures of entire Kubernetes clusters. This ensures continuous data processing, which is critical for real-time financial insights and other time-sensitive applications.
- Reduced Downtime and Data Loss: Automated failover minimizes the mean time to recovery (MTTR) from cluster-wide outages. Instead of manual intervention that could take hours, the system can automatically detect, preserve state, and restart applications in a healthy cluster, significantly reducing downtime and preventing data loss that might occur during prolonged outages. The "exactly-once guarantees" of Flink are maintained across failovers, ensuring data integrity.
- Lower Operational Burden: Platform owners and SRE teams are freed from the reactive, high-stress task of manually recovering jobs after a cluster failure. The unified control plane provided by Karmada, combined with automated failover, reduces the operational overhead associated with managing a multitude of clusters and individual job recoveries. This allows teams to focus on proactive improvements rather than firefighting.
- Simplified Cluster Maintenance and Upgrades: Routine Kubernetes cluster maintenance, which often involves draining or taking clusters offline, becomes far less disruptive. Applications can be gracefully evicted from a cluster undergoing maintenance and automatically rescheduled to healthy ones without requiring heavy coordination or user involvement. This allows for more frequent and less risky updates, improving overall platform security and stability.
- Built-in Disaster Recovery (DR) Testing: The automated failover mechanism inherently provides a robust framework for disaster recovery testing. Organizations can simulate cluster failures or perform planned evictions to validate their DR strategies out-of-the-box, gaining confidence in their ability to recover from major incidents without impacting production.
- Flexible Failover Policies: The ability to define custom
tolerationSecondsand resource interpretation rules allows organizations to tailor failover sensitivity to specific application requirements. Some critical jobs might demand aggressive failover, while others can tolerate longer recovery times. This flexibility enables a nuanced approach to resilience based on business impact.
- Unified Management and Security: Karmada provides a single point of control for managing applications across clusters, simplifying authentication, authorization, and policy enforcement. This unified approach enhances security posture by reducing the attack surface and complexity associated with managing disparate cluster environments.
By integrating Karmada and extending its capabilities with Flink-specific state preservation, Bloomberg has established a resilient architecture that significantly mitigates the risks associated with infrastructure failures, ensuring the continuous and reliable operation of their crucial financial data pipelines.
Key Takeaways
- Karmada enables resilient multi-cluster Kubernetes operations: It provides a unified control plane for managing applications across a federation of clusters, offering intelligent scheduling and automated cross-cluster failover.
- Stateful application failover requires custom extensions: While Karmada offers basic failover, preserving application-specific state (like Flink's Job ID) across cluster boundaries necessitates custom enhancements, such as Karmada's state preservation feature and external tools like Kverno.
- Flink's internal state machine is crucial for health interpretation: Accurately determining the health of a Flink job for failover requires understanding its various internal states (e.g.,
Running,Reconciling,Created) and distinguishing between transient issues and persistent problems. - Kverno webhooks bridge state preservation and application consumption: Mutating webhooks, like those provided by Kverno, are essential for transforming preserved metadata (e.g., a Job ID label) into a format that the application's operator (e.g., Flink operator) can consume to restore state (e.g.,
initialSavepointPath). - Automated cross-cluster failover significantly reduces operational burden: This solution minimizes manual intervention for recovery, simplifies cluster maintenance, and enables robust disaster recovery testing, freeing up platform teams to focus on strategic initiatives.
- Careful tuning of failover parameters is critical: Parameters like
tolerationSecondsmust be meticulously configured to avoid premature or continuous rescheduling, especially considering application initialization times.
About the Speaker(s)
Michas Szacillo is a Senior Software Engineer and Tech Lead at Bloomberg. He works on the streaming platform team, which is responsible for providing Apache Flink to numerous users within Bloomberg. His work focuses on enhancing the platform's capabilities, particularly in areas like supporting stateful application failover and collaborating with open-source communities like Karmada.
Wang Li is a Software Engineer at Bloomberg, also part of the streaming platform team. Alongside Michas Szacillo, he contributes to the development and maintenance of Bloomberg's Apache Flink streaming platform, working on ensuring its reliability, efficiency, and resilience for critical financial data processing.
Reviews
Dr. Zero (Offensive Security Researcher) — MUST SEE
This talk presents an exceptional deep dive into building truly resilient, stateful data pipelines using Apache Flink on multi-cluster Kubernetes, leveraging Karmada and Kverno. The Bloomberg team's collaboration with the Karmada community to engineer custom state preservation and injection mechanisms for Flink's job ID represents a significant, production-grade defensive innovation, solving a critical operational challenge of automated cross-cluster failover for stateful applications. This is not theoretical fluff; it's a meticulously detailed, hard-won solution that will directly inform anyone struggling with similar large-scale, high-availability requirements.
Heather Calloway (CISO) — STRONG ACCEPT
This session from Bloomberg presents a compelling, production-grade solution for achieving stateful failover for critical Apache Flink data pipelines across multiple Kubernetes clusters using Karmada and Kverno. It directly addresses a significant institutional risk: the operational burden and potential data loss associated with manual recovery from full cluster failures. The detailed technical blueprint for automated resilience provides a clear path for organizations managing high-stakes, stateful streaming applications.