Prhub

#46051 [Rust Frontend][Perf] Use dedicated runtime for HTTP/request-processing/ZMQ

原始 PR 作者 BugenZhao 合并时间 2026-06-23 12:03 文件变更 13 提交数 9 评论 4 代码增减 +328 / -15

执行摘要

分离 HTTP/ 请求 /ZMQ 为独立 Tokio runtime,提升高并发性能

PR body 指出,高并发工作负载下 CPU 密集的请求预处理和繁忙的引擎传输路径会阻塞 HTTP 工作线程,导致 /health 等轻量接口的尾部延迟劣化,请求吞吐受到抑制。分离 runtime 可以隔离这些负载,提升整体响应性。

推荐精读此 PR,特别是 BackgroundShutdownRuntime 的设计、Tower middleware offload 模式以及 benchmark 方法。这些实践可复用于其他需要隔离工作负载的场景。团队应考虑将此模式推广到 Python 前端的类似瓶颈。

讨论亮点
  1. ZMQ 线程数选择runtime.rs:46):njhill 提问是否可用 1 线程减少调优复杂度,BugenZhao 回复 benchmark 显示 4 线程在吞吐和延迟上略优,且保留环境变量作为调优入口,结论保留 4 作为默认。
  2. Offload 路径粒度offload.rs:28):njhill 质疑是否应更精细控制(如根据 prompt 长度),BugenZhao 认为路径匹配足以满足当前需求,避免干扰 /health 等控制面,并接受建议将 /inference/v1/generate 替换为 /detokenize,后续 commit 已补充 /detokenize

实现拆解

  1. 创建 BackgroundShutdownRuntime 封装rust/src/engine-core-client/src/runtime.rs):新增 BackgroundShutdownRuntime 结构体,包装 Tokio Runtime 并在 Drop 时调用 shutdown_background() 实现无阻塞关闭,同时提供 Deref/DerefMut 以透明代理。
  2. 构建 ZMQ runtime(同上文件):build_zmq_runtime() 创建多线程 Tokio runtime,线程数默认为 4(可通过 VLLM_RS_ZMQ_WORKER_THREADS 覆盖),用于驱动 ZMQ 发送/接收、输出分发、中止处理和协调任务。
  3. 构建 Request runtimerust/src/server/src/runtime.rs):build_request_runtime() 创建多线程 Tokio runtime,线程数默认为 available_parallelism 上限 32(可通过 VLLM_RS_REQUEST_WORKER_THREADS 覆盖),用于运行 CPU 密集的请求处理路径。
  4. 实现 Tower 中间件进行 offloadrust/src/server/src/middleware/offload.rs):新增 RequestRuntimeService,在 call() 中匹配预定义的 OFFLOADED_PATHS(如 /v1/chat/completions/tokenize),将匹配的请求通过 AbortOnDropHandle spawn 到 Request runtime 上执行,不阻塞 HTTP runtime。
  5. 集成到核心结构EngineCoreClient 中持有一个 BackgroundShutdownRuntime,将输出循环、分发循环和中止循环全部 spawn 到该 runtime 上(client.rs);AppState 增加 OnceLock<BackgroundShutdownRuntime> 懒初始化 Request runtime,并通过 request_runtime() 暴露(state.rs)。
文件 模块 状态 重要度
rust/src/engine-core-client/src/runtime.rs 引擎客户端 added 8.98
rust/src/server/src/middleware/offload.rs 服务器 added 8.97
rust/src/server/src/runtime.rs 服务器 added 7.62
rust/src/engine-core-client/src/client/imp.rs 引擎客户端 modified 6.81
rust/src/engine-core-client/src/client.rs 引擎客户端 modified 6.71
rust/src/server/src/state.rs 服务器 modified 6.02
rust/src/server/src/lib.rs 服务器 modified 5.02
rust/src/engine-core-client/src/error.rs 引擎客户端 modified 4.89
rust/src/engine-core-client/src/lib.rs 引擎客户端 modified 4.67
rust/src/server/src/middleware/mod.rs 服务器 modified 4.52
rust/src/chat/src/multimodal.rs 聊天模块 modified 4.3
rust/src/server/src/routes.rs 服务器 modified 4.3

关键符号

build_zmq_runtime build_request_runtime request_runtime_layer RequestRuntimeService::call should_offload zmq_worker_threads request_worker_threads AppState::request_runtime BackgroundShutdownRuntime::drop

关键源码片段

rust/src/engine-core-client/src/runtime.rs core-logic

新增 ZMQ runtime 构建函数和 BackgroundShutdownRuntime 封装,是 runtime 隔离的基础。

// rust/src/engine-core-client/src/runtime.rs
// 封装 Tokio Runtime,Drop 时后台关闭,防止嵌套阻塞
pub struct BackgroundShutdownRuntime(ManuallyDrop<Runtime>);impl Drop for BackgroundShutdownRuntime {
    fn drop(&mut self) {
        // 安全:仅一次 take,后续不再使用
        let runtime = unsafe { ManuallyDrop::take(&mut self.0) };
        runtime.shutdown_background(); // 非阻塞关闭
    }
}impl Deref for BackgroundShutdownRuntime {
    type Target = Runtime;
    fn deref(&self) -> &Self::Target { &self.0 }
}impl DerefMut for BackgroundShutdownRuntime {
    fn deref_mut(&mut self) -> &mut Self::Target { &mut self.0 }
}impl From<Runtime> for BackgroundShutdownRuntime {
    fn from(runtime: Runtime) -> Self { Self(ManuallyDrop::new(runtime)) }
}// ZMQ runtime 默认 4 线程,多 engine 共享 socket,benchmark 显示足够
const ZMQ_WORKER_THREADS_ENV: &str = "VLLM_RS_ZMQ_WORKER_THREADS";
const DEFAULT_ZMQ_WORKER_THREADS: usize = 4;static ZMQ_RUNTIME_SEQUENCE: OnceLock<AtomicUsize> = OnceLock::new();/// 构建 ZMQ runtime,每次调用生成带独立序列号线程名的 runtime
pub(crate) fn build_zmq_runtime() -> BackgroundShutdownRuntime {
    let sequence = ZMQ_RUNTIME_SEQUENCE
        .get_or_init(|| AtomicUsize::new(0))
        .fetch_add(1, Ordering::Relaxed);    tokio::runtime::Builder::new_multi_thread()
        .worker_threads(zmq_worker_threads())
        .thread_name_fn(move || format!("vllm-zmq-{sequence}"))
        .enable_all()
        .build()
        .expect("failed to build vLLM ZMQ runtime")
        .into()
}fn zmq_worker_threads() -> usize {
    std::env::var(ZMQ_WORKER_THREADS_ENV)
        .ok()
        .and_then(|v| v.parse().ok())
        .filter(|&v| v > 0)
        .unwrap_or(DEFAULT_ZMQ_WORKER_THREADS)
}
rust/src/server/src/middleware/offload.rs core-logic

新增 Tower middleware,将 CPU 密集型请求 offload 到 Request runtime,核心性能优化点。

// rust/src/server/src/middleware/offload.rs
// 需要 offload 的路径列表(HTTP + gRPC)
const OFFLOADED_PATHS: &[&str] = &[
    "/v1/chat/completions",
    "/v1/completions",
    "/tokenize",
    "/detokenize",
    "/inference/v1/generate",
    "/vllm.Generate/Generate",
    "/vllm.Generate/GenerateStream",
];// 返回 Tower layer,包装 inner service
pub(crate) fn request_runtime_layer<S>(
    state: Arc<AppState>,
) -> impl tower::Layer<S, Service = RequestRuntimeService<S>> + Clone {
    layer_fn(move |inner| RequestRuntimeService {
        inner,
        state: state.clone(),
    })
}#[derive(Clone)]
pub(crate) struct RequestRuntimeService<S> {
    inner: S,
    state: Arc<AppState>,
}impl<S, B> Service<Request<B>> for RequestRuntimeService<S>
where /* ... bounds ... */
{
    type Future = BoxFuture<'static, Result<Self::Response, Self::Error>>;    fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
        self.inner.poll_ready(cx)
    }    fn call(&mut self, req: Request<B>) -> Self::Future {
        if !should_offload(req.uri().path()) {
            // 轻量路径直接走 HTTP runtime
            return Box::pin(self.inner.call(req));
        }        // 替换 inner 为克隆,避免 move 冲突
        let clone = self.inner.clone();
        let mut inner = std::mem::replace(&mut self.inner, clone);
        // 在 Request runtime 上执行,AbortOnDropHandle 处理取消
        let task = AbortOnDropHandle::new(
            self.state.request_runtime().spawn(inner.call(req))
        );        Box::pin(async move {
            match task.await {
                Ok(result) => result,
                Err(error) => {
                    error!(%error, "request runtime task failed");
                    Ok(S::Response::request_runtime_error_response())
                }
            }
        })
    }
}// 辅助 trait 生成错误响应
trait RequestRuntimeErrorResponse {
    fn request_runtime_error_response() -> Self;
}impl RequestRuntimeErrorResponse for Response {
    fn request_runtime_error_response() -> Self {
        server_error!("request runtime task failed").into_response()
    }
}impl RequestRuntimeErrorResponse for axum::http::Response<tonic::body::Body> {
    fn request_runtime_error_response() -> Self {
        Status::internal("request runtime task failed").into_http()
    }
}fn should_offload(path: &str) -> bool {
    OFFLOADED_PATHS.contains(&path)
}#[cfg(test)]
mod tests {
    // 测试验证 offload 路径正确性
}
rust/src/server/src/runtime.rs core-logic

新增 Request runtime 构建逻辑,封装了线程数配置和 background shutdown。

// rust/src/server/src/runtime.rs
use tokio::runtime::Builder;
use vllm_engine_core_client::runtime::BackgroundShutdownRuntime;const REQUEST_WORKER_THREADS_ENV: &str = "VLLM_RS_REQUEST_WORKER_THREADS";
const DEFAULT_MAX_REQUEST_WORKER_THREADS: usize = 32;pub(crate) fn build_request_runtime() -> BackgroundShutdownRuntime {
    Builder::new_multi_thread()
        .enable_all()
        .thread_name("vllm-request")
        .worker_threads(request_worker_threads())
        .build()
        .expect("failed to build request runtime")
        .into()
}fn request_worker_threads() -> usize {
    // 优先使用环境变量
    if let Some(value) = std::env::var_os(REQUEST_WORKER_THREADS_ENV) {
        match value.to_string_lossy().parse::<usize>() {
            Ok(worker_threads) if worker_threads > 0 => return worker_threads,
            _ => warn!("ignoring invalid {}: {:?}", REQUEST_WORKER_THREADS_ENV, value),
        }
    }
    // 否则使用可用并行度,上限 32
    std::thread::available_parallelism()
        .map(|p| {
            let available = p.get();
            let capped = available.min(DEFAULT_MAX_REQUEST_WORKER_THREADS);
            if capped < available {
                info!("capping request runtime worker threads from {} to {}, set {} to override",
                    available, capped, REQUEST_WORKER_THREADS_ENV);
            }
            capped
        })
        .unwrap_or(DEFAULT_MAX_REQUEST_WORKER_THREADS)
}

评论区精华

ZMQ runtime worker 线程数默认值 question

njhill 询问是否可以用 1 线程以减少调优参数;BugenZhao 回应 benchmark 显示 4 线程在吞吐和延迟上更优,且保留环境变量供调优。

结论:保留默认 4 线程,并通过环境变量 `VLLM_RS_ZMQ_WORKER_THREADS` 可配。 · 已解决

Offload 路径的粒度控制 设计

njhill 建议更精细地决定哪些请求应该 offload(例如基于 prompt 长度),并提到 gRPC 是否也适用;BugenZhao 认为当前基于路径匹配足够,承诺将 `/inference/v1/generate` 替换为 `/detokenize`,后续已实现。

结论:接受建议,调整 offload 路径列表,暂时保持简单匹配,未来可根据流量优化。 · 已解决

风险与影响

  1. Request runtime 任务失败:若 request_runtime.spawn 的任务 panic 或 abort,中间件会返回 500 错误,可能导致正在处理的请求失败(offload.rsrequest_runtime_error_response())。
  2. 线程资源配置:新增两个 runtime 可能增加线程数和内存开销,默认值(ZMQ 4 线程、Request 最多 32 线程)在高并发下合理,但若环境变量设置过大可能浪费资源。
  3. 路径覆盖遗漏OFFLOADED_PATHS 列表可能遗漏某些 CPU 密集型路径,导致仍阻塞 HTTP runtime,需持续根据业务场景更新。
  4. 背景关闭顺序BackgroundShutdownRuntimeDrop 时后台关闭,可能与其他资源的关闭顺序产生竞态,需确保 runtime 在相关 I/O 完成后释放。

本 PR 仅影响 Rust 前端服务层,不涉及 Python 后端、模型加载或推理核心。用户可见的改善包括高并发下 /health 响应更稳定、请求吞吐提升。团队需要维护新增的两个环境变量配置,并理解三层 runtime 的线程模型。对于仅使用轻量负载的用户,默认配置不产生明显变化;对于高负载场景,性能收益显著。

核心路径变更 运行时资源竞争 线程数配置误导 offload 路径遗漏

关联 Issue

未识别关联 Issue

当前没有检测到明确关联的 Issue 链接,后续同步到相关引用后会出现在这里。

完整报告

参与讨论