Chapter 6 - Large-Scale Resource Management

Updated 4 Oct 2026

Contents

  1. Mesos: A Platform for Fine-Grained Resource Sharing in Data Center
    • Authors: B. Hindman, A. Konwinski, M. Zaharia, A. Ghodsi, A. D. Joseph, R. Katz, S. Shenker, I. Stoica
    • Institution: University of California, Berkeley
    • Conference: NSDI 2011
  2. Omega: Flexible, Scalable Schedulers for Large Compute Clusters
    • Authors: Malte Schwarzkopf, Andy Konwinski, Michael Abd-El-Malek, John Wilkes
    • Conference: EuroSys 2013

Background

Data Center Infrastructure

  • Data centers are built from clusters of commodity hardware
    • Commodity hardware = standard, off-the-shelf servers (not specialized equipment)
    • Cost-effective approach to building large-scale infrastructure

Analogy: Think of commodity hardware like IKEA furniture - standardized, mass-produced, affordable components that can be assembled at scale, versus custom-made furniture (specialized servers) that's expensive and harder to replace.

  • Clusters run a diverse set of applications
    • MapReduce (batch processing)
    • Data Streaming frameworks
    • Storm, S4 (real-time processing)
    • Each application has its own execution framework
  • Each application has its own execution framework
    • Different workload characteristics
    • Different resource requirements
    • Different scheduling needs
  • Multiplexing a cluster between frameworks improves resource utilization and reduces costs
    • Better hardware utilization
    • Reduced idle resources
    • Lower operational costs

Real-world Usage:

  • Google uses cluster sharing extensively across Search, Gmail, YouTube, and Maps services
  • Netflix runs multiple frameworks (Spark, Presto, Flink) on shared AWS/Azure clusters
  • Uber multiplexes real-time analytics, batch processing, and ML training on the same infrastructure

Problem Statement

Core Challenges

  • Single framework approach is inefficient
    • It is very difficult and challenging for a single framework (application-specific) to efficiently manage the resources of clusters
    • Different applications have different needs
    • Hard to optimize for all use cases simultaneously
  • The aim of the Mesos work is to run multiple frameworks in a single cluster to:
    • Maximize utilization - Keep resources busy, avoid idle capacity
    • Share data between frameworks - Avoid data duplication and transfer costs

Mesos = scheduler

Analogy: Imagine a kitchen (cluster) with multiple chefs (frameworks). Having one chef cook everything is inefficient - a pastry chef might waste time cooking steaks, and a grill master struggles with delicate desserts. Better to let specialists work in the same kitchen, sharing the equipment.


Common Solutions for Sharing Clusters

1. Static Partitioning (Coarse-Grained Sharing)

Approach:

  • Statically partition the cluster
  • Run one framework (or VMs) per partition
  • Fixed allocation that doesn't change

IMAGE: Coarse-Grained Sharing diagram - shows Framework 1, Framework 2, Framework 3 each with dedicated machines above Storage System (HDFS)

Disadvantages:

  1. Low utilization
    • Resources allocated but not always used
    • No dynamic reallocation
  2. Rigid partitioning
    • Doesn't match real-time changing application demands
    • Cannot adapt to workload variations
    • Over-provisioning in some partitions, under-provisioning in others
    • เวลาผ่านไป framework อื่น อาจจะต้องการ resource มากขึ้น
  3. Mismatch between allocation granularity
    • Static partitioning allocation granularity doesn't match current application frameworks (like Hadoop)
    • Applications may need flexible, fine-grained resource allocation

Corase-grained, can just share resource to the application level
Fine-grained we can share the resource to task → Maximize utilization (CHECK!!?)

Real-world Example: Early Hadoop deployments often used static partitioning - a company might allocate 30% of cluster to production jobs, 40% to analytics, 30% to development. This led to wasted resources when partitions were idle.


Mesos: A Platform for Fine-Grained Resource Sharing

Overview

  • Apache Mesos is an open-source cluster manager that efficiently manages distributed computing resources across multiple machines
  • Abstracts CPU, memory, storage, and other compute resources
    • Allows applications to treat a cluster of machines as a single large resource pool
    • Resource abstraction layer

Analogy: Mesos is like a hotel concierge who manages all the rooms (resources). When guests (applications) arrive, the concierge assigns rooms based on availability and guest needs, rather than each guest having to find and book their own room directly.

Real-world Usage:

  • Twitter used Mesos to run services, analytics, and ML workloads
  • Airbnb deployed Mesos for their data infrastructure (Spark, Presto, Chronos)
  • Apple runs Siri backend services on Mesos clusters

Challenges and Goals

  1. High Utilization
    • Maximize resource usage across the cluster
    • Minimize idle capacity
    • Dynamic allocation based on demand
  2. Support Diverse Frameworks
    • Each framework has different scheduling needs
    • Different optimization goals (latency vs. throughput)
    • Save computation cost
    • Reduce time to completion
  3. Scalability
    • Must scale to clusters of 1,000s of nodes
    • Running 100s of jobs
    • With 1,000,000s of tasks
    • Scheduling system must handle high throughput
  4. Reliability
    • All applications depend on Mesos
    • System must be fault-tolerant and highly available

Key Definitions:

  • Fault-tolerant: Even though some partial failures occur, the system can still run based on existing resources/services
    • Example: Multiple copies/replications in HDFS, GFS
    • Graceful degradation
  • High Availability (HA): Fail-over architecture
    • Active Machine: Backup Machine
    • Load balancer
    • Automatic failover when primary fails

Analogy:

  • Fault-tolerant is like having spare tires in a truck - if one tire fails, you can still operate
  • High availability is like having a backup driver ready to take over immediately if the main driver is unable to continue

Mesos Definition:

A thin resource sharing layer that enables fine-grained sharing across diverse cluster computing frameworks, by giving frameworks a common interface for accessing cluster resources.


FailOver Cluster Architecture

Clustering Concept:

  • Use of multiple computers and redundant interconnections to form what appears to be a single, highly available system
  • Provides protection against downtime for important applications/services
  • Distributes workload among several computers
  • When failure occurs in one system, the service remains available on another

IMAGE: Failover Clustering diagram showing Node 1, Node 2, Cluster Storage, and Clients connected via shared bus/ISCI

Key Benefits:

  1. High availability
    • Reduces unplanned downtime
    • Increases reliability of services and applications
  2. High scalability
    • Allows administrators to assign multiple nodes to one cluster
    • Enhances performance and availability
    • Enables incremental growth

Cluster Operation:

  • All cluster nodes are aware of all other nodes' shared resources and their availability of services
  • When resource or hardware fails, clustering service automatically transfers workloads from one node to another

Real-world Example:

  • Amazon RDS uses failover clustering - if primary database instance fails, a standby replica automatically becomes the new primary within 60-120 seconds
  • Microsoft SQL Server Always On Availability Groups provide automatic failover for critical databases

Design Elements

1. Fine-Grained Sharing

Key Concept:

  • Allocation at the level of tasks within a job
    • One job can have many tasks
  • Not at the level of entire machines or partitions
  • Dynamic, flexible resource assignment

Benefits:

  • Improves utilization
  • Reduces latency (faster task scheduling)
  • Improves data locality (tasks scheduled near their data)

IMAGE: Comparison diagram showing Coarse-Grained vs Fine-Grained Sharing

Coarse-Grained Sharing:

  • Entire machines dedicated to one framework
  • Framework 1, 2, 3 each get dedicated nodes
  • Inflexible, wasteful

Fine-Grained Sharing (Mesos):

  • Tasks from multiple frameworks share nodes
  • Fw. 1, Fw. 2, Fw. 3 tasks interleaved on same machines
  • Much better utilization
  • Flexible allocation

Analogy:

  • Coarse-grained is like giving each person an entire parking lot - wasteful if they only need one space
  • Fine-grained is like a parking garage where anyone can take any available spot - much more efficient

Real-world Impact: Fine-grained sharing can improve cluster utilization from 20-30% (static partitioning) to 60-80% (dynamic sharing), potentially doubling the effective capacity of your hardware investment.


2. Resource Offers

Concept:

  • Mesos offers available resources to frameworks (VMs)
  • Frameworks choose which resources to use and which tasks to launch
  • Decentralized decision-making

Advantages:

  • Keeps Mesos simple
  • Expandable to future frameworks
  • Frameworks know their own constraints best
  • No need for Mesos to understand every framework's requirements

Disadvantages:

  • Decentralized solutions are not optimal
  • May not achieve global optimization
  • Potential for suboptimal resource allocation

Process:

  1. Mesos offers resources to a framework

    • Example: <slave_id, cpu_cores, memory_gb>
    • "Framework A, slave1 has 4 CPUs and 4GB available"
  2. Framework scheduler evaluates offer

    • Checks if resources match requirements
    • Considers data locality, constraints
  3. Framework responds

    • Accept and launch tasks
    • Reject and wait for better offer
  4. Mesos allocates resources

    • Updates cluster state
    • Instructs slave to launch tasks

Analogy: Resource offers are like a buffet - the restaurant (Mesos) presents all available dishes (resources), and diners (frameworks) choose what they want to eat. The restaurant doesn't need to know each diner's preferences; they just present options.


Mesos Architecture

IMAGE: Mesos Architecture diagram showing Hadoop scheduler, MPI scheduler, Mesos master with standby masters, and Mesos slaves with executors

Components:

  1. Mesos Master
    • Fine-grained sharing coordinator across frameworks
    • Makes resource offers to frameworks
    • Single point of coordination (with standby replicas)
  2. Standby Masters
    • Hot standby for high availability
    • Take over if primary master fails
    • ZooKeeper-based leader election
  3. Mesos Slave (on each node)
    • Reports available resources to master
    • Launches and monitors executor processes
    • Enforces resource isolation
  4. Framework Scheduler
    • Runs application-specific scheduling logic
    • Receives resource offers
    • Decides which tasks to launch on which resources
  5. Executor
    • Framework-specific process running on slave
    • Launches and monitors tasks
    • Reports task status back to scheduler

IMAGE: Mesos Architecture Overview with ZooKeeper

Zookeeper formulate the column (the majority)???

Key Interactions:

  • Leader election and service discovery via ZooKeeper
    • Ensures single active master
    • Maintains cluster configuration
  • Slaves publish available resources to master
    • Periodic heartbeats with resource availability
    • CPU, memory, disk, ports
  • Master sends resource offers to frameworks
    • Based on allocation policy (DRF, priorities)
    • Can offer to multiple frameworks
  • Scheduler replies with tasks and resources needed per task
    • Specifies: task ID, resource requirements, command to execute
    • Can reject offers
  • Master sends tasks to slaves
    • Instructs slave to launch tasks via executor
    • Updates cluster state

Real-world Architecture: This is similar to Kubernetes architecture but with different philosophy - Mesos uses "resource offers" (push model) while Kubernetes uses "requests" (pull model). Companies like Mesosphere DC/OS built enterprise platforms on top of Mesos.

Mesos Architecture - Terminology

  • Mesos master: Manages fine-grained sharing across frameworks
    • Central coordinator (but stateless - state in ZooKeeper)
  • Mesos slave on each node: Resource agent on each machine
  • Frameworks that run tasks on slave nodes: The applications using Mesos
  • Framework schedulers: Application-specific scheduling logic

Key Principle:

"While the Mesos master determines how many resources to offer to each framework, the frameworks' schedulers select which of the offered resources to use."

  • Mesos master: Resource allocation policy (fairness, priorities)
  • Framework scheduler: Resource selection strategy (constraints, locality)

Analogy: The master is like a restaurant maître d' who decides the order in which to serve customers and shows them to tables. The customers (framework schedulers) then decide whether they want that specific table or would prefer to wait for a different one.


Mesos Resource Offer Example

IMAGE: Resource offer flow diagram showing Framework 1 with Job 1 and Job 2, Framework 2, Allocation module, Mesos master, and Slaves

Example Flow:

  1. Slave 1 reports availability
    • <s1, 4cpu, 4gb, ...> → Mesos master
    • Slave has 4 CPU cores and 4GB RAM available
  2. Master sends offer to Framework 1
    • <s1, 4cpu, 4gb, ...> → Framework 1 Scheduler
    • "You can use these resources"
  3. Framework 1 responds with task requests
    • <task1, s1, 2cpu, 1gb, ...>
    • <task2, s1, 1cpu, 2gb, ...>
    • Framework accepts partial offer, launches 2 tasks
  4. Master instructs Slave 1
    • <fw1, task1, 2cpu, 1gb, ...>
    • <fw1, task2, 1cpu, 2gb, ...>
    • Slave launches tasks in executors

Remaining capacity: Slave 1 now has 1cpu, 1gb available for future offers

Mesos Resource Offers - Rejection and Filters

Rejection Handling:

  • Frameworks can reject offers when they don't satisfy constraints
    • Wait for offers that meet requirements
    • Example: "I need 8GB RAM but offer only has 4GB"

Problem:

  • Framework might wait a long time for suitable offer
  • Mesos might send unsuitable offers to many frameworks
  • Time consuming for both Mesos and frameworks

Solution: Filters

  • Mesos frameworks set their own filters
    • Specify offers that will always be rejected
    • Reduces unnecessary offer processing

Example Filters:

  • "Only offer nodes from list L"
    • Geographic constraints
    • Hardware requirements
  • "Only offer nodes with at least R resources free"
    • Minimum resource requirements
    • Example: "Only show me nodes with ≥8GB RAM"

Benefits:

  • Reduces master overhead
  • Reduces network traffic
  • Frameworks get relevant offers faster
  • Can be evaluated quickly by master

Analogy: Filters are like telling a real estate agent "don't show me houses under $500k or outside downtown" - saves everyone time by avoiding unsuitable options upfront.


Mesos Architecture Details

Resource Allocation Module

Allocation Policies:

  1. Fair Sharing
    • Based on generalization of max-min fairness for multiple resources
    • Dominant Resource Fairness (DRF) algorithm
    • Equalizes each framework's share of their dominant resource

DRF: Maximize min⁡i(diti)\boxed{\text{DRF: Maximize } \min_i \left( \frac{d_i}{t_i} \right)}

Where:

  • did_i = dominant resource share of framework ii
  • tit_i = total resources of that type
  1. Strict Priorities
    • Some frameworks get higher priority
    • Production > Development > Testing
    • Higher priority frameworks get offers first
  2. Revocation (Task Killing)
    • Mesos can kill tasks when a greedy framework uses lots of resources for a long time
    • Allows frameworks a "greedy period" initially
    • Prevents indefinite resource hoarding
    • Framework can checkpoint and resume tasks

Real-world Example: In a shared cluster, you might give ML training jobs (batch) lower priority than user-facing web services (latency-sensitive). If the cluster gets busy, training tasks can be preempted and restarted later.


Isolation

Implementation:

  • Achieved using Linux Containers (cgroups)
  • Each task runs in isolated container
  • Resource limits enforced by kernel

Isolation Guarantees:

  • CPU: CPU shares, quota enforcement
  • Memory: Hard memory limits, OOM killer
  • Disk I/O: I/O bandwidth limits
  • Network: Network bandwidth (if configured)

Modern Equivalent: This is essentially what Docker containers provide - similar lightweight isolation using cgroups and namespaces. Mesos was doing container-based isolation before Docker became popular!


Scalable and Robust Resource Offers

Design Principles:

  1. Use of filters for fast evaluation
    • "Only offer nodes from list L"
    • "Only offer nodes with at least R resources free"
    • Filters can be evaluated quickly by master
    • Reduces unnecessary communication
  2. Resource counting towards allocation
    • Mesos counts resources offered to a framework towards its allocation of cluster
    • Prevents frameworks from accumulating offers
    • Ensures fairness even during offer phase
  3. Offer timeout and withdrawal
    • If framework takes too long to respond, Mesos withdraws the offers
    • Asks another framework
    • Prevents one slow framework from blocking others

Scalability Characteristics:

  • Master is stateless (state in ZooKeeper)
  • Can handle 10,000+ slaves
  • Makes 1000s of offers per second
  • Allocation decisions are fast (O(n)O(n) in number of frameworks)

Real-world Scale: Twitter ran Mesos clusters with 10,000+ nodes and hundreds of frameworks, handling millions of tasks per day.

Resource Utilization: Static Partitioning vs. Mesos

IMAGE: Graph comparing resource utilization

Static Partitioning:

  • Hadoop: 0-33% utilization, spiky
  • Pregel: 0-33% utilization, alternating with Hadoop
  • MPI: 0-33% utilization, stable 33%
  • Average utilization: ~33% (1/3 of cluster)
  • Much wasted capacity due to non-overlapping usage

Mesos (Dynamic Sharing):

  • All frameworks share cluster dynamically
  • Combined utilization: 50-100%
  • Much better resource usage
  • Frameworks can use idle capacity from others

Cost Impact: If static partitioning uses only 33% of hardware, you need 3× the hardware to handle the same workload compared to dynamic sharing at 90% utilization. For a 1000-node cluster, that's 2000 extra servers!

Mesos Evaluation

Experiment Setup:

  • Platform: Amazon EC2
  • Cluster Size: 92 Mesos nodes
  • Mix of Frameworks:
    1. Hadoop running mix of small and large jobs
      • Based on Facebook workload trace
      • Realistic job size distribution
    2. Hadoop instance running large batch jobs
      • Long-running background processing
    3. Spark running machine learning jobs
      • Iterative algorithms
      • In-memory processing
    4. Torque running MPI jobs
      • Tightly-coupled HPC applications
      • The Terascale Open-source Resource and QUEue Manager (TORQUE)

Expected Results:

  1. Mesos should:
    • Achieve higher utilization than static partitioning
    • All jobs finish at least as fast as in static partitioning
    • No framework should be significantly penalized

Real-world Validation: This mixed workload is representative of real production clusters that run batch analytics (Hadoop/Spark), ML training, and online services simultaneously.

Mesos Performance For All Frameworks

IMAGE: Timeline graph showing cluster CPU utilization over time with color-coded frameworks

X-axis: Time (seconds), 0-1600 Y-axis: Share of Cluster CPUs (0-1.0)

Color Legend:

  • Blue: Facebook Hadoop Mix
  • Red: Large Hadoop Mix
  • Green: Spark
  • Light Blue: Torque/MPI

Key Observations:

  1. Dynamic resource allocation visible
    • Different frameworks getting resources at different times
    • No static boundaries
  2. High overall utilization
    • Cluster stays busy (50-100% utilization)
    • Little idle capacity
  3. Frameworks scale up and down
    • Around t=400s: Large Hadoop job scales up
    • Around t=600s: Job completes, resources freed
    • Other frameworks immediately use freed capacity
  4. Concurrent execution
    • Multiple frameworks running simultaneously
    • Resources shift between them smoothly

Key Insight: The cluster operates like a dynamic marketplace where resources continuously flow to whoever needs them most, rather than sitting idle in static partitions.


Mesos Performance vs Static Partitioning

[IMAGE: Graph showing Large Hadoop Mix comparison]

Static Partitioning:

  • Flat line at ~22% utilization
  • Resources allocated but not fully used
  • No elasticity

Mesos:

  • Varying utilization: 20-100%
  • Peaks when needed
  • Valleys when idle
  • Can use resources from other frameworks

Results:

MetricStatic PartitioningMesos
Avg Utilization22%60%
Job Completion TimeBaseline2.7× faster
Framework StarvationPossibleNone observed

Key Findings:

  1. Better utilization without sacrificing performance
  2. Jobs complete faster due to elastic scaling
  3. No framework starvation - fair sharing works

Business Impact: 2.7× faster completion means analytics results arrive sooner, ML models train faster, and you can run more workloads on the same hardware.


Omega: Flexible, Scalable Schedulers for Large Compute Clusters

Motivation and Challenges

Google's Context:

  • Google maintains different data centers around the world
    • Geographic distribution
    • Multiple regions for redundancy
  • Clusters and workloads are growing in size
    • 10,000+ machines per cluster
    • Exponential growth in data and computation
  • Diverse workloads
    • Batch analytics (MapReduce)
    • Serving systems (web servers)
    • ML training
    • Infrastructure services
  • Growing job arrival rates
    • 1000s of jobs arriving per minute
    • Sub-second scheduling requirements

Need: A scalable scheduler to tackle these challenges

Real-world Scale: Google processes exabytes of data monthly across thousands of workload types. Traditional monolithic schedulers become bottlenecks at this scale.

Mesos ก็มีบาง Limitation อยู่, Google เลยสร้างเองเลยดีกว่า


Omega - A Next-Generation Cluster Scheduling Approach

Definition:

The Omega Scheduler is a cluster scheduling architecture developed by Google to improve scalability, efficiency, and flexibility in multi-framework environments.

Positioning:

  • Alternative to monolithic schedulers (e.g., Hadoop YARN)
  • Alternative to two-level schedulers (e.g., Apache Mesos)
  • Introduces shared-state, optimistic concurrency approach

Core Innovation:

  • Parallel scheduling with shared state
  • Optimistic concurrency control for conflict resolution
  • Independent schedulers for each framework type

Analogy: Instead of having one traffic controller (monolithic) or a two-tier system (Mesos), Omega is like having multiple air traffic controllers working simultaneously with a shared view of all planes, each handling their own airline's flights.


The Scheduling Problem

IMAGE: Diagram showing Job with Tasks mapping to Machines

Basic Problem:

  • Input: Jobs consisting of multiple tasks
  • Resources: Cluster of machines (10,000s)
  • Goal: Assign tasks to machines efficiently

Constraints:

  • Resource requirements (CPU, RAM, disk)
  • Data locality preferences
  • Anti-affinity rules (don't place together)
  • Priority levels
  • SLAs and deadlines

Analogy: Like Tetris at scale - you have differently shaped pieces (tasks with different resource needs) falling continuously, and you need to fit them efficiently onto your playing field (cluster) without gaps or overflows.


In the Context of...

  1. Diverse workloads
    • Circle, square, hexagon, triangle symbols
    • Different resource profiles
    • Batch vs. serving vs. ML
  2. Increasing cluster sizes
    • Small cluster → Large cluster
    • Scaling to 10,000+ machines
    • Horizontal scaling challenges
  3. Growing job arrival rates
    • Graph showing upward trend
    • 1000s of jobs/minute
    • Need for fast scheduling decisions

Real-world Impact: At Google scale, even 1ms improvement in scheduling latency can mean handling 1000 more jobs per second across the fleet.


Why is it Important to Solve?

IMAGE: Scheduling bottleneck diagram

The Challenge:

  • Arriving jobs and tasks: 1,000s per interval
  • Cluster scheduler: Takes 60+ seconds to make decisions
    • Serialized decision making
    • Complex constraint solving
    • Monolithic bottleneck
  • Cluster machines: 10,000s waiting for instructions

Problems with Slow Scheduling:

  1. Resource underutilization
    • Machines idle while waiting for decisions
    • Wasted capacity
  2. Increased latency
    • Jobs wait in queue
    • User experience degraded
  3. Scalability limits
    • Scheduler becomes bottleneck
    • Cannot grow cluster further

Real-world Impact: A 60-second scheduling delay on a 10,000-node cluster means ~167 machine-hours wasted per scheduling cycle. At cloud prices (0.10/machine−hour),that′s0.10/machine-hour), that's 16.70 lost every minute!


The Actual Problem: Increasing Complexity

IMAGE: Diagram highlighting scheduler complexity

Complexity Factors:

  1. Multiple workload types with different needs
    • Red, blue, green, yellow circles (different job types)
    • Each has unique constraints
  2. Complex placement constraints
    • Affinity rules
    • Anti-affinity rules
    • Locality preferences
    • Hardware requirements
  3. Multiple scheduling objectives
    • Fairness
    • Utilization
    • Latency
    • Throughput
  4. Scale challenges
    • 10,000+ machines
    • 100,000+ tasks
    • Sub-second decision requirements

Analogy: It's like being a restaurant manager who needs to seat thousands of guests per hour, each with dietary restrictions, seating preferences, group sizes, and VIP status - all while maximizing table utilization and minimizing wait times.


Problem Statement

Core Goals:

  1. Break up the cluster scheduler into independent schedulers
    • One scheduler per workload type or team
    • Parallel operation
    • Independent evolution
  2. Arbitrate resources between schedulers
    • Fair allocation
    • Conflict resolution
    • No deadlock or starvation

Why?

  • Scalability: Parallel scheduling scales better
  • Flexibility: Each workload gets custom scheduling logic
  • Maintainability: Teams own their scheduler code
  • Innovation: Can experiment without affecting others

Existing Approaches - Comparison

IMAGE: Three scheduling architectures side by side

1. Monolithic Scheduler

IMAGE: Single scheduler box above machine grid

Characteristics:

  • Single centralized scheduler
  • All logic in one codebase
  • Serial decision making

Disadvantages:

  • ✗ Hard to diversify (one size fits all)
  • ✗ Code growth (complexity increases over time)
  • ✗ Scalability bottleneck (single decision maker)

Example: Early Hadoop YARN


2. Static Partitioning

IMAGE: Three separate partitions with S0, S1, S2 schedulers

Characteristics:

  • Cluster divided into fixed partitions
  • One scheduler per partition
  • No resource sharing between partitions

Disadvantages:

  • ✗ Poor utilization (can't use idle resources in other partitions)
  • ✗ Inflexible (can't adapt to changing demands)

Example: Early HPC cluster setups


3. Two-Level (Mesos Approach)

IMAGE: Resource manager layer below S0, S1, S2 schedulers

Characteristics:

  • Resource manager offers resources
  • Framework schedulers choose
  • Pessimistic concurrency (one offer at a time)

Disadvantages:

  • ✗ Hoarding: Frameworks may hold onto offers
  • ✗ Information hiding: No global view
    • No information about cell state
  • ✗ Offer serialization: Can be slow

Example: Apache Mesos (as discussed earlier)


Existing Approach - Mesos Revised

[IMAGE: Two-level scheduler with Resource Manager]

Disadvantages of Mesos Approach:

  1. Pessimistic Concurrency Control
    • Mesos avoids conflicts by offering a given resource to one framework at-a-time
    • Mesos chooses the order and sizes of offers
    • One framework holds the lock of the resources for the duration of the decision
    • → Slow for frameworks with complex logic
  2. Limited Information
    - A framework does not have access to all the cluster state
    • Only sees offered resources, not entire cluster
    • Cannot support preemption or gang scheduling
    • Cannot see what other frameworks are doing

When Mesos Works Well:

  • Tasks are short-lived and relinquish resources frequently
  • Job sizes are small compared to cluster size
  • Simple placement constraints

When Mesos Struggles:

  • Long-running services
  • Gang scheduling (all-or-nothing placement)
  • Complex global optimization
  • Google cluster workloads do not have these properties

Real-world Example: Try scheduling a 1000-task MapReduce job that needs all tasks on the same rack (gang scheduling) using resource offers - the framework would need to collect enough simultaneous offers for the entire job, which could take many rounds of offer/reject cycles.


Proposed Solution: Omega

IMAGE: Shared-state architecture diagram

Core Innovation: Shared-State Scheduling

Components:

  • S0,S1,S2S_0,S_1,S_2: Independent framework schedulers
    • Each has private local copy of cluster state
    • Each makes scheduling decisions independently
    • Parallel operation
  • Cluster State: Central shared state
    • Single source of truth
    • Maintained in consistent data store
    • Optimistic concurrency control

Key Principles:

  1. No centralized scheduler
    • Each framework maintains its own scheduler
    • Schedulers operate in parallel
  2. Shared state visibility
    • Each scheduler has full view of cluster
    • Can see all resources and allocations
    • Makes informed decisions
  3. Optimistic concurrency
    • Multiple schedulers make decisions simultaneously
    • Conflicts detected at commit time
    • Losers retry with updated state

Analogy: Like Google Docs collaborative editing - multiple people (schedulers) edit the same document (cluster state) simultaneously. When you hit "save," if someone else changed the same paragraph, you get notified and must merge or retry.


Omega Scheduling Example - Step by Step

Step 1: Initial State

IMAGE: Two schedulers S₀ and S₁ with shared cluster state

  • Cluster State:
    • Green circle in cell (3,1)
    • Red square in cell (2,6)
    • Blue circle in cell (3,3)
    • All other cells available
  • Scheduler Views:
    • S₀ (red) sees entire cluster
    • S₁ (blue) sees entire cluster
    • Both have current snapshot

Pending Jobs:

  • S₀: 2 red tasks to schedule
  • S₁: 2 blue tasks to schedule

Step 2: Parallel Scheduling

IMAGE: Schedulers making independent decisions

S₀ Decision:

  • Analyzes cluster state
  • Chooses cells for 2 red tasks
  • Cells: (3,4) and (2,4)
  • Cells: (3,5) and (2,5) also chosen
  • Prepares transaction

S₁ Decision:

  • Analyzes cluster state
  • Chooses cells for 2 blue tasks
  • Cells: (3,5) and (3,6)
  • Cells: (3,7) also chosen
  • Prepares transaction

Note: Both schedulers working simultaneously without coordination

Step 3: Transaction Commit

Conflict Detected:

  • S₀ wants to use cell (3,5)
  • S₁ also wants to use cell (3,5)
  • Overlapping resource claim
  • ⚠️ Conflict!

Omega's Action:

  • Both try to commit atomically
  • Only one transaction succeeds
  • Conflict detected via optimistic concurrency control

State = status of tasks being run in the cluster

Step 4: Resolution

IMAGE: Success and failure outcomes

S1S_1 Transaction: ✓ Success

  • Commit accepted
  • Blue tasks placed at (3,5), (3,6), (3,7)
  • Cluster state updated

S0S_0 Transaction: ✗ Failure

  • Commit rejected (conflict)
  • Red tasks NOT placed
  • S₀ must retry

Next Steps:

  1. S0S_0 resyncs its local copy from shared state
  2. Sees S1S_1's new allocations
  3. Reschedules with updated information
  4. Chooses different cells (avoiding (3,5))
  5. Tries commit again

Omega Shared-State Scheduling - Steps

Formal Process:

  1. Framework scheduler makes a decision
    • Based on local copy of cluster state
    • Applies its scheduling algorithm
    • Selects resources for tasks
  2. Updates the shared cell state copy in an atomic commit
    • Transaction submitted to shared state
    • All-or-nothing semantics
    • Includes validation that resources are still available
  3. At most one such commit will succeed
    • First successful commit wins
    • Concurrent conflicting commits rejected
    • Optimistic concurrency control
  4. Rejected scheduler resyncs its local copy and reschedules
    • Fetches latest cluster state
    • Re-runs scheduling algorithm
    • Submits new transaction
    • Process repeats until success

Transaction Details:

Commit succeeds if: ∀r∈R, state⁡(r)=expected_state⁡(r)\boxed{\text{Commit succeeds if: } \forall r \in \mathcal{R},\ \operatorname{state}(r) = \operatorname{expected\_state}(r)}
  • If any resource has changed state since scheduler's snapshot, commit fails
  • Prevents double-booking
  • Ensures consistency

Omega offer resource hoarding: FALSE (because it never locked the resource?)

Analogy: Like buying concert tickets online - you see available seats, select some, but when you click "purchase," if someone else just bought one of your seats, your purchase fails and you must choose again with the updated availability.


Omega Scheduling - Key Characteristics

1. Parallel Scheduler Operation

  • Multiple schedulers operate in parallel
    • No waiting for other schedulers
    • Each makes independent decisions
    • True parallelism, not time-slicing

2. Incremental Transactions

  • Schedulers typically do incremental transactions
    • Schedule a few tasks at a time (not all)
    • Reduces conflict probability
    • Avoids starvation
    • Job-by-job process??????

Example:

  • Instead of scheduling 1000 tasks in one transaction
  • Schedule 10 tasks per transaction, 100 times
  • If conflict occurs, only lose 10 task placements, not 1000

3. Common Priority Scale

  • Schedulers agree on common scale for expressing relative importance of jobs
    • Production > Batch > Test
    • Numeric priorities (e.g., 0-1000)
    • Used in conflict resolution

Conflict Resolution:

  • If both schedulers want same resource
  • Higher priority job gets resource
  • Lower priority job must reschedule

4. Performance Viability

Critical Factor:

"Performance viability of the shared-state approach is determined by the frequency at which transactions fail and their costs!!!"

Key Metrics:

  • Conflict rate: % of transactions that conflict
  • Retry cost: Time to resync and reschedule
  • Convergence: How many retries until success

Research Finding:

"Our performance evaluation of the Omega model using both lightweight simulations with synthetic workloads, and high-fidelity, trace-based simulations of production workloads at Google, shows that optimistic concurrency over shared state is a viable, attractive approach to cluster scheduling."

Why It Works:

  1. Low conflict rates in practice
    • Most tasks have flexible placement
    • Cluster typically has many available resources
    • Conflicts rare when cluster not fully utilized
  2. Fast retries
    • Resync is fast (read from shared state)
    • Re-scheduling often quick (similar decision)
    • Exponential backoff prevents thrashing
  3. Better than alternatives
    • Faster than serialized scheduling (monolithic)
    • More flexible than static partitioning
    • More informed than two-level (Mesos)

Real-world Validation: Google found conflict rates of 1-5% in production, meaning 95-99% of scheduling decisions succeed on first try. The performance gain from parallelism far outweighs the cost of occasional retries.


Key Characteristics of Omega Scheduler

1. Shared-State Scheduling

Architecture:

  • Unlike Mesos' two-level scheduler, Omega shares a single global view of cluster resources across multiple schedulers
  • Each scheduler can directly see and modify resource assignments without waiting for a central authority to grant them
  • Full cluster visibility
    • See all machines and their state
    • See all running tasks
    • See all pending allocations
    • Make globally-informed decisions

Benefits:

  • Better decisions: Schedulers have complete information
  • No information hiding: Unlike Mesos resource offers
  • Support for preemption: Can see what to preempt
  • Gang scheduling: Can verify entire gang fits

Real-world Impact: With full visibility, a batch scheduler can intelligently avoid nodes running latency-sensitive services, something difficult with limited resource offers.

2. Optimistic Concurrency Control

Concept:

  • Instead of serializing scheduling decisions (like in Mesos and Kubernetes), Omega allows parallel scheduling by different frameworks
  • Uses conflict detection and resolution to handle concurrent scheduling requests
  • If multiple schedulers try to modify the same resources simultaneously, conflicts are detected and resolved by rolling back or retrying

Process:

  1. Read cluster state (no locks)
  2. Compute scheduling decision locally
  3. Attempt commit with version check
  4. If conflict, rollback and retry with new state
  5. If success, update is applied

Concurrency Control Formula:

Success=(Current Version)==(Expected Version)\boxed{\text{Success} = \text{(Current Version)} == \text{(Expected Version)}}

Comparison with Traditional Approaches:

ApproachConcurrencyConflictsThroughput
Pessimistic (Locks)LowNoneLow
Optimistic (Omega)HighSomeHigh

Why Optimistic Works:

  • Low contention in practice (big cluster, many resources)
  • Cost of conflict < Cost of serialization
  • Fast conflict recovery (just retry)

Analogy: Optimistic concurrency is like multiple chefs cooking in a large kitchen - they mostly work independently, but if two reach for the same ingredient simultaneously, one backs off and grabs something else. Much faster than making all chefs wait in line for a single "permission to cook" token.

3. Multi-Tenant & Multi-Framework Support

Capabilities:

  • Supports multiple applications including:
    • Batch jobs: Spark, Hadoop, MapReduce
    • Service-oriented workloads: Kubernetes, Borg (Google's internal system)
    • ML training: TensorFlow, PyTorch jobs
    • Infrastructure services: Monitoring, logging
  • Unlike Mesos where each framework has its own scheduler, Omega allows all schedulers to act concurrently on the same resource pool

Architecture Benefits:

  1. Independent evolution
    • Each team maintains their scheduler
    • Can deploy updates independently
    • No central bottleneck for changes
  2. Specialized optimization
    • Batch scheduler optimizes for throughput
    • Service scheduler optimizes for availability
    • ML scheduler optimizes for GPU placement
  3. Fair sharing
    • All schedulers see same resources
    • Allocation policies enforce fairness
    • No scheduler starves

Real-world Usage: At Google, different product teams (Search, Ads, Gmail, YouTube) each run their own Omega scheduler instance, all sharing the same cluster infrastructure but optimizing for their specific needs.


Omega vs. Mesos vs. Kubernetes

FeatureOmega SchedulerApache MesosKubernetes
Scheduling ModelShared-state, optimistic concurrencyTwo-level schedulingCentralized, declarative
ScalabilityHigh (parallel scheduling)High (multi-framework support)High but focused on containers
Workload SupportMulti-framework (batch, services, containers)Multi-framework (batch, services, containers)Primarily containerized workloads
Conflict HandlingOptimistic concurrency controlMesos master offers resources to frameworksLock-based resource allocation
Container SupportYes (supports multiple runtimes)Yes (Docker, Mesos containers)Primarily Docker & CRI-O
Global ViewFull cluster visibilityLimited (offer-based)Full cluster visibility
Decision SpeedVery fast (parallel)Moderate (serial offers)Moderate (centralized)
Preemption SupportYesLimitedYes
AdoptionGoogle internalTwitter, Airbnb, AppleIndustry standard

Detailed Comparison

Scheduling Philosophy

Omega:

  • "Give everyone full information, let them decide in parallel"
  • Trust but verify (optimistic)
  • Global optimization possible

Mesos:

  • "Central authority offers resources, frameworks choose"
  • Conservative (pessimistic)
  • Local optimization by framework

Kubernetes:

  • "Central scheduler with declarative desired state"
  • Reconciliation loops
  • Constraint satisfaction

When to Use Each

Use Omega Approach When:

  • Need maximum scheduling throughput
  • Have diverse workload types
  • Want team autonomy
  • Can handle occasional conflicts
  • (Note: Omega itself is Google-internal, not open source)

Use Mesos When:

  • Need proven multi-framework support
  • Want framework isolation
  • Have batch + long-running services
  • Don't need ultra-low latency scheduling

Use Kubernetes When:

  • Primarily containerized workloads
  • Want strong ecosystem and tooling
  • Need production-ready platform today
  • Container orchestration is main focus

Industry Trend: Most organizations have moved to Kubernetes due to its maturity, ecosystem, and cloud-native focus. Mesos still used in specific cases (e.g., existing deployments). Omega's ideas influence modern scheduler designs.


Understanding Deadlock vs. Starvation

#Recap

Context: Important concepts for understanding scheduling challenges

Deadlock

Definition:

  • All processes keep waiting for each other to complete
  • None get executed
  • Circular dependency situation

Characteristics:

  1. Process behavior:
    • All processes blocked indefinitely
    • No progress made by anyone
    • System stuck
  2. Resources:
    • Resources are blocked by the processes
    • Cannot be released without external intervention
  3. Necessary conditions (all must be true):
    • Mutual Exclusion: Resource held exclusively
    • Hold and Wait: Process holds resources while waiting for more
    • No Preemption: Resources cannot be forcibly taken
    • Circular Wait: Circular chain of waiting processes
  4. Also known as: Circular wait

Example:

Process A: Holds Resource 1, Needs Resource 2
Process B: Holds Resource 2, Needs Resource 1
→ Deadlock! Neither can proceed.

Starvation

Definition:

  • High priority processes keep executing
  • Low priority processes are blocked indefinitely
  • System makes progress, but unfairly

Characteristics:

  1. Process behavior:

    • High-priority processes execute
    • Low-priority processes never get chance
    • System functions but unfairly
  2. Resources:

    • Resources continuously utilized by high priority processes
    • Low-priority processes wait indefinitely
  3. Cause:

    • Priorities are assigned to the processes
    • Strict priority scheduling
    • No aging or fairness mechanism
  4. Also known as: Lived lock

Example:

Priority 10 jobs: Always running
Priority 5 jobs: Sometimes running
Priority 1 jobs: Never running (starved)

Comparison Table

AspectDeadlockStarvation
ProblemAll processes stuckLow-priority processes stuck
ResourcesLocked by processesContinuously used by high-priority
Necessary ConditionsMutual Exclusion, Hold and Wait, No preemption, Circular WaitPriorities assigned to processes
System StateNo progressSome progress (unfair)
Alternative NameCircular waitLived lock

Prevention in Omega/Mesos Context

Deadlock Prevention:

  • Omega: Optimistic concurrency prevents hold-and-wait

    • No locks held during scheduling computation
    • Transactions are all-or-nothing
    • Cannot create circular dependencies
  • Mesos: Resource offers prevent circular wait

    • Offers are revocable
    • Frameworks cannot hold indefinitely

Starvation Prevention:

  • Omega: Incremental transactions

    • Forces schedulers to make progress gradually
    • Cannot monopolize decision making
  • Mesos: Allocation policies

    • DRF ensures fairness across frameworks
    • Revocation of tasks from greedy frameworks

Real-world Example: Without starvation prevention, a production workload scheduler could monopolize all cluster resources, preventing development teams from testing. DRF and priorities prevent this by ensuring everyone gets their fair share based on importance.


Conclusion and Takeaways

Key Innovations

Mesos Contributions:

  1. Fine-grained sharing - task-level resource allocation
  2. Resource offers - decentralized decision making
  3. Framework abstraction - support for diverse applications
  4. Deployed at scale - proven in production (Twitter, Airbnb, Apple)

Omega Contributions:

  1. Shared-state scheduling - full cluster visibility
  2. Optimistic concurrency - parallel scheduling decisions
  3. Independent schedulers - per-workload optimization
  4. Validated at Google scale - handles largest deployments

Modern Impact

Industry Evolution:

  • 2011: Mesos introduces two-level scheduling
  • 2013: Omega shows shared-state viability
  • 2014+: Kubernetes emerges, combining best ideas
    • Centralized but extensible (like Omega's shared state)
    • Declarative (like Omega's eventual consistency)
    • Practical and production-ready (like Mesos)

Current Best Practices:

  • Most organizations use Kubernetes for container orchestration
  • Legacy Mesos deployments being migrated
  • Omega's ideas influence modern scheduler designs
  • Multi-scheduler patterns emerging (e.g., Volcano for batch on K8s)

Real-world Adoption:

  • Kubernetes: Adopted by majority of cloud-native organizations
  • Mesos: Still running at some large-scale deployments (legacy)
  • Borg/Omega: Google-internal systems, not open source
  • Cloud platforms: AWS ECS, Azure Container Instances use similar concepts

Study Material

Required Reading

  1. Mesos: A Platform for Fine-Grained Resource Sharing in Data Center

    • Authors: B. Hindman, A. Konwinski, M. Zaharia, A. Ghodsi, A. D. Joseph, R. Katz, S. Shenker, I. Stoica
    • From: University of California, Berkeley
    • Conference: NSDI 2011
    • Sections: 1-3.5, 6-6.1.2
  2. Omega: Flexible, Scalable Schedulers for Large Compute Clusters

    • Authors: Malte Schwarzkopf, Andy Konwinski, Michael Abd-El-Malek, John Wilkes
    • Conference: EuroSys 2013
    • Sections: 1-3, 8
  3. Video Resource:


Key Concepts to Master

Mesos:

  • Resource offers mechanism
  • Two-level scheduling architecture
  • DRF (Dominant Resource Fairness)
  • Framework abstraction
  • Fine-grained vs. coarse-grained sharing

Omega:

  • Shared-state architecture
  • Optimistic concurrency control
  • Conflict resolution
  • Parallel scheduling
  • Cell state management

Comparison:

  • Monolithic vs. Two-level vs. Shared-state
  • Pessimistic vs. Optimistic concurrency
  • Scalability tradeoffs
  • Performance implications

Practice Questions

  1. Why does static partitioning lead to poor utilization?

    • Consider workload variability
    • Think about resource fragmentation
    • Analyze opportunity cost
  2. How does Mesos' resource offer mechanism work?

    • Describe the flow from slave to master to framework
    • Explain filter mechanism
    • Discuss rejection handling
  3. What are the advantages and disadvantages of resource offers vs. shared state?

    • Compare information visibility
    • Analyze decision speed
    • Consider conflict rates
  4. How does Omega handle scheduling conflicts?

    • Describe transaction model
    • Explain optimistic concurrency
    • Discuss retry mechanisms
  5. Why is optimistic concurrency viable for cluster scheduling?

    • Consider cluster size vs. scheduling rate
    • Think about conflict probability
    • Analyze cost of conflicts vs. serialization

Real-World Applications to Research

Study these production systems:

  1. Kubernetes Scheduler

    • How does it compare to Mesos/Omega?
    • Scheduling framework and plugins
    • Multi-scheduler support
  2. YARN (Hadoop 2.0)

    • Resource manager architecture
    • How it relates to Mesos concepts
    • ApplicationMaster model
  3. Borg (Google)

    • Predecessor to Omega
    • Lessons learned
    • Relationship to Kubernetes
  4. Modern Cloud Schedulers

    • AWS ECS scheduler
    • Azure Container Instances
    • Google Cloud Run

Additional Notes

Production Deployment Considerations

When implementing cluster scheduling, consider:

  1. Workload Characterization

    • Batch vs. service workloads
    • Resource profiles (CPU-bound, memory-bound, I/O-bound)
    • Priority and SLA requirements
  2. Scalability Requirements

    • Cluster size (10s, 100s, 1000s of nodes?)
    • Job arrival rate
    • Scheduling latency requirements
  3. Operational Complexity

    • Team expertise
    • Maintenance burden
    • Monitoring and debugging
  4. Ecosystem and Tooling

    • Available integrations
    • Community support
    • Documentation quality

Recommendation: For new projects today, start with Kubernetes unless you have specific requirements that mandate otherwise. It has the best ecosystem, tooling, and community support.


Further Reading

Academic Papers:

  • "Large-scale cluster management at Google with Borg" (EuroSys 2015)
  • "Borg, Omega, and Kubernetes" (ACM Queue 2016)
  • "Dominant Resource Fairness: Fair Allocation of Multiple Resource Types" (NSDI 2011)

Industry Resources:

  • Kubernetes documentation on scheduling
  • Mesos architecture documentation
  • Cloud provider scheduling whitepapers

Related Technologies:

  • Apache Spark on Kubernetes
  • Volcano Scheduler (batch scheduling on K8s)
  • Kueue (job queueing for Kubernetes)

Summary

Key Takeaways

✅ Mesos pioneered fine-grained sharing with two-level scheduling ✅ Omega showed shared-state parallel scheduling works at scale ✅ Modern systems (Kubernetes) combine best ideas from both ✅ Scheduling is hard - must balance utilization, fairness, latency ✅ No perfect solution - each approach has tradeoffs

The Future

Emerging Trends:

  • AI/ML scheduling - using ML to predict resource needs
  • Serverless integration - function-as-a-service scheduling
  • Edge computing - scheduling across geo-distributed resources
  • Multi-cluster - federated scheduling across clusters
  • Cost optimization - cloud cost-aware scheduling

Final Thought: Understanding Mesos and Omega gives you the foundation to understand any modern cluster scheduler. The core concepts - resource abstraction, fair sharing, conflict resolution - appear everywhere from cloud platforms to edge computing.


End of Notes