Prhub

#50804 [CI] Stabilize tensor IPC multiprocessing tests

原始 PR 作者 AndreasKaratzas 合并时间 2026-08-12 10:16 文件变更 1 提交数 1 评论 0 代码增减 +271 / -176

执行摘要

加固 tensor IPC 多进程测试,修复 ROCm 超时误报

PR body 指出 #32104 为测试 harness 引入了短的父进程截止时间和不完整的失败清理:在 ROCm 上父进程在子进程仍在导入 torch 时超时,接着销毁同步对象,在 engine test step 产生误导性的 SemLock 错误。本 PR 的目标是在不改变 spawn 语义(vLLM 的 torch_shm 要求与 PyTorch 加速器共享约束)的前提下,让测试在慢启动环境下稳定运行,并让真实失败可诊断。

值得精读。虽然只是单文件测试改动,但它展示了多进程测试中两个常被忽视的设计点:一是通过 get_context('spawn') 局部化 start method,避免全局 set_start_method 的副作用;二是失败路径的有界进程回收与带状态信息的超时诊断。_cleanup_processes 的 join / terminate / kill 阶梯式清理可以推广到 vLLM 其他多进程测试,建议在后续 IPC 与 kernel 相关测试中复用该模式。

讨论亮点

该 PR 没有实质性的 review 讨论:claude[bot] 仅提示仓库配置为手动 review(可评论 @claude review 触发一次性审查),Isotr0py 直接 APPROVED 且无公开评论。所有设计决策(共享截止时间、role 标记、有界清理)都记录在 PR body 中,合并时无未解决疑问。

实现拆解

  1. 显式 spawn context 替换全局 start method:删除模块级 setup_multiprocessing fixture(原实现调用 torch_mp.set_start_method("spawn", force=True)),改为 _MP_CTX = torch_mp.get_context("spawn")。所有测试中的 torch_mp.Queuemp.Queuemp.Eventmp.Barriermp.Process 均改用 _MP_CTX 创建,并为 encoder / decoder 进程命名。这样既保证 spawn 语义,又不影响进程内其他测试的 multiprocessing 默认值,也确保父子进程共享同一套 IPC 原语实现。
  2. 统一的结果收集:新增 _collect_process_results(result_queue, count, processes),用一个 60 秒的共享 deadline 收集全部结果;超时抛出包含各进程 pidis_alive()exitcodeTimeoutError,便于定位是哪个 worker 卡住。
  3. 分层超时策略encoder_ready.waitretrieval_done.wait 从硬编码的 10 秒 / 30 秒统一为 _PROCESS_RESULT_TIMEOUT(60 秒)以容忍慢导入;而 payload_queue.get 保持 5 秒、payload_queue.put 保持 10 秒,确保真实 IPC 停滞仍能快速失败。
  4. 按 role 区分结果encoder_processdecoder_process 的结果 dict 增加 'role' 字段,主进程用 {result['role']: result} 组织结果,不再依赖跨进程队列的到达顺序。
  5. 有界失败清理:新增 _cleanup_processes,先以 10 秒预算 join 所有已启动进程,对仍存活的依次 terminate、join(5s)、kill、join(5s),并返回未正常退出进程列表;测试通过 finally 保证失败或超时路径也能清理,避免僵尸进程污染后续测试或 CI 步骤。多个测试(test_cuda_tensor_queue_basictest_multiple_api_servers_to_engine 等)及 api_server_worker 均接入这些辅助函数。
文件 模块 状态 重要度
tests/v1/test_tensor_ipc_queue.py 多进程测试 modified 6.65

关键符号

setup_multiprocessing _collect_process_results _cleanup_processes encoder_process decoder_process api_server_worker test_cuda_tensor_queue_basic test_multiple_api_servers_to_engine

关键源码片段

tests/v1/test_tensor_ipc_queue.py test-coverage

PR 唯一变更文件,重写多进程测试 harness:显式 spawn context、共享 deadline、role 标记与有界清理,直接消除 ROCm 上误导性的 SemLock 误报。

# 所有 IPC 原语统一从显式 spawn context 创建,避免调用全局 set_start_method,
# 防止影响其他测试的进程模型,同时满足 vLLM 的 torch_shm 与
# PyTorch 加速器共享约束。
_MP_CTX = torch_mp.get_context('spawn')
# 共享的启动 + 结果截止时间:ROCm 环境子进程导入 torch 可能很慢
_PROCESS_RESULT_TIMEOUT = 60.0
# 有界清理总预算,避免失败路径上测试进程无限挂起
_PROCESS_CLEANUP_TIMEOUT = 10.0
​
​
def _collect_process_results(result_queue, count, processes):
    '''在同一个全局 deadline 内收集 count 个 worker 结果。'''
    deadline = time.monotonic() + _PROCESS_RESULT_TIMEOUT
    results = []
    for _ in range(count):
        remaining = deadline - time.monotonic()
        if remaining <= 0:
            # 超时即把每个进程的 pid、存活状态、退出码拼进错误信息,便于 CI 定位
            states = ', '.join(
                f'{p.name}(pid={p.pid}, alive={p.is_alive()}, exit={p.exitcode})'
                for p in processes
            )
            raise TimeoutError(f'Timed out waiting for worker results: {states}')
        try:
            results.append(result_queue.get(timeout=remaining))
        except Empty as exc:
            states = ', '.join(
                f'{p.name}(pid={p.pid}, alive={p.is_alive()}, exit={p.exitcode})'
                for p in processes
            )
            raise TimeoutError(f'Timed out waiting for worker results: {states}') from exc
    return results
​
​
def _cleanup_processes(processes):
    '''先有界 join,再对超时进程逐个 terminate、kill,返回未正常退出进程描述。'''
    started = [p for p in processes if p.pid is not None]
    deadline = time.monotonic() + _PROCESS_CLEANUP_TIMEOUT
    for process in started:
        process.join(timeout=max(0.0, deadline - time.monotonic()))
​
    failures = [
        f'{p.name}(pid={p.pid}, alive={p.is_alive()}, exit={p.exitcode})'
        for p in started
        if p.is_alive() or p.exitcode != 0
    ]
    alive = [p for p in started if p.is_alive()]
    for process in alive:
        process.terminate()
    for process in alive:
        process.join(timeout=5.0)
    for process in alive:
        if process.is_alive():
            process.kill()
            process.join(timeout=5.0)
    return failures

评论区精华

没有提炼出高价值讨论线程

当前评论区没有形成足够清晰的争议点或结论,后续有更多讨论时会体现在这里。

风险与影响

  1. 超时放宽可能掩盖真实问题:worker 启动 / 结果 deadline 从 10 / 30 秒统一放宽到 60 秒共享预算,某些真实卡死可能需要多等约 50 秒才报错;payload 传输路径仍保留 5 秒 / 10 秒短超时,因此 IPC 停滞后仍能较快失败,但启动阶段的异常会被容忍更久。
  2. 硬编码阈值未参数化_PROCESS_RESULT_TIMEOUT_PROCESS_CLEANUP_TIMEOUT 是模块级常量,在极端慢的 CI 或高负载共享机器上仍可能不足,且没有提供环境变量覆盖。
  3. 强清理可能残留资源_cleanup_processes 的 terminate / kill 序列能避免僵尸进程,但 kill() 强制结束可能在极端情况下留下未释放的共享内存句柄,需要观察 ROCm CI 是否出现新的资源泄漏迹象。
  4. 平台覆盖有限:PR body 只报告了本地 12 个测试通过(ROCm 环境),未说明 CUDA / XPU 等平台的 CI 结果;spawn context 的行为在所有平台一致,但超时参数在不同性能环境下可能需要微调。

对运行时无影响(纯测试文件)。对 CI 的影响显著:修复 ROCm 上 engine test step 的 SemLock 误报,提高 tensor IPC 测试在 CUDA / ROCm 上的稳定性;失败时错误信息包含各进程的 pid、存活状态与退出码,排障体验明显改善。对团队的意义在于沉淀了一套多进程测试 harness 模式(显式 context + 共享 deadline + 有界清理),后续其他 IPC 或分布式测试可以直接复用 _MP_CTX_collect_process_results_cleanup_processes

仅测试路径变更 超时阈值未参数化 依赖 terminate/kill 强清理 验证以 ROCm 为主

关联 Issue

未识别关联 Issue

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

完整报告

参与讨论