Skip to main content
Create your own

Scheduling Policies: Impact on TTFT and ITL

Introduction

In our last lesson, you benchmarked our continuous batching scheduler against a static approach, proving with data why modern LLM serving engines have embraced dynamic, iteration-level scheduling. You saw a 2-4x throughput gain by simply keeping the GPU busy with useful work.

However, our scheduler still has a very simple logic for admitting new requests: First-Come, First-Served (FCFS). While fair, FCFS can lead to poor performance in real-world scenarios. Imagine being stuck in a checkout line with one item behind someone with three overflowing carts. This is precisely what happens when a short, quick query gets queued behind a request to summarize a massive document.

Today, we will address this "head-of-line blocking" problem. Your goal is to implement and compare different scheduling policies (e.g., FCFS, preemption) and analyze their effect on Time to First Token (TTFT) and Inter-Token Latency (ITL). By the end of this lesson, you will have evolved your scheduler to make smarter decisions that directly improve user-perceived latency.

1. Beyond Throughput: Defining User-Centric Latency Metrics

While throughput (tokens/sec) is a great measure of system efficiency, it doesn't tell the whole story from a user's perspective. For interactive applications, we care about responsiveness. Two key metrics capture this:

  • Time to First Token (TTFT): The time a user waits from sending a request to seeing the very first token of the response. Low TTFT is crucial for making an application feel responsive.
  • Inter-Token Latency (ITL): The time between consecutive tokens in a streaming response. Low and stable ITL makes the text appear to flow smoothly.

To build a strong mental model of these metrics, please study the following resources.

Inside vLLM: Anatomy of a High-Throughput LLM Inference System

The article 'Inside vLLM' provides an excellent, concise breakdown of the key performance metrics for LLM inference. We'll focus on the section that defines them.

Read the section 'Benchmarks and auto-tuning - latency vs throughput'. Pay close attention to the table defining TTFT, ITL, and other metrics. This will be our vocabulary for evaluating the policies we implement.

LLM Inference Service Token Generation Process: TTFT and ITL
This diagram from NVIDIA clearly visualizes the user's experience. TTFT is the initial wait, while ITL is the delay between each subsequent piece of the response.

With our current FCFS scheduler, a request with a short prompt that could have a low TTFT might get stuck waiting for a long prefill phase of a request that arrived just moments earlier, unnecessarily inflating its TTFT.

2. Policy-Based Scheduling: From FCFS to Shortest Prompt First

To solve this, we can introduce scheduling policies that decide the order in which waiting requests are considered for prefill. The vllm article Inside vLLM you just read mentions a policy setting that can be either FCFS or priority-based. Let's explore this.

Continuous Batching: Optimizing LLM Inference Throughput

The article 'Continuous Batching' discusses the scheduling decisions made at each iteration, including which waiting requests to admit.

Please read the subsection 'Which waiting requests to admit' within the 'Iteration-Level Scheduling' section. It outlines several common policies, including FIFO (same as FCFS) and Shortest Job First.

"Shortest Job First" is a classic scheduling algorithm that minimizes average wait time. In the context of LLMs, we don't know the output length in advance, but we can use the prompt length as a proxy. A Shortest Prompt First (SPF) policy will prioritize requests with shorter prompts for the compute-intensive prefill phase, aiming to get more requests into the decoding phase faster and thus improving the average TTFT.

3. The Memory Wall: Introducing Preemption

Implementing SPF seems like a clear win. But what happens when the system is under heavy load?

Consider this scenario:

  • The GPU's KV cache is almost full with several long-running decode requests.
  • A new, high-priority request with a very short prompt arrives.
  • The scheduler, following the SPF policy, wants to prefill this new request immediately.
  • However, the kv_cache_manager reports there isn't enough memory to even start the prefill.

The system is stuck. The new request must wait, defeating the purpose of the SPF policy. The solution is preemption.

The scheduler can decide to "evict" one or more running requests to free up memory for a higher-priority one. In "recompute preemption," the evicted request is moved back to the waiting queue. Its KV cache is freed, and its prompt will need to be re-processed later. This is a trade-off: we sacrifice the work already done on one request to unblock another.

Inside vLLM: Anatomy of a High-Throughput LLM Inference System

The vLLM article also touches on this exact mechanism, describing how the scheduler handles resource contention.

Re-read the 'Scheduler' section, focusing on step 2 under the allocate_slots function description. It explicitly mentions that if there aren't enough blocks, the engine may attempt 'recompute preemption by evicting low-priority requests'. This is exactly what we are going to implement.

The preemption policy itself—which request to evict—is another important decision. A simple and effective policy is to preempt the request that has made the least progress (i.e., has generated the fewest output tokens), as this minimizes the amount of "sunk" computational cost.

4. Implementation: Evolving the Scheduler

Now, let's put these concepts into code. You will modify the scheduler and simulation from the previous lesson to:

  1. Support pluggable scheduling policies (FCFS and SPF).
  2. Implement a simple preemption mechanism.
  3. Track and report TTFT and average ITL.

Below is the code from our last session, with added data structures for tracking our new metrics.

Your Tasks:

  1. Modify StallFreeScheduler.schedule_next_batch:
    • Add a policy attribute to the scheduler ('fcfs' or 'spf').
    • Before the loop that considers waiting_queue requests, sort a copy of the queue according to the policy. For 'spf', sort by ascending prompt_tokens length.
  2. Implement Preemption:
    • In the schedule_next_batch loop that tries to schedule prefills from the waiting queue, if self.kv_cache_manager.can_allocate returns False, trigger the preemption logic.
    • Your preemption logic should find a suitable victim from the active_pool (e.g., a decoding request that has generated the fewest tokens).
    • If a victim is found, free its resources using kv_cache_manager.free(), remove it from active_pool, and add it back to the waiting_queue. Then, re-attempt to schedule the original prefill request.
  3. Update simulate function:
    • Modify the main simulation loop to record the start_iteration and first_token_iteration for each request to calculate TTFT.
    • Also, track the total time spent on decode steps to calculate the average ITL.
import numpy as np
from collections import deque
from dataclasses import dataclass, field
from typing import List, Deque, Dict, Optional
import math




# --- Data Structures (Now with latency tracking) ---
@dataclass
class Request:
    id: int
    prompt_tokens: List[int]
    output_length: int
    arrival_iteration: int # New: track when the request arrived
    



    # Simulation state
    tokens_processed: int = 0
    output_tokens_generated: int = 0
    start_iteration: Optional[int] = None # When it first got scheduled
    first_token_iteration: Optional[int] = None # When first decode token was generated
    end_iteration: Optional[int] = None # When it completed

    @property
    def is_prefill_done(self) -> bool:
        return self.tokens_processed >= len(self.prompt_tokens)

    @property
    def is_complete(self) -> bool:
        return self.output_tokens_generated >= self.output_length

@dataclass
class ScheduledSequence:
    request_id: int
    is_decode: bool
    num_tokens: int




# --- Mock KVCacheManager (Unchanged) ---
class MockKVCacheManager:
    def __init__(self, num_blocks: int, block_size: int):
        self.num_total_blocks = num_blocks
        self.block_size = block_size
        self.free_blocks = num_blocks
        self.request_blocks: Dict[int, int] = {}

    def _get_num_required_blocks(self, num_tokens: int) -> int:
        return math.ceil(num_tokens / self.block_size)

    def can_allocate(self, request: Request, num_tokens_to_alloc: int) -> bool:
        current_blocks = self.request_blocks.get(request.id, 0)
        current_allocated_tokens = current_blocks * self.block_size
        
        total_tokens_so_far = request.tokens_processed + request.output_tokens_generated
        



        # This logic is simplified for simulation. In reality, it's more complex.
        tokens_in_last_block_unused = (current_allocated_tokens - total_tokens_so_far) % self.block_size
        effective_tokens_to_alloc = max(0, num_tokens_to_alloc - tokens_in_last_block_unused)
        
        required_new_blocks = self._get_num_required_blocks(effective_tokens_to_alloc)
        return self.free_blocks >= required_new_blocks

    def allocate_slots(self, request: Request, num_tokens_to_alloc: int) -> bool:
        if not self.can_allocate(request, num_tokens_to_alloc):
            return False
            
        current_blocks = self.request_blocks.get(request.id, 0)
        current_allocated_tokens = current_blocks * self.block_size
        total_tokens_so_far = request.tokens_processed + request.output_tokens_generated
        tokens_in_last_block_unused = (current_allocated_tokens - total_tokens_so_far) % self.block_size
        effective_tokens_to_alloc = max(0, num_tokens_to_alloc - tokens_in_last_block_unused)
        
        required_new_blocks = self._get_num_required_blocks(effective_tokens_to_alloc)
        
        self.free_blocks -= required_new_blocks
        self.request_blocks[request.id] = self.request_blocks.get(request.id, 0) + required_new_blocks
        return True

    def free(self, request: Request):
        if request.id in self.request_blocks:
            num_freed_blocks = self.request_blocks.pop(request.id)
            self.free_blocks += num_freed_blocks



            # Reset request state upon preemption
            request.tokens_processed = 0
            request.start_iteration = None




# --- Scheduler with Policy and Preemption ---
class StallFreeScheduler:
    def __init__(self, token_budget: int, max_batch_size: int, kv_cache_manager: MockKVCacheManager, policy: str = 'fcfs', allow_preemption: bool = False):
        self.token_budget = token_budget
        self.max_batch_size = max_batch_size
        self.kv_cache_manager = kv_cache_manager
        self.policy = policy
        self.allow_preemption = allow_preemption
        
        self.waiting_queue: Deque[Request] = deque()
        self.active_pool: Dict[int, Request] = {}

    def add_request(self, request: Request):
        self.waiting_queue.append(request)
        
    def _preempt(self) -> bool:
        """Preempts a running request to free memory. Returns True if successful."""
        if not self.allow_preemption:
            return False




        # Victim selection policy: preempt the decoding request that has generated the fewest tokens.
        decode_requests = [r for r in self.active_pool.values() if r.is_prefill_done() and not r.is_complete()]
        if not decode_requests:
            return False # No one to preempt

        victim = min(decode_requests, key=lambda r: r.output_tokens_generated)
        
        print(f"    (Preempting request {victim.id} to free memory)")
        self.kv_cache_manager.free(victim)
        del self.active_pool[victim.id]
        self.waiting_queue.appendleft(victim) # Put it at the front of the queue
        return True

    def schedule_next_batch(self) -> List[ScheduledSequence]:
        next_batch: List[ScheduledSequence] = []
        current_tokens = 0




        # 1. Process completions
        completed_ids = [req_id for req_id, req in self.active_pool.items() if req.is_complete]
        for req_id in completed_ids:



            # We don't free here because it's freed on completion event in main loop
            del self.active_pool[req_id]




        # 2. Schedule decode requests (highest priority)
        # Sort active pool to have a deterministic order for scheduling
        sorted_active_pool = sorted(self.active_pool.values(), key=lambda r: r.arrival_iteration)
        
        for req in sorted_active_pool:
            if req.is_prefill_done and not req.is_complete:
                if current_tokens + 1 <= self.token_budget and len(next_batch) < self.max_batch_size and self.kv_cache_manager.can_allocate(req, 1):
                    self.kv_cache_manager.allocate_slots(req, 1)
                    next_batch.append(ScheduledSequence(request_id=req.id, is_decode=True, num_tokens=1))
                    current_tokens += 1
        



        # 3. Schedule prefill requests from waiting queue
        # TODO: Apply scheduling policy here
        if self.policy == 'spf':



            # Sort a *copy* of the queue by prompt length
            candidate_queue = sorted(self.waiting_queue, key=lambda r: len(r.prompt_tokens))
        else: # 'fcfs'
            candidate_queue = list(self.waiting_queue)

        for req in candidate_queue:
            if len(self.active_pool) >= self.max_batch_size:
                break
            



            # Simple logic: only one prefill per batch for now
            if any(not s.is_decode for s in next_batch):
                break

            prompt_len = len(req.prompt_tokens)
            chunk_size = min(prompt_len, self.token_budget - current_tokens)
            
            if chunk_size <= 0:
                continue

            can_fit = self.kv_cache_manager.can_allocate(req, chunk_size)
            



            # TODO: Implement preemption logic
            if not can_fit:
                if self._preempt():
                    can_fit = self.kv_cache_manager.can_allocate(req, chunk_size) # Retry
                if not can_fit:
                    continue # Still can't fit, skip for this round




            # If we can fit, schedule it
            self.waiting_queue.remove(req)
            self.active_pool[req.id] = req
            
            self.kv_cache_manager.allocate_slots(req, chunk_size)
            next_batch.append(ScheduledSequence(request_id=req.id, is_decode=False, num_tokens=chunk_size))
            current_tokens += chunk_size
            req.tokens_processed += chunk_size
            



            # Mark start iteration
            if req.start_iteration is None:
                req.start_iteration = SIM_ITERATION # Use global simulation iteration

        return next_batch




# --- Simulation ---
SIM_ITERATION = 0

def generate_workload(num_requests: int, seed: int = 42) -> List[Request]:
    np.random.seed(seed)
    prompt_lengths = np.random.lognormal(mean=3.5, sigma=1.0, size=num_requests).astype(int) + 10
    output_lengths = np.random.lognormal(mean=3.0, sigma=0.8, size=num_requests).astype(int) + 10
    arrival_times = np.random.randint(0, 50, size=num_requests) # Staggered arrivals
    
    prompt_lengths = np.clip(prompt_lengths, 10, 512)
    output_lengths = np.clip(output_lengths, 10, 512)

    requests = []
    for i in range(num_requests):
        requests.append(Request(id=i, prompt_tokens=[1] * prompt_lengths[i], output_length=output_lengths[i], arrival_iteration=arrival_times[i]))
    
    return sorted(requests, key=lambda r: r.arrival_iteration)

def simulate(workload: List[Request], scheduler: StallFreeScheduler) -> List[Request]:
    global SIM_ITERATION
    SIM_ITERATION = 0
    
    all_requests = {req.id: req for req in workload}
    pending_requests = deque(workload)
    completed_requests = []

    while pending_requests or scheduler.active_pool:



        # Add new requests that have arrived
        while pending_requests and pending_requests[0].arrival_iteration <= SIM_ITERATION:
            req = pending_requests.popleft()
            scheduler.add_request(req)

        batch = scheduler.schedule_next_batch()




        # Update request states based on batch
        for seq in batch:
            req = scheduler.active_pool[seq.request_id]
            if not seq.is_decode:
                pass # Prefill tokens_processed was already updated
            else:
                if req.first_token_iteration is None:
                    req.first_token_iteration = SIM_ITERATION
                req.output_tokens_generated += 1
        



        # Check for completions
        active_ids = list(scheduler.active_pool.keys())
        for req_id in active_ids:
            req = scheduler.active_pool.get(req_id)
            if req and req.is_complete:
                req.end_iteration = SIM_ITERATION
                scheduler.kv_cache_manager.free(req)



                # del scheduler.active_pool[req.id] # Let scheduler handle this next iter
                completed_requests.append(req)

        SIM_ITERATION += 1
        if SIM_ITERATION > 10000: # Safety break
            print("Simulation timed out!")
            break
            
    return completed_requests


def analyze_results(completed_requests: List[Request]):
    if not completed_requests:
        print("No requests completed.")
        return

    ttfts = []
    itls = []
    
    for req in completed_requests:
        if req.start_iteration is not None and req.first_token_iteration is not None:
            ttfts.append(req.first_token_iteration - req.arrival_iteration)
            
            if req.output_length > 1:
                total_decode_time = req.end_iteration - req.first_token_iteration
                avg_itl = total_decode_time / (req.output_length - 1)
                itls.append(avg_itl)

    print(f"  Avg TTFT: {np.mean(ttfts):.2f} iterations")
    print(f"  P95 TTFT: {np.percentile(ttfts, 95):.2f} iterations")
    print(f"  Avg ITL: {np.mean(itls):.2f} iterations/token")


if __name__ == '__main__':



    # --- Benchmark Parameters ---
    KV_CACHE_BLOCKS = 4096  # Reduced to create memory pressure
    



    # --- Generate Workload ---
    workload = generate_workload(100)
    



    # --- FCFS ---
    print("\n--- Running FCFS Simulation ---")
    kv_manager_fcfs = MockKVCacheManager(num_blocks=KV_CACHE_BLOCKS, block_size=16)
    scheduler_fcfs = StallFreeScheduler(token_budget=2048, max_batch_size=16, kv_cache_manager=kv_manager_fcfs, policy='fcfs')
    results_fcfs = simulate(list(workload), scheduler_fcfs)
    analyze_results(results_fcfs)




    # --- SPF ---
    print("\n--- Running SPF Simulation ---")
    kv_manager_spf = MockKVCacheManager(num_blocks=KV_CACHE_BLOCKS, block_size=16)
    scheduler_spf = StallFreeScheduler(token_budget=2048, max_batch_size=16, kv_cache_manager=kv_manager_spf, policy='spf')
    results_spf = simulate(list(workload), scheduler_spf)
    analyze_results(results_spf)
    



    # --- SPF with Preemption ---
    print("\n--- Running SPF w/ Preemption Simulation ---")
    kv_manager_preempt = MockKVCacheManager(num_blocks=KV_CACHE_BLOCKS, block_size=16)
    scheduler_preempt = StallFreeScheduler(token_budget=2048, max_batch_size=16, kv_cache_manager=kv_manager_preempt, policy='spf', allow_preemption=True)
    results_preempt = simulate(list(workload), scheduler_preempt)
    analyze_results(results_preempt)

5. Analysis of Results

After running your completed simulation, you should observe a clear difference between the policies:

  • FCFS: This will likely have the highest average TTFT. Short requests get stuck behind long ones, and everyone waits longer on average.
  • SPF: This policy should significantly reduce the average TTFT compared to FCFS. By prioritizing short prefills, it quickly moves more requests into the decode phase. However, you might see a higher P95 TTFT if a few very long requests are repeatedly passed over and "starved."
  • SPF with Preemption: This should have a low average TTFT, similar to SPF. Its main advantage appears under memory pressure. Where the plain SPF scheduler might stall, the preemptive scheduler can evict a running job to make progress, improving overall system utilization and responsiveness, especially at the tail (P95).

These trade-offs are fundamental to scheduling. There is no single "best" policy; the optimal choice depends on the service's goals (e.g., minimizing average latency vs. ensuring fairness and preventing starvation).

Conclusion

Congratulations on adding a critical layer of intelligence to your inference engine! You've moved beyond simply processing requests to actively managing them based on policies that impact user experience.

Key Takeaways:

  • Scheduling policies matter: A simple FCFS policy can lead to head-of-line blocking, degrading user-perceived latency.
  • TTFT is a key metric: Policies like Shortest Prompt First (SPF) can drastically improve average TTFT by prioritizing requests that are quick to prefill.
  • Preemption is necessary for robustness: Under memory pressure, preempting lower-priority jobs is essential to unblock higher-priority work and maintain system throughput.
  • There are always trade-offs: Optimizing for average latency (like with SPF) can sometimes lead to starvation for long-running jobs, a classic scheduling dilemma.

Preview of the Next Lesson:

So far, we have built and optimized a sophisticated single-GPU inference engine. But what happens when a model is simply too large to fit on one GPU? The 70B models that are standard today require hundreds of gigabytes of VRAM, far exceeding any single device.

In the next lesson, we will enter the world of distributed systems. You will learn the principles of tensor parallelism and pipeline parallelism, the foundational techniques for sharding a massive model across multiple GPUs and orchestrating their computation. You will start by setting up a multi-GPU communication group and begin implementing the first form of distributed inference.

Can't find a good explanation? Sign up and we'll make it for you

Sign up