# LangChain 两个 Runnable 核心类:RunnableParallel 与 RunnablePassthrough

学 LCEL 时,最容易卡住的地方不是 prompt | model | parser 这种顺序链路,而是数据一旦变成 dict,输入到底会怎么流动。

比如 RAG 链路里经常看到这种写法:

rag_chain = (
    {
        "context": retriever | format_docs,
        "question": RunnablePassthrough(),
    }
    | prompt
    | model
    | parser
)
1
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?"
})
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18

这条链路的数据方向很直观:

dict -> prompt -> model -> parser -> str
1

但真实业务经常需要分叉:

  • 同一个用户问题,一边去检索知识库,一边保留问题本身。
  • 同一段文本,同时生成摘要、标题、关键词。
  • 同一份输入,同时跑规则判断、模型判断、风险判断。
  • 同一个问题,同时请求多个 retriever,再合并上下文。

这时就需要 RunnableParallel。

# RunnableParallel:同一输入,多路执行

RunnableParallel 的核心语义是:

同一份输入 -> 多个 Runnable 分支 -> dict 结果
1

最直接的例子:

from langchain_core.runnables import RunnableParallel


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

result = chain.invoke("langchain")

print(result)
1
2
3
4
5
6
7
8
9
10
11

输出结构类似:

{
    "upper": "LANGCHAIN",
    "length": 9,
}
1
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")
1
2
3
4
5
6

当 dict 出现在 LCEL 链路中时,LangChain 会把它转换为并行 Runnable。这就是很多 RAG 示例里 dict 写法的来源。

# dict 不是普通 dict

在 LCEL 里,下面这段代码不是“提前算出一个 Python dict”:

{
    "summary": summary_chain,
    "keywords": keywords_chain,
}
1
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"])
1
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,
})
1
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"],
})
1
2
3
4
5
6
7

在 LangSmith 或日志里,一个稳定的 run_name 会比匿名 dict 更容易排查。

# RunnablePassthrough:保留原始输入

RunnablePassthrough 的核心语义是:

输入什么,就输出什么
1

最小例子:

from langchain_core.runnables import RunnablePassthrough


passthrough = RunnablePassthrough()

result = passthrough.invoke("hello")

print(result)
1
2
3
4
5
6
7
8

输出就是:

"hello"
1

这看起来好像没用,但一旦进入并行分支,它就很有用。

例如:

from langchain_core.runnables import RunnableParallel, RunnablePassthrough


chain = RunnableParallel({
    "original": RunnablePassthrough(),
    "length": lambda text: len(text),
})

result = chain.invoke("langchain")

print(result)
1
2
3
4
5
6
7
8
9
10
11

输出:

{
    "original": "langchain",
    "length": 9,
}
1
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("数据库迁移上线前要检查什么?")
1
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
1
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",
})
1
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",
}
1
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)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17

输出类似:

{
    "user_id": "u_1001",
    "question": "我可以使用高级报表吗?",
    "profile": "用户 u_1001 是企业版客户。",
}
1
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()
)
1
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"]),
})
1
2
3
4

输出只包含你声明的字段。

RunnablePassthrough.assign 是“保留原始 dict,再追加字段”:

chain = RunnablePassthrough.assign(
    context=lambda payload: retriever.invoke(payload["question"]),
)
1
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()
1
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": "请解释数据库迁移的灰度策略。"
})
1
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"]
1
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()
)
1
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)
1
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)
1
2

并行链路:

for chunk in parallel_chain.stream({"content": "..." }):
    print(chunk)
1
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"],
})
1
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",
        }
    },
)
1
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 ""
1
2
3
4
5
6

再放进并行链路:

retrieval_chain = {
    "product_docs": product_retriever | format_docs,
    "faq_docs": safe_retrieve_faq,
    "question": RunnablePassthrough(),
}
1
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
1
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())
1
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"],
    })
1
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()
    )
1
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,
    }
1
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 是什么?"
1
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 是什么?") == "固定回答"
1
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?"
})
1
2
3

输出还是整个 dict。

如果只想取字段,要显式写:

lambda payload: payload["question"]
1

# 忘记 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 代码,才会既简洁又能长期维护。