Challenges of and Solutions for Migrating Spark From Legacy Hadoop Clu... Neha Singla & Rasik Pandey

Neha Singla, Rasik Pandey

KubeCon + CloudNativeCon Europe 2025 · Session

Overview

This talk, presented by Neha Singla and Rasik Pandey from Apple, delves into the intricate journey of migrating large-scale Apache Spark workloads from traditional bare-metal Hadoop clusters to a modern Kubernetes-based infrastructure. The presentation meticulously details the architectural evolution, the myriad challenges encountered at each stage, and the innovative solutions implemented to achieve a highly efficient, cost-effective, and interactive Spark environment. Given Apple's immense scale and sophisticated data engineering needs, their experiences offer invaluable insights for organizations grappling with similar migrations, particularly concerning the unique demands of interactive Spark applications.

Watch on YouTube

Visual summary for Challenges of and Solutions for Migrating Spark From Legacy Hadoop Clu... Neha Singla & Rasik Pandey by Neha Singla, Rasik Pandey
Visual summary for Challenges of and Solutions for Migrating Spark From Legacy Hadoop Clu... Neha Singla & Rasik Pandey by Neha Singla, Rasik Pandey

Key moments

  1. 0:00 Introduction: Apple's Spark migration to Kubernetes journey
  2. 2:00 Special requirements for interactive Spark workloads
  3. 4:00 Baseline: Traditional Spark on bare metal with YARN
  4. 6:00 Challenges of Spark on bare metal: scaling, coupling, fault tolerance
  5. 7:00 First evolution: Virtualized environment with Mesos
  6. 8:00 Challenges with the virtualized Mesos environment

Challenges of and Solutions for Migrating Spark From Legacy Hadoop Clusters to Kubernetes

Speakers: Neha Singla, Software Engineer, Apple; Rasik Pandey, Apple

Conference: KubeCon EU

YouTube: https://www.youtube.com/watch?v=Q2ct5OXQ8fU

Overview

This talk, presented by Neha Singla and Rasik Pandey from Apple, delves into the intricate journey of migrating large-scale Apache Spark workloads from traditional bare-metal Hadoop clusters to a modern Kubernetes-based infrastructure. The presentation meticulously details the architectural evolution, the myriad challenges encountered at each stage, and the innovative solutions implemented to achieve a highly efficient, cost-effective, and interactive Spark environment. Given Apple's immense scale and sophisticated data engineering needs, their experiences offer invaluable insights for organizations grappling with similar migrations, particularly concerning the unique demands of interactive Spark applications.

The core problem addressed is the inherent mismatch between traditional Kubernetes scheduling paradigms, optimized for microservices, and the resource-intensive, often long-running, and dynamic nature of Spark jobs. The speakers highlight the critical need for features like dynamic resource allocation, fair sharing, gang scheduling, and robust monitoring in a multi-tenant environment. Their narrative spans several generations of infrastructure, from YARN on bare metal to Mesos on virtualized environments, culminating in a sophisticated Kubernetes deployment augmented by the YuniKorn scheduler.

The talk is particularly relevant for data engineers, DevOps professionals, and platform architects who are looking to modernize their big data processing infrastructure. It underscores the importance of a phased migration strategy, careful resource management, and the selection of appropriate scheduling technologies to harness the full potential of Kubernetes for data-intensive workloads. The insights shared by Apple's engineering team provide a practical roadmap for overcoming common pitfalls and optimizing Spark for both batch and interactive use cases in a cloud-native ecosystem.

Background

▶ Watch: Introduction: Apple's Spark migration to Kubernetes journey (0:00)

Apple's journey to a cloud-native Spark environment began with a traditional, bare-metal Hadoop cluster setup, a common architecture for big data processing in the past. This baseline configuration featured a classic NameNode, Resource Manager (YARN), and worker nodes running directly on physical servers. YARN, a core component of the Hadoop ecosystem, was responsible for job scheduling, resource allocation, and monitoring, distributing Spark workloads across available worker nodes. A significant advantage of this setup was YARN's external shuffle service, which facilitated efficient data exchange between executors, enhancing fault tolerance and reducing memory overhead.

However, this traditional model presented significant limitations. Resource provisioning was static and manual, requiring administrators to pre-allocate resources for peak demands, leading to frequent underutilization. Storage and compute were tightly coupled, with each cluster having exclusive access to its HDFS storage. This necessitated complex, distributed copying jobs and manual SR (Site Reliability) support to maintain data consistency across multiple clusters, effectively preventing the development of a unified data lake concept. Fault tolerance was rudimentary; a single job failure could destabilize the entire cluster. Furthermore, the lack of built-in orchestration meant shipping entire Spark distributions manually, a process prone to errors and connectivity issues. For interactive analytics, users were confined to a basic shell interface, severely limiting data scientist productivity.

The first evolutionary step involved migrating to a virtualized environment using VMware. This introduced a layer of virtualization, offering greater flexibility and on-demand scalability, albeit not fully dynamic. The resource scheduler transitioned from YARN to a Mesos-based system, aligning with their new deployment model. While this stage introduced Jupyter Classic for improved interactivity and laid the groundwork for a data lake concept with HDFS clusters gaining network connectivity and supporting Hive Metastore and Iceberg tables, it still faced substantial challenges. Resource management remained problematic, with persistent issues of overallocation and underutilization. Static resource allocation on Mesos meant sessions could not shrink, leading to inefficient use of memory and CPU. Kernel configurations for Jupyter were not persistent or shareable, forcing users to repeatedly define Spark properties for similar experiments. This intermediate stage highlighted that while virtualization offered some gains, it did not fully address the dynamic and efficient resource management needs of interactive Spark.

Key Findings

▶ Watch: Baseline: Traditional Spark on bare metal with YARN (4:00)

The primary finding from Apple's extensive migration journey is that while Kubernetes provides a robust platform for containerized applications, its default scheduling mechanisms are fundamentally ill-suited for the unique demands of interactive, data-intensive Spark workloads. The evolution from bare metal to VMware and then to Kubernetes revealed a progressive understanding of these challenges and led to the crucial insight that a specialized, application-aware scheduler is indispensable for achieving optimal performance, resource utilization, and user experience for Spark on Kubernetes.

The migration demonstrated that separating storage and compute is a foundational requirement, transitioning from tightly coupled HDFS to cloud storage or distributed storage systems accessible via network. This disaggregation unlocked significant flexibility. Furthermore, containerization and dependency management were critical objectives, streamlining deployment and reducing errors.

A key discovery was the inherent limitations of the Kubernetes default scheduler (kube-scheduler) for Spark. Designed for microservices, it lacks crucial features like queuing support, gang scheduling (where all executors for a job must start simultaneously), fairness, preemption, and resource guarantees necessary for interactive, multi-tenant Spark environments. This led to inefficient resource utilization, poor user experience due to launch latency, and a lack of control over job prioritization.

The most significant architectural finding and solution was the integration of YuniKorn as a specialized scheduling layer on top of Kubernetes. YuniKorn effectively bridged the gap by providing YARN-like scheduling capabilities, offering hierarchical resource management, guaranteed resource quotas per queue, priority preemption, FIFO (First-In, First-Out), fair sharing, and resource borrowing. This combination delivered the best of both worlds: Kubernetes' standardization and orchestration alongside YuniKorn's advanced, application-aware scheduling.

Finally, the development of a Jupyter Lab plug-in for kernel configuration management was a critical user experience improvement. This integrated solution addressed the friction of external configuration systems, allowing data scientists to manage all Spark properties—from basic to advanced, and shared to user-specific—directly within their familiar Jupyter Lab environment. This integration, combined with enhanced kernel connectivity monitoring, significantly improved the seamlessness and transparency of interactive Spark experiments, ultimately leading to significant cost savings through better resource utilization.

Technical Deep Dive

▶ Watch: Challenges of Spark on bare metal: scaling, coupling, fault tolerance (6:00)

Apple's architectural evolution for Spark workloads provides a comprehensive technical blueprint for migrating from legacy systems to a highly optimized Kubernetes environment.

1. Bare Metal with YARN (Baseline):

Initially, Spark ran on physical servers managed by YARN (Yet Another Resource Negotiator), the resource manager within the Hadoop ecosystem.

  • Architecture: Classic setup with a NameNode, Resource Manager (YARN), and worker nodes. No virtualization layer.
  • Resource Management: YARN handled job scheduling and resource allocation, distributing Spark workloads. It offered fine-grained access control over CPU and memory via a queue-based architecture.
  • Shuffle Service: Utilized YARN's external shuffle service for efficient data exchange, enhancing fault tolerance and reducing memory overhead.
  • Storage: Tightly coupled HDFS storage, isolated per cluster. Data movement between clusters required manual distributed copying jobs.
  • Challenges: No dynamic scaling, manual node provisioning, error-prone Spark distribution shipment, lack of a data lake, basic shell interface for interactivity, and no built-in orchestration.

2. Virtualized Environment with Mesos (First Evolution):

The migration to VMware introduced virtualization and a new scheduler.

  • Architecture: Virtualized servers, Mesos-based resource scheduler.
  • Resource Management: Mesos replaced YARN. However, it initially supported only static resource allocation, leading to underutilization.
  • Interactivity: Introduced Jupyter Classic, a significant improvement for data scientists.
  • Data Lake Concepts: HDFS clusters gained network connectivity, allowing Spark jobs to access multiple HDFS clusters. Hive Metastore and Iceberg tables were introduced for table-level queries.
  • Challenges: Resource overallocation and underutilization persisted due to static allocation. Inefficient t-shirt sizing (fixed CPU-to-memory ratios) led to resource contention. Jupyter kernel configurations were not persistent or shareable, causing friction for users.

3. Kubernetes with kube-scheduler (Second Evolution):

This marked the move to a fully containerized environment, leveraging Kubernetes.

  • Architecture: Spark clusters running inside a Kubernetes cluster. Spark driver runs in one Kubernetes Pod, and Spark executors run in separate Pods, all managed by Kubernetes.
  • Resource Management: The kube-scheduler, Kubernetes' default scheduler, was initially used. It dynamically managed Pods, scaling resources up or down.
  • kube-scheduler Characteristics: Designed primarily for microservice-type workloads, not data processing. It uses a spread scheduling approach to distribute Pods evenly. Features like taint toleration, node affinity rules, and node selectors were available for basic placement.
  • Interactivity: Introduced Jupyter Lab, an in-browser IDE, allowing data scientists to create notebooks with PySpark, Scala Spark, or pure Python kernels.
  • Kernel Configuration Management: A custom system was built outside Jupyter Lab to manage Spark properties hierarchically, supporting different Spark versions and enabling administrators to hide advanced configurations from general users.
  • Challenges: The kube-scheduler proved inadequate for interactive Spark. It lacked application-aware scheduling, queuing support, gang scheduling (critical for frameworks like XGBoost), fairness, preemption, reservation, and resource guarantees. This resulted in poor resource utilization (lack of bin packing) and a suboptimal user experience. The external kernel configuration system created friction.

4. Kubernetes with YuniKorn (Current Architecture):

The ultimate solution involved integrating a specialized scheduler, YuniKorn, to overcome the kube-scheduler's limitations.

  • Architecture: YuniKorn sits as a specialized layer between the Kubernetes control plane and Spark applications. It replaces the default kube-scheduler.
  • YuniKorn Scheduler Features: Provides YARN-like scheduling capabilities with Kubernetes benefits. Offers hierarchical quota management with guaranteed quotas per queue. Supports priority preemption, FIFO, fair sharing, gang scheduling, and resource borrowing. It works with any workload via Pod annotations and requires no custom resource definitions.
  • Node Pool Strategy: Workloads are divided into three distinct node pools for optimal management:
  1. Node Pool 1 (Jupyter Server): Long-running Jupyter servers, managed by kube-scheduler (as they don't require dynamic scaling or preemption).
  2. Node Pool 2 (Spark Driver): Spark drivers, which are the interactive kernels users connect to. These are critical and non-preemptable, requiring a stable environment.
  3. Node Pool 3 (Spark Executors): Spark executors, which are preemptable and can scale independently. This allows for high churn and dynamic allocation based on workload demands.
  • Benefits of Node Pools: Provides workload separation, grouping similar configurations, and streamlining lifecycle management. Placing drivers and executors in the same network zone improved shuffle latency and overall interactive experience.
  • Jupyter Lab Enhancements: A Jupyter Lab plug-in was developed for integrated kernel configuration management, allowing users to tweak Spark properties directly within the IDE, reducing friction. Enhanced kernel connectivity provided better visibility into the Spark driver's status.
  • Shuffle Service: The talk notes a current reliance on shuffle tracking rather than a fully developed external shuffle service for Kubernetes, indicating an ongoing area of improvement. Challenges with shuffle tracking included fetch failures and executors staying around longer than required, requiring careful tuning of cache timeouts and dynamic allocation properties.
  • Ephemeral Storage: For ephemeral storage requirements of Spark jobs, the recommendation is to use local disk (e.g., emptyDir with specific disk configurations) rather than Persistent Volumes (PVs), as PVs significantly reduced throughput (from 200-300 pods/sec to 50 pods/sec in their experience).
  • CPU/Memory Configuration: A critical learning was to set CPU requests equal to limits to prevent throttling and ensure fair resource sharing, avoiding the "bad neighbor" problem where a pod with high limits could monopolize resources.

This detailed progression highlights Apple's commitment to continuously refining its big data infrastructure, adapting to new technologies, and addressing specific pain points to deliver a robust and efficient platform for its data scientists and engineers.

Demo / Proof of Concept

▶ Watch: First evolution: Virtualized environment with Mesos (7:00)

While the talk did not feature a live, interactive demonstration of the entire Spark on Kubernetes with YuniKorn architecture, the speakers provided a crucial visual insight into the user experience, specifically showcasing the Jupyter Lab kernel configuration interface.

This interface represents a significant proof of concept for improving data scientist productivity and reducing friction in managing complex Spark properties. A screenshot was presented, illustrating how users can now configure their Spark kernels directly within the Jupyter Lab environment. This custom-built plug-in for Jupyter Lab addresses a key challenge faced in earlier stages, where users had to manage Spark properties in an external system, then apply them in Jupyter Lab.

The interface is designed with user experience in mind, offering different "presets." For instance, a "minimal preset" exposes only a smaller, essential set of properties, preventing less experienced users from being overwhelmed. Advanced users or administrators can access a "full preset" to view and tweak a comprehensive range of advanced Spark configurations. This hierarchical and user-centric approach to configuration management, integrated into the familiar Jupyter Lab IDE, demonstrates a tangible improvement in the workflow for data scientists at Apple. The speakers also mentioned plans to open-source this Jupyter Lab plug-in in the future, indicating its value as a generalized solution.

Defensive Implications

▶ Watch: Challenges with the virtualized Mesos environment (8:00)

While this talk isn't about traditional security vulnerabilities, the "defensive implications" here translate to implementing robust, efficient, and resilient practices for running Spark workloads on Kubernetes in an enterprise environment. The strategies outlined by Apple are critical for protecting against resource exhaustion, ensuring data integrity, maintaining service availability, and managing operational costs effectively.

  1. Strategic Scheduler Selection: The most significant defensive measure is to recognize that the default kube-scheduler is insufficient for data-intensive, interactive Spark workloads. Organizations must adopt a specialized batch scheduler like YuniKorn. This defends against:
  • Resource Monopolization: YuniKorn's hierarchical quota management and fair sharing prevent any single user or job from monopolizing cluster resources.
  • Poor Performance/SLA Violations: Features like gang scheduling, priority preemption, and resource guarantees ensure that critical jobs get the resources they need, when they need them, preventing performance bottlenecks and improving interactive response times.
  • Inefficient Utilization: Bin packing capabilities optimize resource allocation, reducing waste and associated costs.
  1. Phased Migration and Testing: A "design a phased strategy" approach is crucial. Incrementally migrating workloads and ensuring rigorous testing on both old and new infrastructures minimizes risks of service disruption and data inconsistencies. This defends against catastrophic failures during large-scale infrastructure changes.
  1. Optimized Resource Allocation and Dynamic Scaling:
  • Dynamic Allocation: Enabling dynamic allocation for Spark executors (where applicable) is a key defense against over-provisioning and under-utilization, leading to significant cost savings. However, understanding when not to use it (e.g., for gang-scheduled jobs like XGBoost) is equally important.
  • Request vs. Limit Configuration: Setting CPU requests equal to limits is a critical best practice. This prevents pods from being throttled and ensures fair resource sharing, acting as a defense against "noisy neighbor" problems that degrade performance across the cluster.
  • Preemption Policies: Implementing preemption policies ensures that high-priority jobs can acquire necessary resources by evicting lower-priority ones, defending against critical job starvation.
  1. Robust Storage and Networking Strategy:
  • Disaggregated Storage: Migrating from local HDFS to distributed cloud storage (e.g., S3-like systems) is vital. This defends against data silos and enables flexible, scalable access to data from multiple compute clusters.
  • Ephemeral Storage Considerations: For high-throughput ephemeral storage needs (e.g., Spark shuffle data), prioritize local disk (e.g., emptyDir volumes) over Kubernetes Persistent Volumes (PVs) to avoid severe performance degradation. This is a defense against I/O bottlenecks.
  • Network Optimization: Careful consideration of network architecture, especially between Spark drivers and executors (e.g., placing them in the same network zone), can significantly reduce shuffle latency, improving performance and stability.
  1. Comprehensive Monitoring and Observability: Setting up extensive monitoring, including YuniKorn's UI metrics, is non-negotiable. This provides visibility into scheduling efficiency, resource usage, and potential bottlenecks. It's a proactive defense against operational issues, allowing for rapid detection and resolution of problems.
  1. Integrated Configuration Management: Integrating kernel configuration management directly into tools like Jupyter Lab reduces user error and friction. This defends against inconsistent configurations across users and experiments, promoting standardization and reproducibility.

By adopting these defensive postures, organizations can build a resilient, high-performance, and cost-effective Spark platform on Kubernetes, capable of supporting demanding data engineering and data science workloads.

Key Takeaways

  • Assess Workload Patterns: Crucially, understand whether your Spark workloads are batch, interactive, or AI-driven, as this dictates the required scheduling capabilities.
  • Phased Migration Strategy: Implement a hybrid, incremental migration approach with robust test plans to ensure smooth transitions and minimize disruption.
  • Optimize Resource Allocation with a Batch Scheduler: The default Kubernetes scheduler is not suitable for batch jobs; leverage a specialized scheduler like YuniKorn for features like queuing, fairness, preemption, and dynamic allocation.
  • Enable Dynamic Allocation (with exceptions): For most interactive workloads, dynamic allocation significantly improves resource efficiency and cost-effectiveness, though specific cases like XGBoost may require static allocation or gang scheduling.
  • Configure Pod Scheduling and Resource Limits Appropriately: Utilize Kubernetes features like node affinity, node selectors, and taint tolerations for optimal pod placement. Crucially, set CPU requests equal to limits to ensure fair resource sharing and prevent throttling.
  • Address Storage and Networking Challenges: Migrate from local disk storage to distributed storage. Carefully plan for networking, especially for shuffle services, and consider local disk for high-throughput ephemeral storage needs to avoid performance bottlenecks with Persistent Volumes.
  • Implement Comprehensive Monitoring: Establish robust monitoring (e.g., YuniKorn UI metrics) to track scheduling efficiency, resource utilization, and quickly identify issues.

About the Speaker(s)

Neha Singla is a Software Engineer at Apple. Her work focuses on addressing the unique requirements of interactive Spark workloads within Apple's infrastructure, including developing solutions for low-latency responses, robust connectivity, and fair resource allocation in multi-tenant environments. She played a key role in the technical evolution of Apple's Spark platform, particularly with the integration of Jupyter Lab and advanced kernel configuration management.

Rasik Pandey is also from Apple. While his specific title isn't mentioned in the transcript, his opening remarks and contributions to the discussion suggest a leadership or architectural role within Apple's data platform team. He provided the overarching context for Apple's migration journey and emphasized the strategic objectives behind moving Spark workloads to Kubernetes, including cost reduction and improved resource utilization.

Reviews

Dr. Zero (Offensive Security Researcher) — MUST SEE

This talk from Apple details a highly valuable and technically deep journey of migrating large-scale interactive Spark workloads from legacy Hadoop to Kubernetes, specifically addressing the critical limitations of kube-scheduler for such applications. The speakers present a comprehensive architectural evolution, highlighting the indispensable role of a specialized scheduler like YuniKorn for achieving optimal performance, resource utilization, and user experience. The insights, including a custom Jupyter Lab plugin for kernel configuration and a sophisticated node pool strategy, provide a practical and actionable roadmap for any organization tackling similar challenges at scale.

Heather Calloway (CISO) — STRONG ACCEPT

This session from Apple delivers a highly credible and actionable account of migrating large-scale Spark workloads to Kubernetes, highlighting critical architectural decisions and their profound impact on operational efficiency and business costs. While deeply technical, the speakers effectively translate complex infrastructure challenges into clear solutions that enhance platform resilience, optimize resource utilization, and significantly improve data scientist productivity, addressing core institutional risks around data processing and resource management. It offers invaluable insights for any organization seeking to modernize its big data infrastructure and ensures critical data…

→ Top-rated talks at KubeCon + CloudNativeCon Europe 2025

All talks from KubeCon + CloudNativeCon Europe 2025