Contents
- 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
- 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:

- Low utilization
- Resources allocated but not always used
- No dynamic reallocation
- 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 มากขึ้น
- 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
- High Utilization
- Maximize resource usage across the cluster
- Minimize idle capacity
- Dynamic allocation based on demand
- Support Diverse Frameworks
- Each framework has different scheduling needs
- Different optimization goals (latency vs. throughput)
- Save computation cost
- Reduce time to completion
- 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
- 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:
- High availability
- Reduces unplanned downtime
- Increases reliability of services and applications
- 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:
-
Mesos offers resources to a framework
- Example:
<slave_id, cpu_cores, memory_gb> - "Framework A, slave1 has 4 CPUs and 4GB available"
- Example:
-
Framework scheduler evaluates offer
- Checks if resources match requirements
- Considers data locality, constraints
-
Framework responds
- Accept and launch tasks
- Reject and wait for better offer
-
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:
- Mesos Master
- Fine-grained sharing coordinator across frameworks
- Makes resource offers to frameworks
- Single point of coordination (with standby replicas)
- Standby Masters
- Hot standby for high availability
- Take over if primary master fails
- ZooKeeper-based leader election
- Mesos Slave (on each node)
- Reports available resources to master
- Launches and monitors executor processes
- Enforces resource isolation
- Framework Scheduler
- Runs application-specific scheduling logic
- Receives resource offers
- Decides which tasks to launch on which resources
- 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:
- Slave 1 reports availability
<s1, 4cpu, 4gb, ...>→ Mesos master- Slave has 4 CPU cores and 4GB RAM available
- Master sends offer to Framework 1
<s1, 4cpu, 4gb, ...>→ Framework 1 Scheduler- "You can use these resources"
- Framework 1 responds with task requests
<task1, s1, 2cpu, 1gb, ...><task2, s1, 1cpu, 2gb, ...>- Framework accepts partial offer, launches 2 tasks
- 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:
- 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
Where:
- = dominant resource share of framework
- = total resources of that type
- Strict Priorities
- Some frameworks get higher priority
- Production > Development > Testing
- Higher priority frameworks get offers first
- 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:
- 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
- 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
- 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 ( 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:
- Hadoop running mix of small and large jobs
- Based on Facebook workload trace
- Realistic job size distribution
- Hadoop instance running large batch jobs
- Long-running background processing
- Spark running machine learning jobs
- Iterative algorithms
- In-memory processing
- Torque running MPI jobs
- Tightly-coupled HPC applications
- The Terascale Open-source Resource and QUEue Manager (TORQUE)
- Hadoop running mix of small and large jobs
Expected Results:
- 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:
- Dynamic resource allocation visible
- Different frameworks getting resources at different times
- No static boundaries
- High overall utilization
- Cluster stays busy (50-100% utilization)
- Little idle capacity
- 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
- 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:
| Metric | Static Partitioning | Mesos |
|---|---|---|
| Avg Utilization | 22% | 60% |
| Job Completion Time | Baseline | 2.7× faster |
| Framework Starvation | Possible | None observed |
Key Findings:
- Better utilization without sacrificing performance
- Jobs complete faster due to elastic scaling
- 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...

- Diverse workloads
- Circle, square, hexagon, triangle symbols
- Different resource profiles
- Batch vs. serving vs. ML
- Increasing cluster sizes
- Small cluster → Large cluster
- Scaling to 10,000+ machines
- Horizontal scaling challenges
- 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:
- Resource underutilization
- Machines idle while waiting for decisions
- Wasted capacity
- Increased latency
- Jobs wait in queue
- User experience degraded
- 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 (16.70 lost every minute!
The Actual Problem: Increasing Complexity
IMAGE: Diagram highlighting scheduler complexity

Complexity Factors:
- Multiple workload types with different needs
- Red, blue, green, yellow circles (different job types)
- Each has unique constraints
- Complex placement constraints
- Affinity rules
- Anti-affinity rules
- Locality preferences
- Hardware requirements
- Multiple scheduling objectives
- Fairness
- Utilization
- Latency
- Throughput
- 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:
- Break up the cluster scheduler into independent schedulers
- One scheduler per workload type or team
- Parallel operation
- Independent evolution
- 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:
- 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
- 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:
- : 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:
- No centralized scheduler
- Each framework maintains its own scheduler
- Schedulers operate in parallel
- Shared state visibility
- Each scheduler has full view of cluster
- Can see all resources and allocations
- Makes informed decisions
- 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

Transaction: ✓ Success
- Commit accepted
- Blue tasks placed at (3,5), (3,6), (3,7)
- Cluster state updated
Transaction: ✗ Failure
- Commit rejected (conflict)
- Red tasks NOT placed
- S₀ must retry
Next Steps:
- resyncs its local copy from shared state
- Sees 's new allocations
- Reschedules with updated information
- Chooses different cells (avoiding (3,5))
- Tries commit again
Omega Shared-State Scheduling - Steps
Formal Process:
- Framework scheduler makes a decision
- Based on local copy of cluster state
- Applies its scheduling algorithm
- Selects resources for tasks
- 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
- At most one such commit will succeed
- First successful commit wins
- Concurrent conflicting commits rejected
- Optimistic concurrency control
- 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:
- 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:
- Low conflict rates in practice
- Most tasks have flexible placement
- Cluster typically has many available resources
- Conflicts rare when cluster not fully utilized
- Fast retries
- Resync is fast (read from shared state)
- Re-scheduling often quick (similar decision)
- Exponential backoff prevents thrashing
- 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:
- Read cluster state (no locks)
- Compute scheduling decision locally
- Attempt commit with version check
- If conflict, rollback and retry with new state
- If success, update is applied
Concurrency Control Formula:
Comparison with Traditional Approaches:
| Approach | Concurrency | Conflicts | Throughput |
|---|---|---|---|
| Pessimistic (Locks) | Low | None | Low |
| Optimistic (Omega) | High | Some | High |
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:
- Independent evolution
- Each team maintains their scheduler
- Can deploy updates independently
- No central bottleneck for changes
- Specialized optimization
- Batch scheduler optimizes for throughput
- Service scheduler optimizes for availability
- ML scheduler optimizes for GPU placement
- 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
| Feature | Omega Scheduler | Apache Mesos | Kubernetes |
|---|---|---|---|
| Scheduling Model | Shared-state, optimistic concurrency | Two-level scheduling | Centralized, declarative |
| Scalability | High (parallel scheduling) | High (multi-framework support) | High but focused on containers |
| Workload Support | Multi-framework (batch, services, containers) | Multi-framework (batch, services, containers) | Primarily containerized workloads |
| Conflict Handling | Optimistic concurrency control | Mesos master offers resources to frameworks | Lock-based resource allocation |
| Container Support | Yes (supports multiple runtimes) | Yes (Docker, Mesos containers) | Primarily Docker & CRI-O |
| Global View | Full cluster visibility | Limited (offer-based) | Full cluster visibility |
| Decision Speed | Very fast (parallel) | Moderate (serial offers) | Moderate (centralized) |
| Preemption Support | Yes | Limited | Yes |
| Adoption | Google internal | Twitter, Airbnb, Apple | Industry 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:
- Process behavior:
- All processes blocked indefinitely
- No progress made by anyone
- System stuck
- Resources:
- Resources are blocked by the processes
- Cannot be released without external intervention
- 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
- 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:
-
Process behavior:
- High-priority processes execute
- Low-priority processes never get chance
- System functions but unfairly
-
Resources:
- Resources continuously utilized by high priority processes
- Low-priority processes wait indefinitely
-
Cause:
- Priorities are assigned to the processes
- Strict priority scheduling
- No aging or fairness mechanism
-
Also known as: Lived lock
Example:
Priority 10 jobs: Always running
Priority 5 jobs: Sometimes running
Priority 1 jobs: Never running (starved)
Comparison Table
| Aspect | Deadlock | Starvation |
|---|---|---|
| Problem | All processes stuck | Low-priority processes stuck |
| Resources | Locked by processes | Continuously used by high-priority |
| Necessary Conditions | Mutual Exclusion, Hold and Wait, No preemption, Circular Wait | Priorities assigned to processes |
| System State | No progress | Some progress (unfair) |
| Alternative Name | Circular wait | Lived 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:
- Fine-grained sharing - task-level resource allocation
- Resource offers - decentralized decision making
- Framework abstraction - support for diverse applications
- Deployed at scale - proven in production (Twitter, Airbnb, Apple)
Omega Contributions:
- Shared-state scheduling - full cluster visibility
- Optimistic concurrency - parallel scheduling decisions
- Independent schedulers - per-workload optimization
- 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
-
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
-
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
-
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
-
Why does static partitioning lead to poor utilization?
- Consider workload variability
- Think about resource fragmentation
- Analyze opportunity cost
-
How does Mesos' resource offer mechanism work?
- Describe the flow from slave to master to framework
- Explain filter mechanism
- Discuss rejection handling
-
What are the advantages and disadvantages of resource offers vs. shared state?
- Compare information visibility
- Analyze decision speed
- Consider conflict rates
-
How does Omega handle scheduling conflicts?
- Describe transaction model
- Explain optimistic concurrency
- Discuss retry mechanisms
-
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:
-
Kubernetes Scheduler
- How does it compare to Mesos/Omega?
- Scheduling framework and plugins
- Multi-scheduler support
-
YARN (Hadoop 2.0)
- Resource manager architecture
- How it relates to Mesos concepts
- ApplicationMaster model
-
Borg (Google)
- Predecessor to Omega
- Lessons learned
- Relationship to Kubernetes
-
Modern Cloud Schedulers
- AWS ECS scheduler
- Azure Container Instances
- Google Cloud Run
Additional Notes
Production Deployment Considerations
When implementing cluster scheduling, consider:
-
Workload Characterization
- Batch vs. service workloads
- Resource profiles (CPU-bound, memory-bound, I/O-bound)
- Priority and SLA requirements
-
Scalability Requirements
- Cluster size (10s, 100s, 1000s of nodes?)
- Job arrival rate
- Scheduling latency requirements
-
Operational Complexity
- Team expertise
- Maintenance burden
- Monitoring and debugging
-
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