乐于分享
好东西不私藏

LangChain源码解析27:RunnableParallel如何同时跑多条链

LangChain源码解析27:RunnableParallel如何同时跑多条链

一段用户输入,既要生成摘要,又要提取关键词、判断情绪,还要保留原文。如果把这四件互不依赖的事全部串成 RunnableSequence,后一条链只能等前一条链结束,最终延迟接近所有分支耗时之和。

RunnableParallel 解决的是另一类组合问题:把同一个输入同时交给多条 Runnable,等待它们各自完成,再按分支名汇合成一个字典。

它的调用方式看起来只是一个普通 Python dict,但源码内部真正处理了六件事:输入扇出、Runnable 自动转换、同步线程池、异步任务汇合、callback 运行树,以及流式 chunk 的交错输出。

RunnableSequence 描述“先做什么、再做什么”,RunnableParallel 描述“同一份输入需要同时得到哪些彼此独立的结果”。

图 1:RunnableParallel 的输入扇出、子运行与结果汇合

一、Sequence 与 Parallel 不是两种速度,而是两种依赖关系

先看最小区别:

组合原语
数据怎样流动
总耗时大致取决于
输出形态
RunnableSequence
上一步输出成为下一步输入
各步骤耗时之和
最后一步的输出
RunnableParallel
所有分支接收同一个输入
最慢分支的耗时
{分支名: 分支结果}

假设三个分支分别耗时 1 秒、2 秒和 4 秒:

Sequence: 1 + 2 + 4 ≈ 7 秒
Parallel: max(1, 2, 4) ≈ 4 秒

这个估算只在分支真正独立、下游资源允许并发时成立。如果三个分支争用同一份模型配额、连接池或数据库锁,墙钟时间不会简单地缩短到最慢分支。

因此,选择 Parallel 的第一道判断不是“想不想更快”,而是:这些计算之间有没有数据依赖?

二、为什么管道里的普通 dict 会自动变成 Parallel

最常见的写法并不需要显式构造 RunnableParallel

from langchain_core.runnables import RunnableLambda

normalize = RunnableLambda(lambda text: text.strip())

chain = normalize | {
    "upper": lambda text: text.upper(),
    "length": lambda text: len(text),
}

result = chain.invoke("  langchain  ")
# {"upper": "LANGCHAIN", "length": 9}

| 运算符会调用 coerce_to_runnable()。它按类型完成四种转换:

已经是 Runnable     -> 原样返回
生成器函数           -> RunnableGenerator
普通/异步 callable   -> RunnableLambda
dict                 -> RunnableParallel

所以这里发生了两层转换:外层 dict 先变成 RunnableParallel,字典里的两个 lambda 再分别变成 RunnableLambda

嵌套字典也会递归形成嵌套 Parallel:

chain = normalize | {
    "text": lambda value: value,
    "stats": {
        "length": len,
        "words": lambda value: len(value.split()),
    },
}

输出结构会保持同样的 key 层级。这里的 key 不只是临时变量名,而是输出 schema、tracing 子运行和下游取值共同依赖的公开契约。

源码还保留了 RunnableMap = RunnableParallel 这个别名。两者没有不同的调度语义;新代码看到 RunnableMap 时,可以直接按 Parallel 理解。

三、构造函数做的核心工作是“收集并标准化分支”

RunnableParallel 同时支持位置字典和关键字参数:

parallel = RunnableParallel(
    {
        "summary": summary_chain,
        "keywords": keyword_chain,
    },
    sentiment=sentiment_chain,
)

构造逻辑可以压缩成下面几行:

merged = {**steps__} if steps__ is not None else {}
merged.update(kwargs)

steps = {
    key: coerce_to_runnable(value)
    for key, value in merged.items()
}

这里有三个容易忽略的细节:

  1. 传入映射会被浅拷贝,外部随后增删原字典不会直接改写这份分支表;
  2. 同名关键字分支会覆盖位置字典中的同名 key;
  3. 每个 value 都必须能被转换成 Runnable,否则构造阶段直接抛出 TypeError

它没有替分支做业务层类型适配。多个分支可以返回完全不同的 Python 类型,但下游必须清楚每个 key 对应什么结果。

四、“同一个输入扇出”不等于“复制多份输入”

同步 invoke() 提交每个分支时,传进去的是同一个 input 对象引用:

                     -> summary(input)
input(同一对象)     -> keywords(input)
                     -> sentiment(input)

框架不会对列表、字典、消息对象或自定义状态做深拷贝。这通常是正确选择:大文档和消息历史如果按分支复制,会产生明显的内存和序列化成本。

但它也给调用方留下一个重要约束:并行分支应该把输入视为只读数据。

如果两个同步分支同时修改同一个可变字典,最终结果会受到线程调度影响;异步分支即使运行在同一个事件循环里,也可能在 await 之间交错修改状态。需要写入时,应该让每个分支创建自己的新对象,或者在进入 Parallel 前显式复制必要字段。

五、同步 invoke() 怎样把分支提交到线程池

同步执行路径可以拆成六步:

1. ensure_config()
2. 创建 RunnableParallel 根 callback run
3. 浅拷贝当前 steps__,固定本轮分支快照
4. 为每个 key 创建 map:key:<key> 子 callback
5. 把所有 step.invoke() 提交给线程池
6. 汇合结果,结束根 run

核心结构接近下面这样:

steps = dict(self.steps__)

with get_executor_for_config(config) as executor:
    futures = [
        executor.submit(invoke_step, step, input, config, key)
        for key, step in steps.items()
    ]
    output = {
        key: future.result()
        for key, future in zip(steps, futures)
    }

get_executor_for_config() 创建的不是最朴素的 ThreadPoolExecutor,而是会复制 contextvars 的 ContextThreadPoolExecutor。这保证 tracing、配置上下文等信息能够跟着任务进入子线程。

执行开始前再做一次 dict(self.steps__),是为了固定这轮调用看到的分支集合。即使另一个调用方在运行期间修改 steps__,本轮也不会一半使用旧映射、一半使用新映射。

六、同步路径里的 max_concurrency 到底限制什么

线程池构造时会读取:

ContextThreadPoolExecutor(
    max_workers=config.get("max_concurrency")
)

因此:

parallel.invoke(
    input_data,
    config={"max_concurrency": 2},
)

当 Parallel 有 8 个同步分支时,8 个任务都会被提交,但同一时刻最多只有 2 个线程执行分支,其余任务留在线程池队列中。

如果不配置,线程池使用 Python 运行时的默认 worker 数量。这个默认值是执行环境细节,不应该被当成稳定的业务并发预算。

还要注意:这里限制的是这一次同步 Parallel 调用内部的线程 worker 数,不是整个进程、整个模型实例或所有请求共享的全局上限。

七、为什么返回字典按声明顺序,而不是完成顺序

虽然各分支同时执行,但普通 invoke() 不是 as-completed API。源码按 steps 与 futures 的声明顺序逐个调用 future.result(),再构造输出字典。

假设声明顺序是:

slow -> fast -> medium

即使 fast 最先完成,最终字典仍然是:

{
    "slow": slow_result,
    "fast": fast_result,
    "medium": medium_result,
}

等待顺序不会让后面的任务停止并发执行,但它会影响异常被主调用观察到的时机:后置分支已经失败时,主线程仍可能先等待排在前面的慢分支。

普通 invoke() 的契约是“一次返回完整字典”。如果消费者需要谁先完成就先处理谁,应该使用流式接口,而不是依赖字典插入顺序猜测完成顺序。

八、异步 ainvoke() 用的是 asyncio.gather()

异步路径没有建立线程池,而是为所有分支构造协程,再交给 asyncio.gather()

results = await asyncio.gather(
    *(
        ainvoke_step(step, input, config, key)
        for key, step in steps.items()
    )
)

output = dict(zip(steps, results))

gather() 返回的结果列表与传入协程顺序一致,所以异步返回字典同样保持分支声明顺序。

这里最需要守住一个源码边界:RunnableParallel.ainvoke() 本身没有用 gather_with_concurrency() 包住这些分支。

也就是说,把 max_concurrency=2 传给一次顶层 parallel.ainvoke(),不会直接把 8 个异步分支截成每次只运行 2 个。config 仍会传给子 Runnable,并可能被子层自己的 batch、线程池或调度器读取,但这一层的 8 个协程会一起交给 asyncio.gather()

同一个配置字段只有在某层源码主动读取时才生效。RunnableConfig 是配置载体,不是自动笼罩整个进程的信号量。

如果异步分支会直接访问昂贵的外部资源,应该在资源入口设置共享信号量、连接池、RateLimiter 或 Provider 侧配额,而不是只依赖顶层 Parallel 的 config。

九、map:key:<key> 怎样把并行分支放进同一棵运行树

每个分支调用前,Parallel 都会 patch 一份 child config:

callbacks=run_manager.get_child(f"map:key:{key}")

因此 tracing 看到的不是几条互不相关的链,而是:

RunnableParallel<summary,keywords,sentiment>
  ├─ map:key:summary
  │   └─ summary_chain
  ├─ map:key:keywords
  │   └─ keyword_chain
  └─ map:key:sentiment
      └─ sentiment_chain

根 run 负责记录完整输入、最终输出和整体成功或失败;子 run 继承 callbacks、tags 与 metadata,同时通过 map:key:* 标记自己属于哪个结果字段。

这套命名让观测系统可以回答三个关键问题:

  1. 哪个分支最慢,决定了整体尾延迟;
  2. 哪个分支失败,导致根 Parallel 失败;
  3. 一次逻辑调用实际扇出了多少模型、检索或工具子调用。

十、一个分支失败,为什么整个 Parallel 不返回“半张答卷”

同步路径在 future.result() 处重新抛出分支异常,异步路径的 asyncio.gather() 也没有使用 return_exceptions=True。外层捕获异常后,会执行根 run 的 on_chain_error(),然后继续向上抛出。

所以默认语义是:

summary 成功
keywords 失败
sentiment 成功
        ↓
Parallel 整体失败,不返回部分结果字典

这不代表其他分支从未执行。异常被观察到之前,它们可能已经完成远程请求、写入日志或产生外部副作用;Parallel 也没有提供跨分支事务回滚。

如果业务允许部分成功,可以把容错边界放到每个分支内部:

parallel = RunnableParallel(
    summary=summary_chain.with_fallbacks([summary_default]),
    keywords=keyword_chain.with_fallbacks([keyword_default]),
)

另一种做法是让每个分支统一返回 {ok, value, error} 结果信封,再由下游决定哪些字段可缺失。不要用一个笼统的外层 try/except 假装拿到了框架没有返回的部分字典。

十一、输出 schema 为什么天然就是一张 key 表

get_output_schema() 会遍历所有分支,把每个 key 与该分支的 OutputType 组合成一个 Pydantic v2 模型:

fields = {
    key: (step.OutputType, ...)
    for key, step in self.steps__.items()
}

这意味着 Parallel 的结构不仅存在于运行时结果中,也可以被静态检查、图展示和工具系统读取。

输入 schema 的处理更谨慎。只要各分支都暴露对象型 JSON Schema,框架会尝试合并它们的字段;遇到有具体类型的 root model 时,则会回退到更合适的单一 schema,避免把根值模型错误拼成普通对象。

get_graph() 同样体现扇出与汇合:它建立一个统一输入节点、一个统一输出节点,把每条分支子图的入口连到输入,把出口连到输出。

图结构描述的是依赖关系,不等于实际线程或 task 数量。运行时并发仍由 invoke()ainvoke()transform() 和它们读取的 config 决定。

十二、流式 Parallel 返回的不是最终字典,而是单 key chunk

模型一旦流式输出,多个分支就不可能同时凑齐一张完整字典再返回,否则最快分支仍要等待最慢分支,streaming 会失去意义。

因此 stream() 的输出更像这样:

{"summary": chunk-1}
{"keywords": chunk-1}
{"summary": chunk-2}
{"sentiment": chunk-1}
{"keywords": chunk-2}

每个 chunk 通常只包含一个分支 key。跨分支顺序取决于谁先产出下一块,因此不稳定;同一分支内部仍保持自己的原始顺序。

消费者应该按 key 更新各自的 UI 区域或聚合缓冲区,而不是假设第一个 chunk 一定来自声明在最前面的分支。

十三、safetee() 为什么是同步流式扇出的起点

transform() 面对的不是一个静态 input,而是一条输入迭代器。多个分支如果直接对同一个迭代器调用 next(),会把上游元素互相抢走:第一个输入被 A 取走,第二个输入可能被 B 取走,任何分支都看不到完整输入流。

源码先执行:

input_copies = list(
    safetee(inputs, len(steps), lock=threading.Lock())
)

safetee() 为每个分支建立一个逻辑副本,并用线程锁协调对共享上游的读取。某个分支先向前推进时,新元素会被暂存在其他副本的缓冲区里,保证所有分支最终看到相同顺序的输入。

这里复制的是迭代进度和缓冲,不是把每个输入对象做深拷贝。可变对象仍然应该按只读数据处理;分支速度差距过大时,落后分支对应的缓冲还可能持续增长。

十四、FIRST_COMPLETED 怎样让最快的 chunk 先出来

每个分支生成器创建后,Parallel 会先为它提交一次 next()

futures = {
    executor.submit(next, generator): (step_name, generator)
    for step_name, generator in named_generators
}

随后进入循环:

等待任意 future 完成(FIRST_COMPLETED)
  -> 取出这个分支的 chunk
  -> 包成 {step_name: chunk}
  -> 立即向外 yield
  -> 再为同一个生成器提交下一次 next()

每个分支同一时刻只保留一次“取下一块”的请求,因此分支内部不会乱序。不同分支则互相竞争,谁先准备好,谁的 chunk 就先交给消费者。

如果多个 future 在同一轮等待中同时完成,它们之间的遍历顺序也不应被视为稳定契约。

十五、AddableDict 怎样把交错 chunk 重新拼起来

源码输出的单 key chunk 实际是 AddableDict

chunk = AddableDict({step_name: future.result()})

它的加法规则是:

  • 新 key 直接加入;
  • 已有 key 的两个值如果支持 +,就相加;
  • 值不支持相加时,用较新的值替换;
  • None
     不会无意义地覆盖已有结果。

因此字符串、AIMessageChunk 等可加对象可以自然聚合:

combined = None

for chunk in parallel.stream(input_data):
    combined = chunk if combined is None else combined + chunk

这也是 LCEL 流式组合能够继续向下游传递部分结构的基础:外层按 key 交错交付,聚合时再把同一 key 的连续内容相加。

十六、异步流式路径是 atee() 加 asyncio.wait()

astream() 保持相同设计,只是换成异步原语:

safetee + threading.Lock
    ↓
atee + asyncio.Lock

executor.submit(next)
    ↓
asyncio.create_task(anext)

wait(..., FIRST_COMPLETED)
    ↓
asyncio.wait(..., FIRST_COMPLETED)

每个异步分支先创建一个读取下一块的 task。任意 task 完成后,框架输出单 key AddableDict,再为这个分支创建新的 task。

与 ainvoke() 一样,这一层会为所有分支各创建一个 task,没有在 _atransform() 内用 max_concurrency 建立统一 semaphore。外部 I/O 容量仍然要在真正的资源边界控制。

图 2:batch 输入并发、Parallel 分支扇出与流式交错

十七、batch() 套 Parallel,为什么并发可能乘起来

Parallel 解决的是一个输入内部的分支并发,batch 解决的是多个输入之间的并发。两层组合后,不能再只看一个数字。

假设:

输入数量 = 100
Parallel 分支数 = 8
max_concurrency = 5

异步 abatch() 会用 max_concurrency=5 控制同时推进的顶层输入;但每个 RunnableParallel.ainvoke() 又会把 8 个分支一起交给 asyncio.gather()。于是理论上可能同时存在:

5 个输入 × 8 个异步分支 = 40 个子调用

同步 batch() 也有两层调度:外层线程池控制同时处理多少输入,每个同步 parallel.invoke() 内部又会创建自己的分支线程池。相同 config 可能同时被两层读取,形成嵌套线程池与更多排队任务。

如果 Parallel 里面还有嵌套 Parallel、模型自身 retry 或工具内部并发,真实 fan-out 会继续扩大。评估容量时应该画出完整调用树,而不是只在最外层看到 max_concurrency=5 就认为系统最多发出 5 条下游请求。

十八、上一层的 RateLimiter 为什么仍然重要

max_concurrency 控制同时推进多少工作,RateLimiter 控制请求以多快的节奏开始。Parallel 扇出后,两者的边界更清楚:

batch 输入槽位
  -> Parallel 分支槽位或异步 task
  -> 模型 cache
  -> 模型 RateLimiter
  -> Provider 请求

如果 8 个分支共享同一个模型和同一个 limiter,它们会竞争同一只请求桶;如果每个分支各自创建 limiter,参数即使完全相同,也不会自动共享额度。

Parallel 还会放大 retry 的最坏请求数。8 个模型分支,每个最多 3 次 attempt,一次逻辑输入理论上可能产生 24 次 Provider 尝试。可靠性、并发与限流必须按同一棵运行树估算。

十九、哪些场景适合使用 RunnableParallel

1. 一份内容做多维提取

同一篇文档同时生成摘要、关键词、分类、风险标签,各字段之间没有依赖。

2. RAG 中保留问题并获取上下文

一个分支把问题原样传递,另一个分支执行 retriever,随后统一交给 prompt:

{
    "question": RunnablePassthrough(),
    "context": retriever,
} | prompt | model

3. 多模型或多提示词对照

同一输入交给不同模型、温度或 prompt,输出按 key 保留,便于评测和人工比较。

4. 同时计算原值与派生值

保留原始对象的同时计算统计信息、embedding 或路由信号,让后续步骤获得一张结构化结果表。

5. 多个独立 I/O 查询

同时查询互不依赖的数据源,但前提是连接池、超时、限流和失败语义已经明确。

二十、哪些写法看似并行,实际是在制造问题

1. 把有依赖的步骤硬塞进不同分支

B 必须使用 A 的结果时,它们应该组成 Sequence,或者先 Parallel 再做统一汇合;不能靠完成时间碰运气。

2. 多分支修改同一个可变输入

Parallel 不做深拷贝。共享写入会制造数据竞争和不可复现结果。

3. 把有副作用的分支当成一个事务

一个分支失败不会自动撤销其他分支已经发送的邮件、写入的数据库或调用的工具。

4. 认为普通 invoke() 会先返回最快结果

它只在所有分支成功后返回完整字典。真正的完成顺序交付在 stream() / astream() 中。

5. 认为一个 max_concurrency 能覆盖所有层

同步线程池、异步 gather、batch、嵌套 Runnable 和 Provider SDK 各自拥有调度边界。必须确认具体哪一层读取了配置。

6. 无限制地动态生成分支 key

Parallel 在开始调用时会对全部分支扇出。来自用户输入的超大动态分支集合,可能瞬间创建大量 future、task、callback 和外部请求。

二十一、一份可落地的生产检查清单

  1. 每个分支是否真的互不依赖?
  2. 分支是否把输入视为只读对象?
  3. 输出 key 是否稳定,并已成为下游明确契约?
  4. 最慢分支的 P95/P99 是否决定了不可接受的整体尾延迟?
  5. 同步调用的线程池 worker 是否与连接池容量匹配?
  6. 异步 ainvoke() 的所有分支同时调度是否可承受?
  7. batch
     输入并发与 Parallel 分支数相乘后,最大子调用数是多少?
  8. 嵌套 Parallel、retry 和 fallback 是否继续放大 fan-out?
  9. 多个模型分支使用共享 limiter,还是各自独立配额?
  10. 某个分支失败时,业务需要整体失败还是部分成功?
  11. 有副作用的分支是否具备幂等键、补偿或事务边界?
  12. stream 消费者是否按 key 聚合,而不是依赖跨分支顺序?
  13. 快慢分支差距过大时,tee 缓冲是否可能增长?
  14. tracing 中是否能按 map:key:* 观察每条分支的耗时、错误和请求数?
  15. 取消、超时和服务关闭时,仍在执行的分支怎样清理?

二十二、Parallel 的本质是一张“同输入、多结果”的运行图

RunnableParallel 的表面 API 很轻:一个字典、几个 key、几条 Runnable。真正有价值的是它把并行组合固定成了一套可观察、可建模的运行语义:

同一个输入
  -> 按 key 扇出
  -> 每个分支拥有独立 child run
  -> 同步用线程池,异步用 gather
  -> 普通调用按 key 汇合完整字典
  -> 流式调用按 FIRST_COMPLETED 交错输出
  -> AddableDict 再把 chunk 聚合回来

它不会替业务判断分支是否独立,也不会提供全局并发预算、事务回滚或部分成功协议。框架只把“如何扇出、如何追踪、如何汇合”做成统一原语,资源边界与失败语义仍然属于应用设计。

最终需要记住的是:

Parallel 的性能收益来自独立分支重叠等待时间;它的工程风险也来自同一个事实:一次输入会在很短时间内变成多条真实工作。

系列链接

第 1 篇:LangChain源码解析01:先看懂Agent工程骨架

第 2 篇:LangChain源码解析02:Runnable把一切串起来

第 3 篇:LangChain源码解析03:RunnableConfig如何追踪到底

第 4 篇:LangChain源码解析04:Message不只是字符串

第 5 篇:LangChain源码解析05:Tool如何从函数变成契约

第 6 篇:LangChain源码解析06:Prompt和Parser守住两端

第 7 篇:LangChain源码解析07:BaseChatModel如何统一模型调用

第 8 篇:LangChain源码解析08:init_chat_model如何动态切换模型

第 9 篇:LangChain源码解析09:create_agent如何编译Agent运行图

第 10 篇:LangChain源码解析10:Agent条件边如何决定下一步

第 11 篇:LangChain源码解析11:middleware如何接管Agent执行链

第 12 篇:LangChain源码解析12:结构化输出如何在两种策略间切换

第 13 篇:LangChain源码解析13:ToolRuntime如何把上下文注入工具

第 14 篇:LangChain源码解析14:stream_events如何汇合Agent多路事件

第 15 篇:LangChain源码解析15:Agent为什么最终变成StateGraph

第 16 篇:LangChain源码解析16:AgentExecutor如何驱动经典Agent循环

第 17 篇:LangChain源码解析17:ChatOpenAI如何统一两套API

第 18 篇:LangChain源码解析18:ChatAnthropic如何编排思考与工具

第 19 篇:LangChain源码解析19:一套测试如何验收所有模型

第 20 篇:LangChain源码解析20:长文档如何切成可检索的块

第 21 篇:LangChain源码解析21:模型能力如何驱动运行时决策

第 22 篇:LangChain源码解析22:一次Agent调用如何变成运行树

第 23 篇:LangChain源码解析23:新Provider如何通过框架验收

第 24 篇:LangChain源码解析24:LangChain到底设计对了什么

第 25 篇:LangChain源码解析25:with_retry与with_fallbacks如何接住失败

第 26 篇:LangChain源码解析26:RateLimiter如何决定请求能不能发

源码参考: GitHub: https://github.com/langchain-ai/langchain

当一张并行结果表还要继续保留原输入、追加派生字段并逐层整形时,怎样才能避免到处手写字典拆装?