并行执行并等待全部结果汇合

并行查询再汇合

本章目标

同时执行两个独立查询,并在两者完成后生成汇总结果。

完整代码:examples/parallel_join.py

builder.add_edge(START, "fanout")
builder.add_edge("fanout", "query_kb")
builder.add_edge("fanout", "query_order")
builder.add_edge(["query_kb", "query_order"], "merge")
builder.add_edge("merge", END)

列表形式的 join edge 明确表达:merge 需要等待两个上游节点完成。两个查询分别写 kb_resultorder_result,因此不需要 Reducer。

两路查询一起跑

如果多个并行节点都写 findings,应将其声明为:

汇合前先定规则

findings: Annotated[list[str], operator.add]

运行验证:

python examples/parallel_join.py

预期:

退款政策允许处理;订单尚未发货

生产查询节点应写成 async def。同步 HTTP 或数据库调用会阻塞事件循环,抵消并行带来的延迟收益。

用时间测量证明并行真的发生

不要仅凭图形上有两条边就断言延迟下降。给两个测试节点各加入约 200 毫秒异步等待,测量总耗时应接近较慢分支,而不是两者之和:

import asyncio
from time import perf_counter


async def query_kb(state):
    await asyncio.sleep(0.2)
    return {"kb_result": "policy"}


async def query_order(state):
    await asyncio.sleep(0.2)
    return {"order_result": "paid"}


started = perf_counter()
result = await graph.ainvoke({"question": "退款"})
elapsed = perf_counter() - started
assert elapsed < 0.35

持续集成环境存在调度抖动,性能断言要留余量;更稳妥的做法是用事件或屏障验证两个节点都在 join 前完成。

并行分支失败策略

并行不等于忽略失败。为每个外部查询预先选择策略:

策略适用场景结果
fail-fast缺少任一结果都不能决策终止本次运行并返回可重试错误
有界技术重试超时、限流等短暂故障节点内部按策略重试
降级为部分结果辅助数据缺失仍能安全回答State 明确记录 degraded_sources
转人工缺失结果会影响资金或权限停止自动副作用并创建人工任务

不要把异常转换成空字符串后继续。"" 无法区分“服务正常返回空结果”和“服务调用失败”。可以返回结构化状态:

return {
    "order_result": None,
    "order_error": {"code": "TIMEOUT", "retryable": True},
}

取消与超时

为 HTTP、数据库和模型调用设置比整次运行预算更短的节点超时。客户端断开时是否取消后台运行,应由业务语义决定:只读问答可以取消,已经创建审批单或开始退款的流程通常必须转为后台运行并允许稍后查询状态。

并行度必须有上限

图内只有两条并行边,不代表系统最多只有两个并发调用。二百个运行同时到达时,会形成四百个下游请求。生产部署必须让 Worker 并发数、HTTP 连接池、数据库连接池和下游限流额度保持一致,并在入口设置队列上限。对动态 fan-out,还要限制单次运行可展开的任务数量;超过预算时分批处理或转异步队列,不能一次创建几千个分支。

容量测试应记录三个维度:单次运行的关键路径延迟、全局在途节点数、下游连接池等待时间。若连接池等待已经占据大部分延迟,继续增加应用实例只会把压力转移到下游。此时应降低并发、批量查询或增加缓存,而不是让 RetryPolicy 放大流量。

本章验收

  • 用测试证明 join 节点在两个分支完成后才执行。
  • 用计时或事件证明两个 I/O 分支不是串行运行。
  • 为每个分支记录超时、重试、降级或转人工策略。
  • 能区分“空业务结果”和“查询失败”。
  • 压测时并行请求数不会超过下游连接池和限流预算。