Skip to content

Preempt y prefill throttle: defensa propia cuando la caché KV está ajustada

源码版本v0.25.1

Responsabilidades

La caché KV (KV cache) es el recurso más escaso del planificador v1: la cantidad de bloques de VRAM en una GPU se fija durante la fase de profiling. Cuando un request de la cola running necesita avanzar más tokens pero KVCacheManager.allocate_slots no obtiene bloques nuevos, el planificador (scheduler) debe decidir: echar a un request que ya está en running, devolver sus bloques KV y reintentar. Este mecanismo se llama preempt (preemption) y su entrada es Scheduler._preempt_request (_preempt_request:1140).

Otro mecanismo relacionado pero más sutil es el prefill throttle: en escenarios DP (data parallel), los pasos de prefill de cada rank deben alinearse, o el rank más rápido arrastra al lento a esperar. Cuando el planificador recibe throttle_prefills=True, pospone el prefill chunk en curso y los nuevos prefills a un paso alineado, dejando que el decode llene el paso (defer_prefills:436-438). Ambos se tratan juntos porque son comportamientos de contención: "renunciar a avanzar el prefill en cierto paso".

Motivación de diseño

¿Por qué no simplemente rechazar el request nuevo y hacer preemption en su lugar?

  • Mantener el progreso compensa: un request que ya ha hecho medio prefill y se rechaza pierde todo el KV calculado. Con preemption, sus bloques KV se devuelven, el estado pasa a PREEMPTED y num_computed_tokens se pone a cero (preempt 状态:1152-1153), pero el request no se descarta y se vuelve a poner a la cabeza de waiting (prepend_request:1161); el siguiente paso puede hacer prefill de nuevo (la parte con hit en prefix cache se reutiliza).
  • Política de preemption intercambiable: el modo priority elige como víctima el running con (priority, arrival_time) máximo (priority preempt:547-552); el modo FCFS hace directamente self.running.pop() (FCFS preempt:572). La semántica es coherente con la de desencolado.
  • El paso con preemption no acepta requests nuevos: la entrada del bucle waiting en schedule() chequea not preempted_reqs (waiting 入口:637), evitando que un request devuelto a waiting sea sacado otra vez en el mismo paso y compita con el running que se está preventando.
  • El prefill throttle es equidad entre DP: los pasos de prefill entre ranks DP deben alinearse; si no, un rank ya en decode mientras otro sigue en prefill arrastra el throughput. El throttle hace que el paso no alineado solo ejecute decode y pospone el prefill al paso alineado (defer_prefills:436-438).
  • Cobertura automática en capacidad limitada: prefill_capacity_bound (prefill_capacity_bound:290) registra si el último paso alineado ya vació waiting; si lo vació, no hace falta volver a hacer throttle, porque continuar throttlendo solo deja la GPU ociosa.
  • Reserva separada para async KV load: _inflight_prefill_reserved_blocks (_inflight_prefill_reserved_blocks:2400) suma los bloques que aún necesitan todos los preflights en vuelo; los requests con carga asíncrona no se preemptan y deben tener bloques reservados suficientes para entrar, evitando deadlock.

Archivos clave

Flujo de datos

El disparador de preemption está dentro del bucle running, cuando allocate_slots devuelve None:

python
# The request cannot be scheduled.
# Preempt the lowest-priority request.
if self.policy == SchedulingPolicy.PRIORITY:
    preempted_req = max(
        self.running,
        key=lambda r: (r.priority, r.arrival_time),
    )
    self.running.remove(preempted_req)
    if preempted_req in scheduled_running_reqs:
        # 已经在本步分到 token 的牺牲者,把预算还回去
        preempted_req_id = preempted_req.request_id
        scheduled_running_reqs.remove(preempted_req)
        token_budget += num_scheduled_tokens.pop(preempted_req_id)
        req_to_new_blocks.pop(preempted_req_id)
        ...
        req_index -= 1
else:
    preempted_req = self.running.pop()

self._preempt_request(preempted_req, scheduled_timestamp)
preempted_reqs.append(preempted_req)
if preempted_req == request:
    # No more request to preempt. Cannot schedule this request.
    break

Esto está en scheduler.py:545-582. En el modo priority hay un detalle: si la víctima seleccionada ya había sido planificada en este paso (está en scheduled_running_reqs), hay que devolver su presupuesto de tokens, sus bloques nuevos, sus spec tokens y su presupuesto de encoder (还预算:553-570), y hacer req_index -= 1 porque la lista running perdió un elemento. preempted_req == request es la condición de parada: si la preemption llega a uno mismo, se admite el fallo y se rompe el bucle.

El request echado pasa por _preempt_request:

python
def _preempt_request(self, request: Request, timestamp: float) -> None:
    """Preempt a request and put it back to the waiting queue.

    NOTE: The request should be popped from the running queue outside of this
    method.
    """
    assert request.status == RequestStatus.RUNNING, (
        "Only running requests can be preempted"
    )
    self._free_request_blocks(request)
    self.encoder_cache_manager.free(request)
    self._inflight_prefills.discard(request)
    request.status = RequestStatus.PREEMPTED
    request.num_computed_tokens = 0
    if request.spec_token_ids:
        request.spec_token_ids = []
    request.num_preemptions += 1
    ...
    # Put the request back to the waiting queue.
    self.waiting.prepend_request(request)
    self.reset_preempted_req_ids.add(request.request_id)

Esto está en _preempt_request:1140-1162. Ojo a tres acciones: (1) _free_request_blocks devuelve los bloques KV al block pool (con defer_block_free van a la cola deferred_frees esperando un fence) (_free_request_blocks:2130-2143); (2) encoder_cache_manager.free libera la caché del encoder; (3) num_computed_tokens = 0 resetea el progreso, pero los bloques con hit de prefix cache se preservan dentro de _free_request_blocks (porque son compartidos y no pertenecen en exclusiva al request), así que el siguiente prefill no necesariamente recalcula todo.

La ruta del prefill throttle es independiente. defer_prefills se calcula una vez al inicio de schedule() (defer_prefills:436-438): throttle_prefills and not self.prefill_capacity_bound and any(not r.is_prefill_chunk for r in self.running). Hace falta que se cumplan las tres condiciones para posponer: la flag throttle viene del engine core DP, prefill_capacity_bound indica que en el paso previo se vació waiting, y la última condición garantiza que solo se posponga prefill cuando hay decode en running (decode tiene menor prioridad que un prefill chunk en curso). Luego surte efecto en dos sitios: el prefill chunk en curso dentro de running se salta (running 推迟:467-471), y los prefills nuevos de waiting se rompen directamente (waiting 推迟:801-804). Al terminar un paso no throttle se actualiza prefill_capacity_bound = bool(self.waiting) (capacity_bound 更新:1019-1020); en el siguiente paso throttle, si ve que waiting ya está vacío, no vuelve a hacer throttle.

Límites y fallos

  • Solo se puede preemptar RUNNING: al inicio de _preempt_request se asume request.status == RequestStatus.RUNNING (assert:1146-1148). Los requests en waiting no se pueden "re-preemptar" porque no tienen bloques KV.
  • Assert: el request ya se sacó de running: el docstring deja claro "popped from the running queue outside of this method" (NOTE:1143-1144). _preempt_request solo toca el estado y el regreso a waiting; no manipula la lista self.running, el llamador debe hacer el pop primero.
  • Se vacían los spec_token_ids: al ser preemptado, request.spec_token_ids = [] (spec 清空:1154-1155) evita que el spec decoder devuelva draft tokens caducos en la siguiente re-planificación.
  • num_preemptions se incrementa: en cada preemption request.num_preemptions += 1 (num_preemptions:1156); este contador se usa para marcar preempted=True en las estadísticas de prefix cache (preempted 标记:721) con fines observabilidad.
  • Deferred free evita escrituras concurrentes: cuando defer_block_free=True y last_sched_seq > processed_step_seq, _free_request_blocks no libera inmediatamente sino que encola en deferred_frees esperando un fence (_free_request_blocks:2130-2143), porque bajo async scheduling otro paso en vuelo puede estar escribiendo esos bloques.
  • reset_preempted_req_ids notifica al worker: el id del request preemptado entra en reset_preempted_req_ids (reset_preempted:1162) y se añade a SchedulerOutput.preempted_req_ids (SchedulerOutput.preempted_req_ids:1100) para que el worker sepa que debe limpiar el estado de caché KV correspondiente.
  • El throttle no afecta a decode: defer_prefills solo se salta los request con is_prefill_chunk (running 推迟条件:467); los running en puro decode siguen planificándose con normalidad, garantizando que la GPU no se quede ociosa en el paso con throttle.
  • Async KV load no es preemptable: los requests que van por load_kv_async entran en estado WAITING_FOR_REMOTE_KVS (WAITING_FOR_REMOTE_KVS:950) y no pasan a running; se reservan bloques vía _inflight_prefill_reserved_blocks (_inflight_prefill_reserved_blocks:2400) para evitar que entren a running y luego sean preemptados, lo que destrozaría el estado de la transferencia KV.

Resumen

La preemption es la defensa propia cuando la caché KV está ajustada: se echa de running a un request de baja prioridad, se liberan bloques para el de alta prioridad, y el echado vuelve a la cabeza de waiting para re-planificarse en el siguiente paso. El prefill throttle es otra contención bajo DP: el paso no alineado solo corre decode y pospone el prefill al paso alineado. Ambos renuncian activamente a "meter un prefill más" a cambio de estabilidad del sistema. Para ver dónde encaja la preemption en el bucle principal del planificador, leer Scheduler.schedule; para ver cómo se vuelve a encolar el request echado, leer RequestQueue.

Véase la documentación oficial: Documentación de vLLM · README.