Prhub

#52491 [Bugfix][EPD] Fix encoder round-robin fan-out

原始 PR 作者 AnkitNakhawa 合并时间 2026-08-20 15:30 文件变更 2 提交数 2 评论 6 代码增减 +110 / -2

执行摘要

修复 EPD 代理 encoder 轮询重置问题

Issue #52484 明确指出问题本质:"This is only round-robin within a single incoming request. Since i always starts at 0 for every request, e_urls[0] receives the first multimodal item of every request." 当多数请求只含 1 个多模态 item 时,全部 encoder 流量都会集中到第一个实例,3 个 encoder 的示例中 E0 永远被选中、E2 永远空闲。PR body 承接该 issue,目标是让代理 "distribute encoder sub-requests fairly across encoder instances over time, not restart from e_urls[0] for every incoming request"。

值得快速精读,重点关注两点:(1) 把游标数学抽成纯函数并让测试直接驱动真实模块,避免测试复制生产逻辑的常见反模式;(2) 在单事件循环 async 服务中,用一把锁保护跨请求共享状态的最小化做法。若后续要扩展 EPD 代理(如 encoder 健康检查、加权分配),可在此基础上演进,并建议补一次真实集群端到端验证。

讨论亮点

review 交互非常精简:claude[bot] 提示 fork PR 的自动化 review 被禁用("This pull request is from a fork — automated review is disabled"),最终由 gty111 人工批准并留言 "Thanks for the fix! CC @Isotr0py",Isotr0py 直接批准。真正的技术方案论证发生在 Issue #52484:报告方给出两种修复方向 —— 维护全局轮询游标("because the proxy is async, this should probably be protected with an asyncio.Lock or replaced with another concurrency-safe strategy"),或按 item 随机选择("avoids hot-spotting e_urls[0], although it is not deterministic round-robin")。PR 选择了确定性的全局游标方案,与 issue 建议的第一种方案完全一致,并通过抽出纯函数规避了并发代码难以测试的问题。

实现拆解

  1. 新增模块级游标状态:在 examples/disaggregated/disaggregated_encoder/disagg_epd_proxy.py 顶部新增 encoder_rr_idx = 0(跨请求轮询位置)与 encoder_rr_lock = asyncio.Lock()(并发保护)。代理在单事件循环上并发处理请求,asyncio.Lock 足以保证游标读写原子性。
  2. 抽出纯函数 encoder_rr_assignment:函数接收 e_urls、起始游标 start 与 item 数量 count,返回 URL 列表与下一个游标 (start + count) % len(e_urls)。纯函数化让分配逻辑可以在不启动 HTTP 服务的情况下直接做单元测试,是后续测试配套的关键前提。
  3. 改造 fanout_encoder_primer:删除每次请求从 0 开始的生成器 url_cycle = (e_urls[i % len(e_urls)] for i in range(len(mm_items))),改为 global encoder_rr_idx 声明后在 async with encoder_rr_lock 内一次性调用 encoder_rr_assignment 得到 url_cycle 并推进游标。返回值从生成器变为列表,下游 zip(mm_items, url_cycle) 用法不变。
  4. 新增回归测试:新增 tests/v1/ec_connector/unit/test_epd_proxy_round_robin.py(82 行),参照 tests/v1/kv_connector/unit/test_moriio_proxy_routing.py 的先例,用 importlib.util 从 examples/ 路径加载真实代理模块,fixture 暴露 encoder_rr_assignment_drive 模拟多请求推进游标。4 个测试覆盖多 encoder 数量下的均匀覆盖、bug 原始场景(3 个 encoder 单 item 请求轮转)、变长 item 数下游标连续性、单 encoder 退化场景。
  5. 测试与 CI 配套:本地 pytest 7 passed,ruff checkruff format --diff 均通过,py_compile 通过;本地需 --noconftest 绕开顶层 conftest 的 GPU 依赖,CI 环境无此问题。CI 由 gty111 触发 Buildkite CI #84756。无配置、schema 或部署配套改动。
文件 模块 状态 重要度
examples/disaggregated/disaggregated_encoder/disagg_epd_proxy.py EPD 代理 modified 6.5
tests/v1/ec_connector/unit/test_epd_proxy_round_robin.py EPD 代理 added 7.03

关键符号

encoder_rr_assignment fanout_encoder_primer _load_proxy_module _drive

关键源码片段

examples/disaggregated/disaggregated_encoder/disagg_epd_proxy.py core-logic

修复的核心文件:原轮询生成器每次请求都从 e_urls[0] 开始,导致单 item 请求全部打向第一个 encoder;改为跨请求持久游标 + asyncio.Lock 保护,并抽出可单测的纯函数 encoder_rr_assignment。

# 模块级状态:游标跨请求保持,避免每次请求都从 e_urls[0] 重新开始
encoder_rr_idx = 0
encoder_rr_lock = asyncio.Lock()
​
​
def encoder_rr_assignment(
    e_urls: list[str], start: int, count: int
) -> tuple[list[str], int]:
    """把 count 个多模态 item 从游标 start 开始分配到 encoder URL。    返回每个 item 对应的 URL 列表,以及下一次调用应使用的游标值:
    连续请求之间的分配保持衔接(contiguous),而不是每次都从
    e_urls[0] 重启,单 item 请求因此能轮流命中所有 encoder。
    """
    urls = [e_urls[(start + i) % len(e_urls)] for i in range(count)]
    next_start = (start + count) % len(e_urls)
    return urls, next_start
​
​
async def fanout_encoder_primer(orig_request, e_urls, req_id):
    # ... 前置逻辑:解析 mm_items,初始化 tasks / item_uuids / item_meta ...
    # 关键修复点:在锁内一次性完成本次请求的 URL 分配并推进全局游标,
    # 替换原来每次请求都从 0 开始、按 item 逐个生成的写法。
    global encoder_rr_idx
    async with encoder_rr_lock:
        url_cycle, encoder_rr_idx = encoder_rr_assignment(
            e_urls, encoder_rr_idx, len(mm_items)
        )
    # ... 后续逻辑:为每个 item 构造 encoder 子请求并并发发送 ...
tests/v1/ec_connector/unit/test_epd_proxy_round_robin.py test-coverage

新增回归测试:用 importlib 加载 examples 下真实代理模块,驱动 encoder_rr_assignment 验证跨请求轮询,为 bug 场景(单 item 请求打满 e_urls[0])提供精确覆盖。

# 从 examples/ 真实路径加载代理模块,测试中不重写任何分配逻辑,
# 未来若改动 encoder_rr_assignment,测试会立即验证真实行为。
def _load_proxy_module():
    path = Path(__file__).parents[4] / PROXY_REL
    spec = importlib.util.spec_from_file_location("disagg_epd_proxy_under_test", path)
    assert spec is not None and spec.loader is not None
    module = importlib.util.module_from_spec(spec)
    spec.loader.exec_module(module)
    return module
​
​
@pytest.fixture(scope="module")
def assign():
    return _load_proxy_module().encoder_rr_assignment
​
​
def _drive(assign, e_urls, counts):
    """逐个请求推进游标,模拟 fanout_encoder_primer 每次调用的效果。"""
    cursor = 0
    all_urls = []
    for count in counts:
        urls, cursor = assign(e_urls, cursor, count)
        assert len(urls) == count
        all_urls.append(urls)
    return all_urls
​
​
@pytest.mark.parametrize("n_urls", [1, 2, 3, 5])
def test_full_url_space_is_covered_uniformly(assign, n_urls):
    # 多轮单 item 请求后,所有 encoder URL 都应被命中且次数完全相等,
    # 这正是原 bug(永远打向 e_urls[0])会失败的地方。
    e_urls = [f"E{i}" for i in range(n_urls)]
    urls = _drive(assign, e_urls, counts=[1] * (n_urls * 3))
    hits = Counter(u for req in urls for u in req)
    assert set(hits) == set(e_urls)
    assert max(hits.values()) == min(hits.values())

评论区精华

全局轮询游标 vs 随机选择的修复方向 设计

Issue #52484 给出两种修复方向:维护全局轮询游标("because the proxy is async, this should probably be protected with an asyncio.Lock or replaced with another concurrency-safe strategy"),或按 item 随机选择("avoids hot-spotting e_urls[0], although it is not deterministic round-robin")。PR 选择了确定性游标方案,并把游标数学抽成纯函数 encoder_rr_assignment 以降低并发代码的测试难度。

结论:采用全局游标 + asyncio.Lock 方案,与 issue 建议的第一种方案一致;review 中无异议,gty111 与 Isotr0py 均批准。 · 已解决

测试直接加载 examples 真实模块的模式 测试

PR body 说明参照 tests/v1/kv_connector/unit/test_moriio_proxy_routing.py 的先例,用 importlib 加载真实 disagg_epd_proxy.py,避免在测试中重写路由逻辑;本地用 --noconftest 绕开顶层 conftest 的 GPU 依赖,并需预装 aiohttp / fastapi / uvicorn 等第三方包。

结论:该模式被接受;代价是测试与 examples 路径及第三方依赖耦合,若目录调整或依赖缺失会导致测试收集失败。 · 已解决

fork PR 的自动化 review 状态 other

claude[bot] 提示 "This pull request is from a fork — automated review is disabled. A repository maintainer can comment @claude review to run a one-time review." 随后 gty111 直接批准并留言,Isotr0py 亦批准。

结论:人工 review 流程完成,无需触发自动 review;无遗留讨论项。 · 已解决

风险与影响

  • 并发模型假设:锁只保护单事件循环内的游标读写;若代理未来改为多进程 / 多 worker 部署,模块级 encoder_rr_idx 会在进程间失步,需要换成分布式或至少进程内共享状态。
  • 缺少端到端验证:PR body 明确 "Not run end-to-end against live encoder/decode/prefill clusters",锁在真实并发流量下的行为未经实测,现有单测只覆盖分配逻辑、未覆盖 HTTP 层的并发交错。
  • 测试耦合 examples 路径与第三方依赖:测试通过 importlib 加载的脚本顶层导入 aiohttp / fastapi / uvicorn / pybase64,若 examples 目录结构调整或环境缺依赖,测试收集即失败;本地必须 --noconftest 才能运行。
  • 回归面:变更仅限示例脚本,不触及 vLLM 核心推理、调度与 KV 缓存路径,核心回归风险极低。
  • 用户影响:使用 EPD 解耦示例(多模态 encoder 集群)的用户直接受益,单 item 请求不再把流量全部压到第一个 encoder,各 encoder 实例负载趋于均匀,缓解热点与实例闲置。
  • 系统影响:无核心系统变更,不改变 API 与协议;示例脚本行为变化仅体现在 encoder 选择顺序上。
  • 团队影响:测试确立了一种可复用模式——用 importlib 从 examples/ 加载真实模块做单测,避免测试复制生产逻辑;测试放置于 tests/v1/ec_connector/unit/,与近期 kv-connector 测试布局一致。
并发保护依赖单事件循环假设 缺少端到端集群验证 测试耦合 examples 路径与第三方依赖

关联 Issue

#52484 [Bug]: disagg_epd_proxy encoder fan-out restarts round-robin per request, overloading first encoder

完整报告

参与讨论