Deep Dive into Volcano Scheduler: Core Concepts and Source Code Analysis

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 VolcanoJob CRD contains TaskSpecs, which define pod templates.
  • In the Scheduler layer, the internal JobInfo struct acts as a wrapper for a PodGroup. The TaskInfo struct 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:

  1. Actions: These represent the stages of the scheduling lifecycle. Common actions include enqueue, allocate, and backfill.
  2. 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.

  1. It iterates over all jobs and organizes them into their respective Queues.
  2. It uses a priority queue to determine which Queue should be processed first.
  3. It updates the status of the PodGroup to Inqueue if 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:

  1. Select a Queue based on priority.
  2. Select a Job from that Queue based on priority.
  3. Select a Task (Pod) from that Job.
  4. Filter available nodes using predicate functions.
  5. Score the remaining nodes to find the best fit.
  6. 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.

Tags: Volcano kubernetes batch-scheduling Go source-code-analysis

Posted on Thu, 01 Oct 2026 16:54:00 +0000 by Whear