Core Concepts and Resource Model
Volcano extends the Kubernetes API with three primary Custom Resource Definitions (CRDs):
- Queue: A logical grouping mechanism for resource allocation. Queues manage resource quotas and priority, allowing administrators to partition cluster resources among different teams or projects.
- PodGroup: A collection of pods that belong to the same job. This is the fundamental unit for gang scheduling, ensuring that pods are treated as a set.
- VolcanoJob: An abstraction similar to the native Kubernetes Job but with enhanced features for batch processing. It manages the lifecycle of a group of pods (PodGroup) and handles task dependencies.
Terminology Disambiguation: CRs vs. Scheduler Internals
When analyzing the Volcano source code, it is crucial to distinguish between the CRD definitions and the scheduler's internal data structures:
- In the Controller layer, a
VolcanoJobCRD containsTaskSpecs, which define pod templates. - In the Scheduler layer, the internal
JobInfostruct acts as a wrapper for aPodGroup. TheTaskInfostruct inside the scheduler represents a single Pod.
Thereofre, when the scheduler logic refers to a "Job," it is processing a PodGroup. When it refers to a "Task," it is processing a standard Kubernetes Pod wrapper.
Scheduling Framework Architecture
Volcano implements a plugin-based scheduling framework. The scheduling process is divided into distinct phases, defined by two main concepts:
- Actions: These represent the stages of the scheduling lifecycle. Common actions include
enqueue,allocate, andbackfill. - Plugins: These implement specific algorithms within each action, such as priority sorting, resource predicates, or node scoring.
The relationship is straightforward: Actions define what needs to be done, while Plugins define how to do it.
Source Code Analysis
Action Registration
Actions are registered during the program initialization phase. Located in the scheduler's action factory, the init() function registers specific action implementations into a global registry.
// pkg/scheduler/actions/factory.go
func init() {
// Registering the core scheduling actions
framework.RegisterAction(reclaim.New())
framework.RegisterAction(allocate.New())
framework.RegisterAction(backfill.New())
framework.RegisterAction(preempt.New())
framework.RegisterAction(enqueue.New())
framework.RegisterAction(shuffle.New())
}
Each Action implements a standard interface, ensuring consistent execution flow:
// pkg/scheduler/framework/interface.go
type Action interface {
Name() string
Initialize()
Execute(ssn *Session)
UnInitialize()
}
Scheduler Initialization and Main Loop
The scheduler starts by initializing a cache and watching for configuration changes. The core scheduling logic is triggered periodically by the runOnce method.
// pkg/scheduler/scheduler.go
func (s *Scheduler) Run(stopCh <-chan struct{}) {
s.loadSchedulerConf()
go s.watchSchedulerConf(stopCh)
s.cache.Run(stopCh)
s.cache.WaitForCacheSync(stopCh)
// Execute the scheduling loop periodically
go wait.Until(s.runOnce, s.schedulePeriod, stopCh)
}
The Session Mechanism
The Session is a critical context object created for every scheduling cycle. It holds a snapshot of the cluster state (Jobs, Nodes, Queues) and registered Plugin functions. When a Session opens, plugins register their callback functions (e.g., job ordering functions, predicate functions) into the Session.
// Simplified flow of a scheduling cycle
func (s *Scheduler) runOnce() {
// Create a new session context for this cycle
ssn := framework.OpenSession(s.cache, s.plugins, s.configurations)
defer framework.CloseSession(ssn)
// Execute configured actions sequentially
for _, action := range s.actions {
action.Execute(ssn)
}
}
Action Analysis: Enqueue
The enqueue action is responsible for moving jobs from a pending state into a queue where they become candidates for resource allocation.
- It iterates over all jobs and organizes them into their respective Queues.
- It uses a priority queue to determine which Queue should be processed first.
- It updates the status of the PodGroup to
Inqueueif the job is eligible.
// Simplified logic for enqueue
func (a *Action) Execute(ssn *framework.Session) {
queues := util.NewPriorityQueue(ssn.QueueOrderFn)
jobsMap := make(map[api.QueueID]*util.PriorityQueue)
// Categorize pending jobs into queues
for _, job := range ssn.Jobs {
if job.IsPending() {
// Add job to the specific queue's job list
jobsMap[job.Queue].Push(job)
}
}
// Process queues by priority
for !queues.Empty() {
q := queues.Pop().(*api.QueueInfo)
jobs := jobsMap[q.UID]
if jobs != nil && !jobs.Empty() {
job := jobs.Pop().(*api.JobInfo)
// Mark job as Inqueue
job.PodGroup.Status.Phase = scheduling.PodGroupInqueue
}
}
}
Action Analysis: Allocate
The allocate action performs the actual binding of pods to nodes. This is the most complex action, involving resource prediction, predicate filtering, and priority scoring.
1. Workflow Overview:
- Select a Queue based on priority.
- Select a Job from that Queue based on priority.
- Select a Task (Pod) from that Job.
- Filter available nodes using predicate functions.
- Score the remaining nodes to find the best fit.
- Allocate resources and bind the Pod to the Node.
2. Predicate and Scoring Logic:
The action constructs a predicate functon to filter out nodes that do not meet the task's requirements (e.g., insufficient resources, node affinity mismatch).
// Predicate function logic
predicateFn := func(task *api.TaskInfo, node *api.NodeInfo) error {
// Check if node has enough idle resources
if !task.InitResreq.LessEqual(node.FutureIdle(), api.Zero) {
return fmt.Errorf("insufficient resources")
}
// Execute plugin-registered predicate functions
return ssn.PredicateFn(task, node)
}
Once nodes are filtered, the scheduler scores them to find the optimal placement. It distinguishes between Idle resources (currently free) and Future Idle resources (expected to be freed).
// Main allocation loop skeleton
for !queues.Empty() {
queue := queues.Pop().(*api.QueueInfo)
jobs := jobsMap[queue.UID]
job := jobs.Pop().(*api.JobInfo)
for _, task := range job.Tasks {
// Filter nodes
predicateNodes, _ := ph.PredicateNodes(task, allNodes, predicateFn)
var bestNode *api.NodeInfo
// Score nodes to select the best one
if len(predicateNodes) > 1 {
nodeScores := util.PrioritizeNodes(task, predicateNodes, ssn.NodeOrderFn)
bestNode = util.SelectBestNode(nodeScores)
} else if len(predicateNodes) == 1 {
bestNode = predicateNodes[0]
}
// Bind task to node
if bestNode != nil {
stmt.Allocate(task, bestNode)
}
}
}
Action Analysis: Backfill
The backfill action handles best-effort tasks (tasks without strict resource requests) or fills in small resource gaps on nodes. It ensures that resources are maximized by scheduling tasks that can fit into the remaining slack space of the cluster.