Pidoku

One Scheduling Step

Intermediate 1h Difficulty 3/5 Lesson 01 of 04

Prerequisites The Engine Core Loop

The idea in one minute#

Scheduler.schedule() is called once per engine step and answers one question: which tokens does the GPU process next? It has a budget of tokens and a pool of memory blocks. It walks the requests that are already running and gives each what it needs; then, if budget and memory remain, it admits requests from the waiting queue. For each request it decides a single number — how many tokens to compute this step — and reserves the KV cache blocks those tokens will occupy. The result is a SchedulerOutput. After the GPU has run, update_from_output() records what was produced. Those two functions are the scheduler; this lesson reads both.

A picture#

flowchart TB
  S[":i-list-checks: <b>schedule()</b><br/><small>token_budget = max_num_batched_tokens</small>"] --> P1["<b>Pass 1 — each RUNNING request</b><br/><small>tokens = min(have − computed, budget)</small>"]
  P1 --> A1{"allocate_slots()"}
  A1 -- "pool empty" --> PRE[":i-triangle-alert: preempt a running request<br/><small>then try again</small>"]
  PRE --> A1
  A1 -- "blocks granted" --> P2["<b>Pass 2 — head of the WAITING queue</b><br/><small>gates: budget left, running &lt; max_num_seqs,<br/>request ready, LoRA slot</small>"]
  P2 --> C["prefix-cache lookup<br/><small>skip tokens already cached</small>"]
  C --> A2{"allocate_slots()<br/><small>whole prompt must fit</small>"}
  A2 -- "granted" --> OK2["move to RUNNING<br/><small>then look at the next in the queue</small>"]
  A2 -- "refused" --> STOP["stop admitting this step"]
  OK2 --> OUT[":i-package: <b>SchedulerOutput</b>"]
  STOP --> OUT
  class S,P1,P2 queue
  class C,OK2 neutral
  class A1,A2 memory
  class PRE,STOP warn
  class OUT io

How it really works#

What the scheduler holds#

Python
self.requests: dict[str, Request]   # every unfinished request, by ID
self.waiting: RequestQueue          # not yet running (a deque, or a heap for priority)
self.running: list[Request]         # admitted; holds KV blocks
self.kv_cache_manager               # the block pool and prefix cache

And each Request carries the counters the algorithm is built on:

FieldMeaning
num_prompt_tokensLength of the prompt
num_tokensPrompt plus output tokens so far
num_tokens_with_specnum_tokens plus any speculative draft tokens
num_computed_tokensHow many tokens have had their KV computed — or are scheduled to
max_tokensThe output limit
statusWAITING, RUNNING, PREEMPTED, one of the FINISHED_* states, and a few “waiting for something” states

The scheduler’s whole job is to make num_computed_tokens catch up with num_tokens_with_spec for as many requests as the budget allows.

The budgets#

Python
token_budget = self.max_num_scheduled_tokens
input_budget = self.scheduler_config.max_num_batched_tokens

max_num_batched_tokens is the flag you set. max_num_scheduled_tokens defaults to the same value; it is smaller only when the model itself may add tokens to the batch, as speculative decoding does. For this lesson treat them as one number: the most tokens the GPU processes in one forward pass.

There is a second, separate limit: max_num_seqs, the most requests that may be in running at once. A step is therefore bounded both ways — at most max_num_seqs sequences and at most max_num_batched_tokens tokens.

Pass 1: running requests#

Python
req_index = 0
while req_index < len(self.running) and token_budget > 0:
    request = self.running[req_index]
    ...
    num_new_tokens = (
        request.num_tokens_with_spec
        + request.num_output_placeholders
        - request.num_computed_tokens
    )
    if 0 < long_prefill_token_threshold < num_new_tokens:
        num_new_tokens = long_prefill_token_threshold
    num_new_tokens = min(num_new_tokens, token_budget, input_budget - draft_slots)
    # Make sure the input position does not exceed the max model len.
    num_new_tokens = min(
        num_new_tokens,
        self.max_model_len - request.num_computed_tokens - self.num_sampled_tokens_per_step,
    )

Read it as: the gap, then three caps.

  • A request in the middle of its prompt has a large gap. It gets as much as the budget allows. This is one chunk of a chunked prefill.
  • A request that is generating has a gap of 1: the token it sampled last step has not been through the model yet.

(num_output_placeholders is zero unless async scheduling is on; it gets its own lesson.)

Running requests go first, in the order they were admitted. This is what makes the engine fair to requests already in progress: a newly arrived 50,000-token prompt cannot stall anyone’s generation, because all the running requests have taken their share before it is considered.

If nothing can be scheduled for a request this step, the code does something deliberate:

Python
if num_new_tokens == 0:
    # NOTE(woosuk): Here, by doing `continue` instead of `break`,
    # we do not strictly follow the FCFS scheduling policy and
    # allow the lower-priority requests to be scheduled.
    req_index += 1
    continue

It skips the request instead of stopping. One request that cannot move must not hold up the rest of the batch.

Then memory:

Python
while True:
    new_blocks = self.kv_cache_manager.allocate_slots(
        request, num_new_tokens, num_lookahead_tokens=self.num_lookahead_tokens
    )
    if new_blocks is not None:
        break                      # the request can be scheduled
    ...                            # otherwise: preempt someone and try again

allocate_slots returns the blocks needed to hold num_new_tokens more tokens, or None if the pool cannot supply them. A generating request needs a new block only once every 16 tokens, so most calls return an empty list at no cost. On None, a running request is preempted to free memory (Priorities, Preemption and Queues).

Finally the bookkeeping for a scheduled request:

Python
scheduled_running_reqs.append(request)
req_to_new_blocks[request_id] = new_blocks
num_scheduled_tokens[request_id] = num_new_tokens
token_budget -= num_new_tokens

Pass 2: waiting requests#

Pass 2 runs only if pass 1 preempted nobody. If memory was so short that a running request had to be evicted, admitting more work would be self-defeating:

Python
if not preempted_reqs and self._pause_state == PauseState.UNPAUSED:
    while token_budget > 0:
        request_queue = self.kv_holding_waiting or self.waiting
        if not request_queue or input_budget <= draft_slots:
            break
        if num_running >= self.max_num_active_reqs:
            break
        request = request_queue.peek_request()
        ...

For the request at the head of the queue, in order:

GatePasses whenOtherwise
Budgettoken_budget > 0Stop admitting this step
Slotslen(running) < max_num_seqsStop admitting this step
ReadyThe request is not waiting for a grammar to compile, KV data to arrive from elsewhere, or the next chunk of a streaming inputSkip it and look at the next
LoRAIts adapter is already in this batch, or the batch has fewer than max_loras distinct adaptersSkip it
Cache lookup—Find how much of the prompt is already in the prefix cache
Memoryallocate_slots(...) returns blocksStop admitting this step

Two different reactions appear in that table, and the difference matters. A request that is not ready is skipped: it is set aside and put back at the front of the queue afterwards, so the requests behind it are not delayed. A request that does not fit in memory stops admission: the scheduler does not look further down the queue for a smaller one. That is head-of-line blocking, chosen deliberately. Letting small requests overtake a large one would starve the large one for as long as small requests keep arriving.

The cache lookup happens exactly once per request, when it has computed nothing yet:

Python
if request.num_computed_tokens == 0:
    (new_computed_blocks, num_new_local_computed_tokens, ...) = \
        self._get_local_prefix_cache_hit(request)

If the first 2,048 tokens of a 2,100-token prompt are already in the cache, the request starts with num_computed_tokens = 2048 and only 52 tokens are charged to this step’s budget (Prefix Caching).

The amount to schedule uses the same rule as pass 1:

Python
num_new_tokens = request.num_tokens - num_computed_tokens
if 0 < long_prefill_token_threshold < num_new_tokens:
    num_new_tokens = long_prefill_token_threshold
...
num_new_tokens = min(num_new_tokens, request_token_budget)

If the prompt does not fit in what is left of the budget, it gets what is left and continues in later steps. That is chunked prefill, and it needs no code of its own: the request simply re-enters pass 1 next step with a smaller gap.

The memory check for a new request is stricter than for a running one:

Python
new_blocks = self.kv_cache_manager.allocate_slots(
    request, num_tokens_past_hit,
    num_new_computed_tokens=num_new_local_computed_tokens,
    new_computed_blocks=new_computed_blocks,
    full_sequence_must_fit=self.scheduler_reserve_full_isl,
    ...
)

scheduler_reserve_full_isl defaults to True: the request is admitted only if blocks for its whole prompt are free now, not just for the first chunk. Without that check, several long prompts could each be admitted on the strength of a first chunk and then fight for memory halfway through, preempting one another. The config docstring calls that “over-admission and KV cache thrashing”.

On success the request changes state:

Python
request = request_queue.pop_request()
self.running.append(request)
request.status = RequestStatus.RUNNING
request.num_computed_tokens = num_computed_tokens

What the scheduler sends to the worker#

Python
scheduler_output = SchedulerOutput(
    scheduled_new_reqs=new_reqs_data,          # full data, sent once per request
    scheduled_cached_reqs=cached_reqs_data,    # only the changes since last step
    num_scheduled_tokens=num_scheduled_tokens, # req_id -> tokens this step
    total_num_scheduled_tokens=...,
    finished_req_ids=self.finished_req_ids,    # so the worker can forget them
    new_block_ids_to_zero=...,
    ...
)

The split between new and cached is a bandwidth optimisation. The first time a request is scheduled, the worker receives its prompt token IDs, sampling parameters and block IDs (NewRequestData). Every later step it receives only the difference — the newly allocated block IDs and the updated token count (CachedRequestData). A request with a 100,000-token prompt sends those tokens to the worker once, not once per generated token.

finished_req_ids rides along so the worker can drop its copy of a request’s state. The worker never decides on its own that a request is over.

After scheduling: advance the counter immediately#

Python
def _update_after_schedule(self, scheduler_output):
    # Advance the number of computed tokens for the request AFTER
    # the request is scheduled.
    for req_id, num_scheduled_token in num_scheduled_tokens.items():
        request = self.requests[req_id]
        request.num_computed_tokens += num_scheduled_token
        request.num_in_flight_tokens += num_scheduled_token
        request.is_prefill_chunk = request.num_computed_tokens < (
            request.num_tokens + request.num_output_placeholders
        )

num_computed_tokens is moved forward when the work is scheduled, not when it is done. This is what allows the next schedule() call to run before the GPU has finished: the scheduler’s view already shows this step’s tokens as computed, so it can plan the following chunk. num_in_flight_tokens records how much of that is still a promise.

update_from_output(): closing the step#

When the worker returns sampled token IDs, the scheduler loops over every request that was in the batch:

Python
for req_id, num_tokens_scheduled in num_scheduled_tokens.items():
    request = self.requests.get(req_id)
    request.num_in_flight_tokens -= num_tokens_scheduled
    if request is None or request.is_finished():
        continue        # aborted while the GPU was working on it
    generated_token_ids = sampled_token_ids[req_index] if sampled_token_ids else []
    ...
    if new_token_ids:
        new_token_ids, stopped = self._update_request_with_output(request, new_token_ids)
    ...
    if stopped:
        finished = self._handle_stopped_request(request)
        if finished:
            self._free_request(request)
    outputs[request.client_index].append(EngineCoreOutput(request_id=req_id, ...))

For each request:

  1. Skip it if it has gone. An abort may have arrived while the forward pass ran.
  2. Take its sampled tokens. A request still in the middle of its prompt has none: tokens were computed but nothing was sampled.
  3. Append and check for stop (check_stop in vllm/v1/core/sched/utils.py), in this order: end-of-sequence token; a stop token ID; length (num_tokens >= max_model_len or num_output_tokens >= max_tokens); then, once min_tokens is satisfied, repetition detection.
  4. If finished, free its blocks back to the pool and add its ID to finished_req_ids for the next SchedulerOutput.
  5. Emit an EngineCoreOutput, filed under the API server that owns the request.

The comment above this loop is a warning to contributors: “As len(num_scheduled_tokens) can be up to 1K or more, the below loop can be a performance bottleneck.” Everything in it is paid once per request per token.

The life of a request#

flowchart LR
  W["WAITING"] -->|"admitted in pass 2"| R["RUNNING"]
  W -->|"grammar compiling,<br/>KV arriving"| B["WAITING_FOR_…"]
  B -->|"ready"| W
  R -->|"blocks taken away"| P["PREEMPTED"]
  P -->|"re-admitted, recomputes"| R
  R -->|"EOS or stop token"| F1["FINISHED_STOPPED"]
  R -->|"max_tokens or context full"| F2["FINISHED_LENGTH_CAPPED"]
  R -->|"client disconnected"| F3["FINISHED_ABORTED"]
  W --> F3
  class W,B,P queue
  class R compute
  class F1,F2,F3 neutral

The statuses are an IntEnum, and one comparison defines “finished”:

Python
@staticmethod
def is_finished(status: "RequestStatus") -> bool:
    return status > RequestStatus.PREEMPTED

Code#

A scheduler with both limits and a block pool. It implements the two passes, the whole-prompt admission check and the “stop at the first request that does not fit” rule.

Go
package main

import (
	"fmt"
	"strings"
)

const (
	blockSize   = 16
	tokenBudget = 64 // max_num_batched_tokens
	maxSeqs     = 3  // max_num_seqs
	totalBlocks = 12 // the whole KV cache: 192 tokens
)

type request struct {
	id       string
	prompt   int
	maxNew   int
	output   int
	computed int
	blocks   int
}

func (r *request) total() int { return r.prompt + r.output }

func blocksFor(tokens int) int { return (tokens + blockSize - 1) / blockSize }

func main() {
	waiting := []*request{
		{id: "A", prompt: 40, maxNew: 4},
		{id: "B", prompt: 100, maxNew: 2},
		{id: "C", prompt: 90, maxNew: 2},
		{id: "D", prompt: 10, maxNew: 2},
	}
	var running []*request
	free := totalBlocks

	for step := 1; len(waiting)+len(running) > 0 && step < 20; step++ {
		budget := tokenBudget
		var log []string
		scheduled := map[*request]int{}
		note := ""

		// Pass 1: running requests, in admission order.
		for _, r := range running {
			if budget == 0 {
				break
			}
			n := min(r.total()-r.computed, budget)
			need := blocksFor(r.computed+n) - r.blocks
			if need > free {
				note = r.id + " cannot grow (a real scheduler would preempt here)"
				break
			}
			free -= need
			r.blocks += need
			scheduled[r] = n
			budget -= n
			log = append(log, fmt.Sprintf("%s:%d", r.id, n))
		}

		// Pass 2: admit from the head of the waiting queue.
		for len(waiting) > 0 && budget > 0 && note == "" {
			if len(running) >= maxSeqs {
				note = "max_num_seqs reached"
				break
			}
			r := waiting[0]
			// The whole prompt must fit, not just the first chunk.
			if blocksFor(r.prompt) > free {
				note = fmt.Sprintf("%s needs %d blocks, %d free: stop admitting",
					r.id, blocksFor(r.prompt), free)
				break
			}
			waiting = waiting[1:]
			running = append(running, r)
			n := min(r.total()-r.computed, budget)
			need := blocksFor(r.computed+n) - r.blocks
			free -= need
			r.blocks += need
			scheduled[r] = n
			budget -= n
			log = append(log, fmt.Sprintf("%s:%d", r.id, n))
		}

		// "Forward pass" and update_from_output.
		keep := running[:0]
		for _, r := range running {
			if n, ok := scheduled[r]; ok {
				r.computed += n
				if r.computed == r.total() {
					r.output++ // sampled a token
				}
			}
			if r.output == r.maxNew {
				free += r.blocks // finished: blocks go back to the pool
				log = append(log, r.id+" done")
				continue
			}
			keep = append(keep, r)
		}
		running = keep

		fmt.Printf("step %-2d [%-22s] used %2d/%d  free blocks %2d  %s\n",
			step, strings.Join(log, " "), tokenBudget-budget, tokenBudget, free, note)
	}
}

Things to find in the output: B’s 100-token prompt is read in three chunks while A keeps generating; C is refused although the token budget has room, because its whole prompt does not fit in the free blocks; and D, a tiny request, waits behind C even though it would fit. That last line is head-of-line blocking, exactly as the real scheduler behaves.

Remember this#

  • schedule() has two passes: running requests first, then the head of the waiting queue.
  • Each request gets min(gap, budget) tokens; the gap is num_tokens - num_computed_tokens.
  • Two independent limits bound a step: max_num_batched_tokens and max_num_seqs.
  • Running requests that cannot move are skipped; a waiting request that cannot fit stops admission.
  • A new request is admitted only if its whole prompt fits in free blocks.
  • The worker gets full data for a request once, and only differences afterwards.
  • num_computed_tokens advances at scheduling time; update_from_output() settles the result.

Try it#

  1. Set totalBlocks to 20 and rerun. Does C get in earlier? Does D?
  2. Change the admission rule to check only the first chunk (blocksFor(min(r.prompt, budget))). Find a step where a running request “cannot grow”. That is the thrashing scheduler_reserve_full_isl prevents.
  3. Reorder the waiting queue so D arrives before C. How much sooner does D finish? This is why arrival order matters under FCFS.

Check yourself#

  1. Why are running requests scheduled before waiting ones?
  2. What is the difference between skipping a waiting request and stopping admission, and when does each happen?
  3. Why does the worker receive a request’s prompt token IDs only once?

Sources#

Checked on 5 October 2026 against main at commit 0c16eee.

↑↓ navigate↵ openesc close

drag to pan · scroll to zoom