Prhub

#44726 [Bugfix][Core] Close underlying iterator in merge_async_iterators single-iterator fast path

原始 PR 作者 Sunt-ing 合并时间 2026-07-07 20:13 文件变更 2 提交数 4 评论 2 代码增减 +31 / -2

执行摘要

修复异步迭代器合并时资源泄漏

当合并的生成器被关闭(如消费者取消或断开连接),单迭代器快速路径未关闭底层生成器,导致 AsyncLLM.generate() 持续运行,请求在 OutputProcessor 中滞留,消耗调度器/引擎资源直到自然完成。单 prompt 是 OpenAI completion/chat、pooling、scoring 等功能的常见路径。

建议精读并合并此 PR。修复了一个重要的资源泄漏 bug,且实现简洁高效,设计模式(try/finally 确保清理)值得学习。

讨论亮点

评审者 njhill 提出优化:避免在正常路径上引入额外的 asyncio 边界。作者接受并调整了实现,将原本可能因 aclose() 调用而产生的额外异步等待消除,仅保留资源泄漏时的关闭操作。

实现拆解

  1. 修改 vllm/utils/async_utils.py 中的 merge_async_iterators 函数:将单迭代器快速路径的 async for 循环包裹在 try/finally 块中。在 finally 块中检查底层迭代器是否仍非 None,若是则调用 aclose() 关闭。同时引入局部变量 iterator 保存迭代器引用,以便在正常迭代结束时置为 None,避免重复关闭。
  2. 新增测试用例 tests/utils_/test_async_utils.py:添加 test_merge_async_iterators_single_closes_underlying 测试,验证单迭代器路径下关闭合并生成器后底层生成器的 finally 块被执行。测试使用一个带有 finally 块的异步生成器,消费一项后关闭合并器,断言 closed 标志为 True
文件 模块 状态 重要度
vllm/utils/async_utils.py 异步工具 modified 6.71
tests/utils_/test_async_utils.py 异步工具 modified 5.95

关键符号

merge_async_iterators test_merge_async_iterators_single_closes_underlying

关键源码片段

vllm/utils/async_utils.py core-logic

核心逻辑变更:修复单迭代器快速路径未关闭底层迭代器的 bug。

async def merge_async_iterators(
    *iterators: AsyncGenerator[T, None],
) -> AsyncGenerator[tuple[int, T], None]:
    # 单迭代器快速路径:需要确保在生成器关闭时清理底层迭代器
    if len(iterators) == 1:
        iterator: AsyncGenerator[T, None] | None = iterators[0]
        try:
            async for item in iterator: # type: ignore[union-attr]
                yield 0, item
            iterator = None # 正常结束,不再需要关闭
        finally:
            # 如果迭代器未正常结束(如 aclose 或异常),则显式关闭
            if iterator is not None:
                with contextlib.suppress(BaseException):
                    await iterator.aclose()
        return
    # 多迭代器路径代码保持不变 ...
tests/utils_/test_async_utils.py test-coverage

新增回归测试,验证单迭代器路径下底层迭代器被正确关闭。

@pytest.mark.asyncio
async def test_merge_async_iterators_single_closes_underlying():
    # 验证单迭代器路径下关闭合并器后底层生成器被正确清理
    closed = False
​
    async def gen():
        nonlocal closed
        try:
            while True:
                yield "x"
                await asyncio.sleep(0.01)
        finally:
            closed = True # 标记生成器已被关闭
​
    merged = merge_async_iterators(gen())
    assert await anext(merged) == (0, "x")
    await merged.aclose()
    assert closed # 底层生成器应已被关闭

评论区精华

避免正常路径的额外 await 性能

评审者 njhill 建议避免在正常路径上引入额外的 asyncio 边界,即不要在正常完成时调用 aclose。

结论:作者采纳建议,通过设置 `iterator = None` 来标记正常完成,`finally` 块仅当 `iterator is not None` 时调用 `aclose()`。 · 已解决

风险与影响

风险较低。变更仅影响单迭代器快速路径,且通过 try/finally 和条件检查确保仅在需要时执行 aclose(),不会影响正常流程。但需确保 iterator.aclose()BaseException 被抑制时不会引入其他副作用(已使用 contextlib.suppress)。

修复了单 prompt 请求在消费者取消后资源泄漏的问题,影响范围包括 OpenAI completion/chat、pooling、scoring 等功能。对正常流程无性能影响,因为 aclose() 仅在异常或提前关闭时被调用。

核心路径变更

关联 Issue

未识别关联 Issue

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

完整报告

参与讨论