# LangChain 两个 Runnable 核心类:RunnableParallel 与 RunnablePassthrough
学 LCEL 时,最容易卡住的地方不是 prompt | model | parser 这种顺序链路,而是数据一旦变成 dict,输入到底会怎么流动。
比如 RAG 链路里经常看到这种写法:
rag_chain = (
{
"context": retriever | format_docs,
"question": RunnablePassthrough(),
}
| prompt
| model
| parser
)
2
3
4
5
6
7
8
9
第一眼看会有点绕:retriever 怎么拿到问题?RunnablePassthrough() 又为什么能把问题放到 question 字段里?这个 dict 到底是在创建字典,还是在并行执行?
答案是:在 LCEL 中,dict 会被自动转换成 RunnableParallel;RunnablePassthrough 则负责把原始输入原样传下去,或者在原始 dict 上追加字段。理解这两个类,LCEL 的数据流就会清楚很多。
这篇重点讲两个核心类:
RunnableParallel:把同一份输入分发给多个分支,并把分支结果合并成 dict。RunnablePassthrough:把输入原样传递,常用于保留原始问题,或通过assign给输入追加字段。
它们看起来像语法糖,但在生产链路里很关键:RAG、并行打标、多模型对比、上下文组装、用户画像注入、trace 可读性,都会用到这两个能力。
# 先看一个普通顺序链
普通 LCEL 链路是线性的:
from langchain_core.output_parsers import StrOutputParser
from langchain_core.prompts import ChatPromptTemplate
from langchain_openai import ChatOpenAI
prompt = ChatPromptTemplate.from_messages([
("system", "你是一个简洁、准确的技术助手。"),
("human", "{question}"),
])
model = ChatOpenAI(model="gpt-4o-mini", temperature=0)
parser = StrOutputParser()
chain = prompt | model | parser
answer = chain.invoke({
"question": "什么是 RunnableParallel?"
})
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
这条链路的数据方向很直观:
dict -> prompt -> model -> parser -> str
但真实业务经常需要分叉:
- 同一个用户问题,一边去检索知识库,一边保留问题本身。
- 同一段文本,同时生成摘要、标题、关键词。
- 同一份输入,同时跑规则判断、模型判断、风险判断。
- 同一个问题,同时请求多个 retriever,再合并上下文。
这时就需要 RunnableParallel。
# RunnableParallel:同一输入,多路执行
RunnableParallel 的核心语义是:
同一份输入 -> 多个 Runnable 分支 -> dict 结果
最直接的例子:
from langchain_core.runnables import RunnableParallel
chain = RunnableParallel({
"upper": lambda text: text.upper(),
"length": lambda text: len(text),
})
result = chain.invoke("langchain")
print(result)
2
3
4
5
6
7
8
9
10
11
输出结构类似:
{
"upper": "LANGCHAIN",
"length": 9,
}
2
3
4
两个分支拿到的是同一个输入 "langchain"。upper 分支负责转大写,length 分支负责计算长度,最终结果按 key 合成一个 dict。
在 LCEL 里,还可以直接用 dict 字面量表达同样的事:
chain = {
"upper": lambda text: text.upper(),
"length": lambda text: len(text),
}
result = chain.invoke("langchain")
2
3
4
5
6
当 dict 出现在 LCEL 链路中时,LangChain 会把它转换为并行 Runnable。这就是很多 RAG 示例里 dict 写法的来源。
# dict 不是普通 dict
在 LCEL 里,下面这段代码不是“提前算出一个 Python dict”:
{
"summary": summary_chain,
"keywords": keywords_chain,
}
2
3
4
它表示一组并行分支。真正执行发生在 invoke、batch、stream 时。
例如:
from langchain_core.output_parsers import StrOutputParser
from langchain_core.prompts import ChatPromptTemplate
from langchain_openai import ChatOpenAI
model = ChatOpenAI(model="gpt-4o-mini", temperature=0)
parser = StrOutputParser()
summary_chain = (
ChatPromptTemplate.from_template("请用 50 字总结:{content}")
| model
| parser
)
keywords_chain = (
ChatPromptTemplate.from_template("请提取 5 个关键词,用逗号分隔:{content}")
| model
| parser
)
analysis_chain = {
"summary": summary_chain,
"keywords": keywords_chain,
}
result = analysis_chain.invoke({
"content": "LCEL 可以把多个 Runnable 组合成清晰的数据流。"
})
print(result["summary"])
print(result["keywords"])
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
这里 summary_chain 和 keywords_chain 都收到同一份输入 dict。每个分支自己从 {content} 里取值。
生产中这类写法适合文档处理、内容审核、运营自动化、离线标注等场景。但要注意并发带来的成本和限流压力。两个分支就是两次模型调用,十个分支就是十次模型调用,不会因为写在一个 dict 里就免费。
# RunnableParallel 的显式写法
dict 写法很简洁,但复杂链路里我更建议关键位置显式使用 RunnableParallel。
from langchain_core.runnables import RunnableParallel
analysis_chain = RunnableParallel({
"summary": summary_chain,
"keywords": keywords_chain,
})
2
3
4
5
6
7
显式写法的好处是阅读成本更低,尤其适合:
- 团队里有人刚接触 LCEL。
- 并行分支很多。
- 每个分支都有独立业务含义。
- 需要给这段链路命名或加配置。
例如:
analysis_chain = RunnableParallel({
"summary": summary_chain,
"keywords": keywords_chain,
}).with_config({
"run_name": "document_parallel_analysis",
"tags": ["document", "parallel"],
})
2
3
4
5
6
7
在 LangSmith 或日志里,一个稳定的 run_name 会比匿名 dict 更容易排查。
# RunnablePassthrough:保留原始输入
RunnablePassthrough 的核心语义是:
输入什么,就输出什么
最小例子:
from langchain_core.runnables import RunnablePassthrough
passthrough = RunnablePassthrough()
result = passthrough.invoke("hello")
print(result)
2
3
4
5
6
7
8
输出就是:
"hello"
这看起来好像没用,但一旦进入并行分支,它就很有用。
例如:
from langchain_core.runnables import RunnableParallel, RunnablePassthrough
chain = RunnableParallel({
"original": RunnablePassthrough(),
"length": lambda text: len(text),
})
result = chain.invoke("langchain")
print(result)
2
3
4
5
6
7
8
9
10
11
输出:
{
"original": "langchain",
"length": 9,
}
2
3
4
它的作用不是改变输入,而是把原始输入保留下来,给后续步骤使用。
# RAG 里的经典用法
回到最常见的 RAG:
from langchain_core.output_parsers import StrOutputParser
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.runnables import RunnablePassthrough
from langchain_openai import ChatOpenAI
def format_docs(docs):
return "\n\n".join(doc.page_content for doc in docs)
prompt = ChatPromptTemplate.from_messages([
(
"system",
"""
你是企业知识库问答助手。
请只根据上下文回答问题。
如果上下文中没有答案,请明确说不知道。
上下文:
{context}
""",
),
("human", "{question}"),
])
model = ChatOpenAI(model="gpt-4o-mini", temperature=0)
rag_chain = (
{
"context": retriever | format_docs,
"question": RunnablePassthrough(),
}
| prompt
| model
| StrOutputParser()
)
answer = rag_chain.invoke("数据库迁移上线前要检查什么?")
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
执行过程可以这样理解:
输入问题
-> context 分支:问题 -> retriever -> docs -> format_docs -> context
-> question 分支:问题 -> RunnablePassthrough -> question
-> 合并成 {"context": "...", "question": "..."}
-> prompt -> model -> parser
2
3
4
5
retriever 能拿到问题,是因为 RunnableParallel 会把同一份输入传给每个分支。RunnablePassthrough() 能把问题放到 question 字段,是因为它原样返回输入。
这就是那段 RAG 链路的核心。
# 输入是 dict 时怎么用
上面的例子输入是字符串。如果输入本身就是 dict,RunnablePassthrough 会原样传递整个 dict。
from langchain_core.runnables import RunnableParallel, RunnablePassthrough
chain = RunnableParallel({
"payload": RunnablePassthrough(),
"question": lambda payload: payload["question"],
"tenant_id": lambda payload: payload["tenant_id"],
})
result = chain.invoke({
"question": "如何使用 LCEL?",
"tenant_id": "tenant_001",
})
2
3
4
5
6
7
8
9
10
11
12
13
输出类似:
{
"payload": {
"question": "如何使用 LCEL?",
"tenant_id": "tenant_001",
},
"question": "如何使用 LCEL?",
"tenant_id": "tenant_001",
}
2
3
4
5
6
7
8
生产里,输入 dict 通常会包含很多字段:用户 ID、租户 ID、问题、渠道、语言、权限、实验分组。不要把整个 dict 都塞进 Prompt。应该明确哪些字段给模型看,哪些字段只放到 config.metadata 做观测。
# assign:在原始 dict 上追加字段
RunnablePassthrough.assign 是非常实用的能力。它要求输入是 dict,然后在原始 dict 上追加新字段。
from langchain_core.runnables import RunnablePassthrough
def load_user_profile(payload: dict) -> str:
return f"用户 {payload['user_id']} 是企业版客户。"
chain = RunnablePassthrough.assign(
profile=load_user_profile,
)
result = chain.invoke({
"user_id": "u_1001",
"question": "我可以使用高级报表吗?",
})
print(result)
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
输出类似:
{
"user_id": "u_1001",
"question": "我可以使用高级报表吗?",
"profile": "用户 u_1001 是企业版客户。",
}
2
3
4
5
它适合做上下文补充:
- 根据
user_id加载用户画像。 - 根据
tenant_id加载租户配置。 - 根据
question检索知识库上下文。 - 根据权限计算可用工具列表。
- 根据输入内容计算风险等级。
然后再进入 Prompt:
prompt = ChatPromptTemplate.from_messages([
(
"system",
"""
你是产品客服。
请结合用户画像和知识库上下文回答。
用户画像:
{profile}
上下文:
{context}
""",
),
("human", "{question}"),
])
chain = (
RunnablePassthrough.assign(
profile=load_user_profile,
context=lambda payload: retriever.invoke(payload["question"]),
)
| prompt
| model
| StrOutputParser()
)
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
相比先手动构造一个大 dict,assign 的优势是链路意图更清楚:先保留原始输入,再追加运行时上下文。
# RunnableParallel 和 assign 的区别
两者都能生成 dict,但语义不一样。
RunnableParallel 是“分支结果组成一个新 dict”:
chain = RunnableParallel({
"question": lambda payload: payload["question"],
"context": lambda payload: retriever.invoke(payload["question"]),
})
2
3
4
输出只包含你声明的字段。
RunnablePassthrough.assign 是“保留原始 dict,再追加字段”:
chain = RunnablePassthrough.assign(
context=lambda payload: retriever.invoke(payload["question"]),
)
2
3
输出包含原始字段和新增字段。
选择规则很简单:
| 需求 | 推荐 |
|---|---|
| 只想构造 Prompt 需要的几个字段 | RunnableParallel |
| 想保留完整原始输入并追加字段 | RunnablePassthrough.assign |
| 输入是字符串,需要一边检索一边保留问题 | RunnableParallel + RunnablePassthrough() |
| 输入是 dict,需要增加 profile、context、risk | RunnablePassthrough.assign |
生产里我更偏向在“输入 dict 已经是业务对象”的情况下使用 assign,因为它不会丢失原始字段,后面做审计、trace 和错误处理更方便。但 Prompt 输入前仍然要控制字段,不要让不该被模型看到的字段混进去。
# 多检索器并行
RunnableParallel 很适合多检索器场景。
例如一个企业知识库可能有:
- 产品文档 retriever。
- 工单历史 retriever。
- FAQ retriever。
可以并行查:
from langchain_core.runnables import RunnableParallel, RunnablePassthrough
def format_docs(docs):
return "\n\n".join(doc.page_content for doc in docs)
retrieval_chain = RunnableParallel({
"product_docs": product_retriever | format_docs,
"ticket_docs": ticket_retriever | format_docs,
"faq_docs": faq_retriever | format_docs,
"question": RunnablePassthrough(),
})
prompt = ChatPromptTemplate.from_messages([
(
"system",
"""
你是企业知识库助手。
请综合不同来源回答,并优先使用产品文档中的正式说法。
产品文档:
{product_docs}
历史工单:
{ticket_docs}
FAQ:
{faq_docs}
""",
),
("human", "{question}"),
])
chain = retrieval_chain | prompt | model | StrOutputParser()
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
这个写法比把多个 retriever 手动写在 service 里更直观,也更容易在 trace 中看到每个分支的耗时。
但多检索器并行会带来几个问题:
- 上下文可能过长。
- 不同来源质量不一致。
- 检索结果可能互相冲突。
- 多个检索器同时失败时错误处理更复杂。
- 总延迟取决于最慢分支。
生产里要给每个来源设置 top k、token 上限、超时和降级策略。不能把所有检索结果无脑塞进 Prompt。
# 多模型并行对比
RunnableParallel 也可以用于模型对比。
from langchain_core.runnables import RunnableParallel
cheap_model_chain = prompt | cheap_model | parser
strong_model_chain = prompt | strong_model | parser
compare_chain = RunnableParallel({
"cheap": cheap_model_chain,
"strong": strong_model_chain,
})
result = compare_chain.invoke({
"question": "请解释数据库迁移的灰度策略。"
})
2
3
4
5
6
7
8
9
10
11
12
13
14
这适合:
- 离线评估。
- Prompt 调优。
- 灰度对比。
- 成本和质量分析。
不建议在普通线上请求里长期这么做。双模型并行意味着双倍成本,还可能触发更高的并发压力。线上对比更适合采样、异步评估或影子流量,而不是每次用户请求都同步跑两套模型。
# 风险判断与主链路并行
一个比较实用的模式是:回答生成和风险判断并行。
answer_chain = answer_prompt | model | StrOutputParser()
risk_chain = risk_prompt | model.with_structured_output(RiskResult)
chain = RunnableParallel({
"answer": answer_chain,
"risk": risk_chain,
})
result = chain.invoke({
"question": "帮我删除这个用户的所有数据。",
})
if result["risk"].requires_review:
return "这个操作需要人工确认。"
return result["answer"]
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
但高风险业务不能只靠模型判断。风险链路可以作为辅助信号,最终还要结合权限系统、操作类型、资源状态和审计策略。
# 和 RunnableLambda 搭配
很多时候需要在并行前后做一点数据处理,可以搭配 RunnableLambda。
from langchain_core.runnables import RunnableLambda, RunnableParallel
def normalize(payload: dict) -> dict:
return {
"question": payload["question"].strip(),
"locale": payload.get("locale", "zh-CN"),
"tenant_id": payload["tenant_id"],
}
def select_question(payload: dict) -> str:
return payload["question"]
chain = (
RunnableLambda(normalize)
| RunnableParallel({
"question": select_question,
"context": lambda payload: retriever.invoke(payload["question"]),
})
| prompt
| model
| StrOutputParser()
)
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
注意这里的 RunnableLambda 应该保持轻量、纯粹、可测试。不要把数据库写入、扣费、发消息这类副作用藏在 lambda 里。
# batch 时会发生什么
RunnableParallel 本身也是 Runnable,所以支持 batch:
payloads = [
{"content": "第一篇文章..."},
{"content": "第二篇文章..."},
{"content": "第三篇文章..."},
]
results = analysis_chain.batch(payloads)
2
3
4
5
6
7
每个输入都会进入并行分支。也就是说,如果 analysis_chain 有 3 个模型分支,batch 有 100 条输入,理论上就是 300 次模型调用。
生产批处理要特别注意:
- 并发上限。
- provider rate limit。
- 单条失败隔离。
- 结果幂等写入。
- 成本预算。
- 任务取消和恢复。
不要只因为 batch 写起来很短,就把几千条数据直接丢进去。大规模任务应该用队列、worker、限速器和断点续跑。
# stream 时要注意什么
并行链路也可以 stream,但输出事件会比线性链路复杂。
简单文本链路:
for chunk in answer_chain.stream({"question": "什么是 RunnablePassthrough?"}):
print(chunk, end="", flush=True)
2
并行链路:
for chunk in parallel_chain.stream({"content": "..." }):
print(chunk)
2
并行分支的 chunk 可能按 key 返回增量结果。前端如果要展示,需要知道每个事件属于哪个分支。
生产里建议:
- 普通聊天使用单一文本流。
- 多分支过程使用事件流。
- 最终结构化结果完整返回。
- 前端明确处理 start、delta、end、error。
- 后端不要把内部 trace 原样暴露给用户。
流式不是越多越好。对结构化决策和风险判断,完整结果比逐字输出更重要。
# 可观测性
RunnableParallel 的一个好处是 trace 更清楚。每个分支可以被单独看到,耗时、错误、token 使用也更容易定位。
建议给关键分支加稳定名字:
retrieval_chain = RunnableParallel({
"product_docs": product_retriever.with_config({"run_name": "retrieve_product_docs"}) | format_docs,
"ticket_docs": ticket_retriever.with_config({"run_name": "retrieve_ticket_docs"}) | format_docs,
"question": RunnablePassthrough(),
}).with_config({
"run_name": "parallel_retrieval",
"tags": ["rag", "retrieval"],
})
2
3
4
5
6
7
8
调用时把请求级 metadata 放到 config:
answer = chain.invoke(
"如何申请企业版发票?",
config={
"metadata": {
"tenant_id": "tenant_001",
"request_id": "req_20260806_001",
"prompt_version": "kb-qa-2026-08-06-v1",
}
},
)
2
3
4
5
6
7
8
9
10
不要把 tenant_id、request_id 这类只用于观测的字段塞进 Prompt,除非模型确实需要知道。业务输入和运行配置分开,链路会干净很多。
# 错误处理
并行分支中任何一个分支失败,都可能导致整个链路失败。
例如:
- 产品文档 retriever 超时。
- FAQ retriever 返回空。
- 某个模型分支限流。
- 某个 lambda 字段缺失。
- parser 解析失败。
处理策略取决于业务:
- 必要分支失败:整条链路失败,返回稳定错误码。
- 可选分支失败:记录错误,用空上下文或降级结果继续。
- 批处理失败:隔离单条,不影响其他输入。
- 高风险判断失败:默认进入人工确认,而不是放行。
不要把所有分支都当成同等重要。生产链路要区分 required context 和 optional context。
可以把可选分支包成更稳的函数:
def safe_retrieve_faq(question: str) -> str:
try:
docs = faq_retriever.invoke(question)
return format_docs(docs)
except Exception:
return ""
2
3
4
5
6
再放进并行链路:
retrieval_chain = {
"product_docs": product_retriever | format_docs,
"faq_docs": safe_retrieve_faq,
"question": RunnablePassthrough(),
}
2
3
4
5
关键是错误要进入日志或 trace。返回空字符串可以作为降级,但不能让团队完全不知道 FAQ 检索已经坏了。
# 类型边界
RunnableParallel 和 RunnablePassthrough 都很灵活,灵活也意味着容易乱。
建议为重要链路定义输入输出模型:
from pydantic import BaseModel
class QAInput(BaseModel):
question: str
tenant_id: str
locale: str = "zh-CN"
class RetrievalContext(BaseModel):
question: str
product_docs: str
faq_docs: str
2
3
4
5
6
7
8
9
10
11
12
13
在服务边界先校验:
def answer(payload: dict) -> str:
qa_input = QAInput.model_validate(payload)
return chain.invoke(qa_input.model_dump())
2
3
4
不要让前端随便传进来的 dict 直接进入 LCEL 链路。字段缺失、类型不一致、空字符串、超长文本,都会在更深的位置变成难查的问题。
# 目录组织建议
可以把可复用分支单独命名:
# internal/ai/chains/retrieval.py
from langchain_core.runnables import RunnableParallel, RunnablePassthrough
def build_retrieval_context_chain(product_retriever, faq_retriever):
return RunnableParallel({
"product_docs": product_retriever | format_docs,
"faq_docs": faq_retriever | format_docs,
"question": RunnablePassthrough(),
}).with_config({
"run_name": "retrieval_context",
"tags": ["rag", "retrieval"],
})
2
3
4
5
6
7
8
9
10
11
12
13
主链路只负责组合:
# internal/ai/chains/qa.py
from langchain_core.output_parsers import StrOutputParser
from internal.ai.chains.retrieval import build_retrieval_context_chain
from internal.ai.models import get_default_chat_model
from internal.ai.prompts.qa import qa_prompt
def build_qa_chain(product_retriever, faq_retriever):
retrieval_context = build_retrieval_context_chain(
product_retriever,
faq_retriever,
)
return (
retrieval_context
| qa_prompt
| get_default_chat_model()
| StrOutputParser()
)
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
这样做有几个好处:
- 检索上下文链路可以单独测试。
- 主 QA 链路更短。
- trace 名称更稳定。
- 后续新增来源不需要改 route。
- 可选分支和必选分支更容易区分。
# 测试建议
先测试 RunnableParallel 的数据形状:
def test_parallel_context_shape():
chain = {
"question": RunnablePassthrough(),
"length": lambda text: len(text),
}
result = chain.invoke("hello")
assert result == {
"question": "hello",
"length": 5,
}
2
3
4
5
6
7
8
9
10
11
12
再测试 assign 是否保留原始字段:
def test_assign_keeps_original_fields():
chain = RunnablePassthrough.assign(
context=lambda payload: f"context for {payload['question']}"
)
result = chain.invoke({
"question": "LCEL 是什么?",
"tenant_id": "tenant_001",
})
assert result["question"] == "LCEL 是什么?"
assert result["tenant_id"] == "tenant_001"
assert result["context"] == "context for LCEL 是什么?"
2
3
4
5
6
7
8
9
10
11
12
13
最后测试主链路时,可以用假 retriever 和假 model,不要每次都真实调用模型:
from langchain_core.runnables import RunnableLambda
fake_retriever = RunnableLambda(lambda question: [
FakeDoc("RunnableParallel 会把同一输入分发给多个分支。")
])
fake_model = RunnableLambda(lambda _: "固定回答")
chain = (
{
"context": fake_retriever | format_docs,
"question": RunnablePassthrough(),
}
| prompt
| fake_model
| StrOutputParser()
)
def test_rag_chain():
assert chain.invoke("RunnableParallel 是什么?") == "固定回答"
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
LLM 应用不是不能测试,而是要把模型调用和链路结构拆开测。Runnable 的统一协议正好让这件事变简单。
# 常见坑
# 以为 RunnablePassthrough 会自动取字段
RunnablePassthrough() 是原样传递输入,不会自动帮你从 dict 里取 question。
如果输入是 dict:
RunnablePassthrough().invoke({
"question": "什么是 LCEL?"
})
2
3
输出还是整个 dict。
如果只想取字段,要显式写:
lambda payload: payload["question"]
# 忘记 dict 分支拿到的是同一份输入
RunnableParallel 的每个分支都拿到同一份输入。不是上一个 key 的输出传给下一个 key。
如果分支之间有先后依赖,不要写成并行 dict,要用顺序链路或 assign。
# 并行分支过多
并行分支多了,trace 会变复杂,成本和限流压力也会上来。能合并的轻量规则可以合并,真正需要独立观测和独立失败处理的分支再拆出来。
# 把敏感字段透传进 Prompt
assign 会保留原始 dict。如果原始输入里有 token、内部权限、成本策略、风控字段,后续 Prompt 可能无意中拿到它们。进入 Prompt 前要明确字段白名单。
# 在分支里隐藏外部副作用
并行分支适合读取和计算,不适合偷偷做写操作。发消息、扣费、删除、改权限这类副作用不应该藏在 RunnableParallel 里。
# 生产 Checklist
使用这两个类前,可以问自己几个问题:
- 每个分支是否真的可以并行?
- 每个分支拿到的输入类型是否一致?
- 输出 dict 的字段是否有清晰命名?
- Prompt 需要的字段是否明确?
- 原始输入里是否有不该给模型看的字段?
- 可选分支失败时是否能降级?
- 必要分支失败时是否能快速暴露?
- 多分支模型调用是否会打爆 rate limit?
- batch 时总调用次数是否可控?
- trace 里能否看懂每个分支的作用?
把这些想清楚,RunnableParallel 和 RunnablePassthrough 就不是炫技写法,而是非常实用的数据流工具。
# 小结
RunnableParallel 解决的是分流问题:同一份输入同时进入多个分支,最后合并成一个 dict。
RunnablePassthrough 解决的是保留问题:把原始输入原样带到后续链路,也可以通过 assign 在原始 dict 上追加字段。
这两个类通常一起出现在 RAG、并行检索、上下文组装、多模型对比和风险判断里。它们本身不复杂,但数据流一旦理解错,LCEL 链路就会变得很难调试。
生产里真正要守住的是边界:并行分支要可观测,透传字段要可控,失败策略要明确,高风险副作用要离开 LCEL 链路。这样写出来的 LangChain 代码,才会既简洁又能长期维护。