一段用户输入,既要生成摘要,又要提取关键词、判断情绪,还要保留原文。如果把这四件互不依赖的事全部串成 RunnableSequence,后一条链只能等前一条链结束,最终延迟接近所有分支耗时之和。
RunnableParallel 解决的是另一类组合问题:把同一个输入同时交给多条 Runnable,等待它们各自完成,再按分支名汇合成一个字典。
它的调用方式看起来只是一个普通 Python dict,但源码内部真正处理了六件事:输入扇出、Runnable 自动转换、同步线程池、异步任务汇合、callback 运行树,以及流式 chunk 的交错输出。
RunnableSequence描述“先做什么、再做什么”,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()
}
这里有三个容易忽略的细节:
传入映射会被浅拷贝,外部随后增删原字典不会直接改写这份分支表; 同名关键字分支会覆盖位置字典中的同名 key; 每个 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:* 标记自己属于哪个结果字段。
这套命名让观测系统可以回答三个关键问题:
哪个分支最慢,决定了整体尾延迟; 哪个分支失败,导致根 Parallel 失败; 一次逻辑调用实际扇出了多少模型、检索或工具子调用。
十、一个分支失败,为什么整个 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 容量仍然要在真正的资源边界控制。

十七、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 和外部请求。
二十一、一份可落地的生产检查清单
每个分支是否真的互不依赖? 分支是否把输入视为只读对象? 输出 key 是否稳定,并已成为下游明确契约? 最慢分支的 P95/P99 是否决定了不可接受的整体尾延迟? 同步调用的线程池 worker 是否与连接池容量匹配? 异步 ainvoke()的所有分支同时调度是否可承受?batch输入并发与 Parallel 分支数相乘后,最大子调用数是多少? 嵌套 Parallel、retry 和 fallback 是否继续放大 fan-out? 多个模型分支使用共享 limiter,还是各自独立配额? 某个分支失败时,业务需要整体失败还是部分成功? 有副作用的分支是否具备幂等键、补偿或事务边界? stream 消费者是否按 key 聚合,而不是依赖跨分支顺序? 快慢分支差距过大时,tee 缓冲是否可能增长? tracing 中是否能按 map:key:*观察每条分支的耗时、错误和请求数?取消、超时和服务关闭时,仍在执行的分支怎样清理?
二十二、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
当一张并行结果表还要继续保留原输入、追加派生字段并逐层整形时,怎样才能避免到处手写字典拆装?
夜雨聆风