乐于分享
好东西不私藏

vllm -- 源码剖析 (monitor_engine_liveness)

vllm -- 源码剖析 (monitor_engine_liveness)

声明: 纯兴趣爱好,如有疏漏敬请谅解。坚持每日钻研一点、每日精进一点,在持续实践中梳理技术逻辑,汇总编写文档,沉淀可复用的技术经验。

源码版本: v0.21.0

学习相关源码路径: 

https://github.com/vllm-project/vllm/blob/v0.21.0/vllm/entrypoints/cli/main.py

概述

本文将针对run_headless启动模式下的 CoreEngineProcManager,对其引擎存活监控方法 monitor_engine_liveness 进行完整流程拆解与源码剖析。

monitor_engine_liveness函数

代码:

def monitor_engine_liveness(self) -> None:    """Monitor engine core process liveness."""    sentinel_to_proc = {proc.sentinel: proc for proc in self.processes}    sentinels = set(sentinel_to_proc.keys())    while sentinels and not self.manager_stopped.is_set():        died_sentinels = connection.wait(sentinels, timeout=1)        for sentinel in died_sentinels:            proc = sentinel_to_proc.pop(cast(int, sentinel))            exitcode = proc.exitcode            if exitcode != 0 and not self.manager_stopped.is_set():                self.failed_proc_name = proc.name        if died_sentinels:           # Any engine exit currently triggers a shutdown. Future           # work (e.g., Elastic and fault-tolerant EP) will add finer-grained           # handling for different exit scenarios.           break    self.shutdown()

1.  构建哨兵对象和进程对象字典 -- sentinel_to_proc (key: 哨兵对象, value: 进程对象)

         key: 哨兵对象, 本质上是int。通常是一个文件描述符。

        self.processes为list[BaseProcess], 至于如何初始化,我们在下一篇文章初始化CoreEngineProcManager流程中详细介绍。

2.  构建哨兵对象集合 -- sentinels

3.  循环 -- 条件: 哨兵对象非空 (sentinels)  并且  threading.Event无事件信号

        3.1 监听进程结束事件(sentinels) 

        3.2 处理结束的进程

                从sentinel_to_proc删除key并获取进程对象,

                获取进程退出码exitcode

                在意外情况下,记录首个失败进程名

        3.3 有任意一个进程结束,跳出循环

4. shutdown退出函数

流程图:

知识点:

BaseProcess:

BaseProcess是Python multiprocessing库中的进程基类,用于表示和管理操作系统级别的子进程。

核心属性:

属性
类型
说明
name
str
进程的名字(如 "Engine_0")
pid
int
进程 ID(只读,启动后可用)
ident
int
进程标识符(启动后与 pid 相同)
exitcode
int | None
退出码(None = 运行中,0 = 正常,非0 = 异常)
sentinel
int/handle
哨兵对象(用于监听进程结束)
daemon
bool
是否为守护进程

核心方法:

方法
说明
示例
.start()
启动进程
proc.start()
.join(timeout)
阻塞等待进程结束
proc.join(timeout=5)
.is_alive()
检查进程是否运行中
if proc.is_alive(): ...
.terminate()
温和地终止进程(SIGTERM)
proc.terminate()
.kill()
强制杀死进程(SIGKILL)
proc.kill()
.close()
关闭进程对象
proc.close()

threading:

threading是Python标准库的线程模块,用于实现多线程编程。线程 = 轻量级的执行流,可以在同一进程内并发运行

事例:

import threading# 创建一个线程thread = threading.Thread(target=函数名)# 启动线程thread.start()

事件

event = threading.Event()# 检查状态event.is_set()      # False(未设置)# 设置事件(发信号)event.set()         # True# 清除事件event.clear()       # False# 等待事件event.wait()        # 阻塞等待直到 set() 被调用event.wait(timeout=5)  # 等最多5秒

vllm中用法

self.manager_stopped = threading.Event()  # 创建停止信号# 监控线程中while not self.manager_stopped.is_set():  # 检查是否停止# 主线程调用self.manager_stopped.set()  # 发送停止信号

常用方法对比表:

类/方法
用途
用法
Thread
创建和管理线程
Thread(target=func).start()
Event
简单的 True/False 信号
.set()
 / .is_set()
Lock
保护共享资源
with lock:
Condition
复杂的条件同步
.wait()
 / .notify()
Semaphore
计数同步
Semaphore(N)
RLock
可重入锁
同 Lock

connection.wait

connection.wait(sentinels, timeout=1)是multiprocessing模块中的一个函数,用于高效地监听多个进程的结束事件。

参数:

参数
类型
说明
sentinels
list/set
要监听的 sentinel 对象列表
timeout
float
最多等待多少秒(可选,None表示永久等待)

返回值:

返回已死亡进程的sentinels列表