Préemption et prefill throttle : se préserver quand le KV cache est tendu
Responsabilités
Le KV cache est la ressource la plus rare de l'ordonnanceur v1 — le nombre de blocs mémoire sur un GPU est figé dès la phase de profiling. Quand une request de la file running doit calculer davantage de tokens mais que KVCacheManager.allocate_slots ne peut pas fournir de nouveaux blocs, l'ordonnanceur (scheduler) doit décider : évincer une request déjà présente dans running, rendre ses blocs KV, puis réessayer. Ce mécanisme s'appelle la préemption (preemption), et son point d'entrée est Scheduler._preempt_request(_preempt_request:1140).
Un autre mécanisme, plus discret, est le prefill throttle : en DP (parallélisme de données), les étapes prefill de chaque rank doivent s'aligner, sinon le rank rapide traîne en attendant le lent. Quand l'ordonnanceur reçoit throttle_prefills=True, il diffère les chunks prefill en cours et les nouveaux prefill vers une étape d'alignement, laissant le decode remplir le step(defer_prefills:436-438). Ces deux mécanismes sont présentés ensemble car ils relèvent tous deux d'une même retenue : « à une étape donnée, renoncer volontairement à pousser un prefill supplémentaire ».
Motivation de conception
Pourquoi préempter plutôt que de simplement refuser les nouvelles requêtes ?
- Préserver la progression est plus rentable : une request a déjà préfillé plus de la moitié. Si on la refuse, tout le KV calculé est perdu ; la préemption (preemption) rend ses blocs KV, repasse l'état à
PREEMPTED, et remetnum_computed_tokensà zéro(preempt 状态:1152-1153), sans toutefois jeter la request elle-même : elle retourne en tête dewaiting(prepend_request:1161), et pourra être re-prefillée à l'étape suivante (les portions touchant la cache de préfixe sont réutilisables). - Stratégie de préemption interchangeable : en mode priority on évince le running avec le
(priority, arrival_time)maximal(priority preempt:547-552) ; en mode FCFS on faitself.running.pop()(FCFS preempt:572), ce qui reste cohérent avec la sémantique de sortie de file. - Un step ayant préempté n'accepte plus de nouvelles requêtes : l'entrée de la boucle waiting de
schedule()vérifienot preempted_reqs(waiting 入口:637) afin qu'une request réinjectée dans waiting ne soit pas immédiatement extraite pour concourir aux ressources des running en cours de préemption. - prefill throttle = équité DP : les étapes prefill entre ranks DP doivent s'aligner, sinon un rank commence déjà son decode tandis qu'un autre est encore en prefill, et le débit chute. Le throttle fait en sorte que les étapes non alignées ne fassent que du decode, les prefill étant repoussés à l'étape d'alignement commune(
defer_prefills:436-438). - Capacity-bound s'auto-désactive :
prefill_capacity_bound(prefill_capacity_bound:290) enregistre si la dernière étape alignée a vidé waiting ; si oui, plus besoin de throttler — continuer ne ferait qu'occuper le GPU à vide. - Réservation dédiée à l'async KV load :
_inflight_prefill_reserved_blocks(_inflight_prefill_reserved_blocks:2400) recense combien de blocs les prefill en cours nécessitent ; les requests en chargement asynchrone ne sont pas préemptées et ne peuvent entrer que si les blocs réservés suffisent, afin d'éviter un deadlock.
Fichiers clés
_preempt_request:1140— entrée de la préemption : libère les blocs, nettoie l'état, remet dans waiting._free_request_blocks:1149— libère les blocs KV de la request ; sidefer_block_freeest activé, ils entrent dans la file différée.PREEMPTED:1152-1153— étatPREEMPTED,num_computed_tokensremis à zéro.prepend_request:1161-1162— remet en tête dewaiting, ajoute àreset_preempted_req_idspour notifier le worker.抢占循环:534-575— bouclewhile Truede préemption déclenchée à l'échec deallocate_slots.priority 选牺牲者:547-552— en mode priority, sélectionne le running avec le plus haut(priority, arrival_time).FCFS 选牺牲者:572— en FCFS, simple pop en fin de running.defer_prefills:436-438— commutateur principal du DP prefill throttle ; throttle et non saturation => diffère les prefill.running 里推迟 prefill:467-471— les chunks prefill en cours sont sautés sur l'étape throttle.waiting 里推迟 prefill:801-804— sur l'étape throttle, les nouveaux prefill de waiting font directement break.prefill_capacity_bound 更新:1019-1020— met à jour ce flag à la fin de chaque étape non throttle.prefill_capacity_bound 初始化:290— dans__init__, False par défaut._inflight_prefill_reserved_blocks:2400— statistique de blocs réservés pour l'async KV load, évite le deadlock.waiting 入口保护:637— si une préemption a eu lieu sur ce step, on n'accepte plus de waiting.
Flux de données
Le déclencheur de la préemption se trouve dans la boucle running, quand allocate_slots renvoie None :
# 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.
breakCe bloc est dans scheduler.py:545-582. Détail du mode priority : si la victime sélectionnée a déjà été planifiée ce step (présente dans scheduled_running_reqs), il faut rendre son budget de tokens, ses nouveaux blocs, ses spec tokens et son budget encoder(还预算:553-570) et faire req_index -= 1, car la liste running a perdu un élément. preempted_req == request est la condition d'arrêt : si la préemption finit par s'évincer soi-même, on admet l'échec et on break.
La request évincée passe ensuite par _preempt_request :
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)Ce bloc est dans _preempt_request:1140-1162. Trois actions à noter : (1) _free_request_blocks rend les blocs KV au block pool (avec defer_block_free, ils entrent dans la file deferred_frees en attente de fence)(_free_request_blocks:2130-2143) ; (2) encoder_cache_manager.free libère la cache encoder en parallèle ; (3) num_computed_tokens = 0 remet la progression à zéro — toutefois, les blocs touchés par un hit prefix cache sont conservés à l'intérieur de _free_request_blocks (ces blocs sont partagés et n'appartiennent pas exclusivement à la request), si bien que le prochain prefill ne sera pas nécessairement tout recalculé.
Le chemin du prefill throttle est distinct. defer_prefills est calculé une fois en tête 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). Les trois conditions doivent être réunies pour différer : throttle provient du engine core DP, prefill_capacity_bound indique que waiting a été vidé à la dernière étape, et la dernière condition garantit qu'on ne diffère le prefill que lorsque running contient du decode (le decode a une priorité inférieure aux chunks prefill en cours). Puis deux points d'effet : dans running, les chunks prefill en cours sont sautés(running 推迟:467-471) ; dans waiting, les nouveaux prefill font directement break(waiting 推迟:801-804). À la fin d'une étape non throttle, prefill_capacity_bound = bool(self.waiting) est mis à jour(capacity_bound 更新:1019-1020) : à l'étape throttle suivante, si waiting est déjà vide, on cesse de throttler.
Limites et échecs
- Ne peut préempter que RUNNING :
_preempt_requestdébute par une assertionrequest.status == RequestStatus.RUNNING(assert:1146-1148). Une request de waiting ne peut pas être « re-préemptée » car elle ne détient pas de blocs KV. - Assertion : la request a déjà été retirée de running : la docstring précise « popped from the running queue outside of this method »(
NOTE:1143-1144)._preempt_requestne gère que l'état et la réinjection dans waiting ; il ne touche pas àself.running. L'appelant doit d'abord pop. - Spéc_token_ids nettoyés : à la préemption,
request.spec_token_ids = [](spec 清空:1154-1155) pour éviter que le spec decoder ne fournisse des draft tokens périmés lors de la prochaine planification. - Incrémentation de
num_preemptions: à chaque préemption,request.num_preemptions += 1(num_preemptions:1156). Ce compteur est utilisé pour marquerpreempted=Truedans les stats prefix cache(preempted 标记:721), à des fins d'observabilité. - Deferred free contre les écritures concurrentes : dans
_free_request_blocks, sidefer_block_free=Trueetlast_sched_seq > processed_step_seq, la libération n'est pas immédiate : les blocs sont pushés dansdeferred_freesen attente de fence(_free_request_blocks:2130-2143) car, en scheduling asynchrone, un autre step en vol peut encore écrire dans ces blocs. reset_preempted_req_idsnotifie le worker : l'id de la request préemptée est ajouté àreset_preempted_req_ids(reset_preempted:1162) puis placé dansSchedulerOutput.preempted_req_ids(SchedulerOutput.preempted_req_ids:1100), informant le worker qu'il doit nettoyer l'état KV cache correspondant.- throttle n'affecte pas le decode :
defer_prefillsne saute que les requestsis_prefill_chunk(running 推迟条件:467) ; les running de pur decode continuent à être planifiées normalement, afin que le GPU ne reste pas inactif sur l'étape throttle. - async KV load non préemptable : une request passant par
load_kv_asyncentre dans l'étatWAITING_FOR_REMOTE_KVS(WAITING_FOR_REMOTE_KVS:950) ; elle ne rentre pas dans running et s'appuie sur_inflight_prefill_reserved_blocks(_inflight_prefill_reserved_blocks:2400) pour réserver les blocs, évitant qu'elle ne soit préemptée après être entrée dans running, ce qui perturberait l'état du transfert KV.
Résumé
La préemption est l'auto-protection lorsque le KV cache est tendu : on évince de running une request de basse priorité pour libérer des blocs au profit d'une request plus prioritaire, et la victime retourne en tête de waiting en attendant d'être re-planifiée au step suivant. Le prefill throttle est une autre retenue, en scénario DP : les étapes non alignées ne font que du decode, les prefill étant regroupés sur l'étape d'alignement. Dans les deux cas, l'ordonnanceur renonce volontairement à « planifier un prefill de plus » au bénéfice de la stabilité du système. Pour replacer la préemption dans la boucle principale de l'ordonnanceur, lire Scheduler.schedule ; pour comprendre comment une request évincée est remise en file, lire RequestQueue.
Voir la documentation officielle : Documentation vLLM · README