分支链与并行链
构建 RunnableParallel 和 RunnableBranch 结构,同时运行多条链,或根据动态条件将输入路由到不同的链。
分支链与并行链 是 CoddyKit 上的免费 AI Engineering Academy 课时。 这是第 3 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 AI Engineering Academy 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 AI Engineering Academy 课程共包含 4 节课。
并行与分支为何重要
在真实世界的 LLM 流水线中,通常需要同时执行多项任务,或根据内容以不同方式路由请求。并行链会同时运行多个分支,在任务彼此独立时可以降低延迟。分支链则会根据动态条件,将输入路由到不同的专用链。LCEL 通过 RunnableParallel 和 RunnableBranch 原生支持这两种模式。
RunnableParallel 基础
RunnableParallel 接收一个字典,其中每个键都映射到一个可运行对象。调用时,它会并发运行所有分支,并返回一个字典,其中每个键对应其分支的结果。当您希望从同一个输入生成多个输出时,这种方式非常理想,例如同时生成摘要并提取关键词。
from langchain_core.runnables import RunnableParallel
from langchain_openai import ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser
model = ChatOpenAI(model='gpt-4o-mini')
parser = StrOutputParser()
parallel = RunnableParallel(
summary=(
ChatPromptTemplate.from_template('Summarize: {text}') | model | parser
),
keywords=(
ChatPromptTemplate.from_template('Extract keywords from: {text}') | model | parser
)
)
result = parallel.invoke({'text': 'Long document text here...'})
print(result['summary'])
print(result['keywords'])使用字典简写实现并行
LCEL 提供了一种便捷的简写方式:将普通字典作为管道链中的一个步骤传入时,系统会自动将其封装为 RunnableParallel。这样,并行分支的使用更加自然,也无需显式实例化类。字典的键会成为结果的键,而值则是并发运行的分支。
from langchain_core.runnables import RunnablePassthrough
# Dict shorthand creates RunnableParallel automatically
chain = (
RunnablePassthrough.assign(
sentiment=(
ChatPromptTemplate.from_template('Sentiment of: {review}')
| model | parser
),
aspects=(
ChatPromptTemplate.from_template('List aspects mentioned in: {review}')
| model | parser
)
)
)
result = chain.invoke({'review': 'Great battery but poor camera quality.'})
print(result['sentiment'])
print(result['aspects'])理解 RunnableBranch
RunnableBranch 会根据条件将输入路由到不同的链。您需要提供一个由 (condition, runnable) 对组成的列表,以及一个默认可运行对象。分支会按顺序评估条件,并运行第一个匹配的分支。这支持基于意图的路由,例如将客户服务查询发送到支持链,将技术问题发送到文档链。
from langchain_core.runnables import RunnableBranch
technical_chain = (
ChatPromptTemplate.from_template('Technical answer: {query}') | model | parser
)
general_chain = (
ChatPromptTemplate.from_template('General answer: {query}') | model | parser
)
branch = RunnableBranch(
(lambda x: 'error' in x['query'].lower() or 'bug' in x['query'].lower(),
technical_chain),
general_chain # default branch
)
result = branch.invoke({'query': 'I got a TypeError in my code'})
# Routes to technical_chain because 'error' is in the query使用 LLM 分类进行语义路由
一种更灵活的路由模式是使用分类器 LLM 调用来确定应使用哪个分支。路由器首先调用一个较小的模型,对输入的意图进行分类,然后根据分类结果将请求路由到相应的专用链。这种方式可以处理关键词匹配无法识别的细微情况,但代价是需要额外进行一次 API 调用。
from langchain_core.output_parsers import StrOutputParser
# Step 1: classify intent
classify_prompt = ChatPromptTemplate.from_template(
'Classify this query as exactly one of: billing, technical, general.\nQuery: {query}'
)
classifier = classify_prompt | model | StrOutputParser()
# Step 2: route based on classification
def route(classification_result: dict):
topic = classification_result['topic'].strip().lower()
if topic == 'billing':
return billing_chain
elif topic == 'technical':
return technical_chain
return general_chain
full_chain = (
RunnablePassthrough.assign(topic=lambda x: classifier.invoke(x))
| RunnableLambda(route)
)并行 RAG:多个检索器
在高级 RAG 系统中,您可能需要同时从多个数据源检索信息,然后合并结果。RunnableParallel 允许您同时查询产品数据库、FAQ 存储和文档索引。随后,合并步骤会整合排名靠前的结果,再将上下文传递给 LLM,为模型提供更丰富的信息基础。
from langchain_core.runnables import RunnableParallel, RunnablePassthrough
# Assume these retrievers are already set up
faq_retriever = faq_vectorstore.as_retriever(search_kwargs={'k': 3})
doc_retriever = doc_vectorstore.as_retriever(search_kwargs={'k': 3})
retrieval = RunnableParallel(
faq_results=faq_retriever,
doc_results=doc_retriever
)
def merge_docs(retrieved: dict) -> str:
all_docs = retrieved['faq_results'] + retrieved['doc_results']
return '\n\n'.join(d.page_content for d in all_docs)
pipeline = (
retrieval
| RunnableLambda(merge_docs)
| ChatPromptTemplate.from_template('Context: {context}\nAnswer: {question}')
| model | parser
)使用 itemgetter 的条件链
当您需要从并行结果中选择特定键,或只将部分上下文传递给下一步时,Python 的 operator.itemgetter 可以作为一种轻量级 Runnable 选择器。并行步骤之后尤其适合使用它:当不同分支生成不同的键,而您只需提取与下一阶段处理相关的键时,它会非常有用。
from operator import itemgetter
from langchain_core.runnables import RunnablePassthrough
# After a parallel step, extract just the summary
chain = (
RunnableParallel(
summary=summary_chain,
sentiment=sentiment_chain
)
| itemgetter('summary') # pass only the summary onward
| translate_chain
)
# itemgetter works because dict.__getitem__ is a valid transform衡量并行加速效果
RunnableParallel 的主要优势在于减少实际耗时。三个依次执行、每个耗时 2 秒的 LLM 调用,总共需要 6 秒。而并行执行时,它们大约在 2 秒内完成——也就是最慢分支所需的时间。不过,并行调用会同时增加令牌用量,因此请注意速率限制。如有需要,请在 batch() 中使用 max_concurrency,或为每个键设置速率限制。
import time
# Measure sequential time
start = time.time()
result1 = chain_a.invoke(input_data)
result2 = chain_b.invoke(input_data)
result3 = chain_c.invoke(input_data)
seq_time = time.time() - start
print(f'Sequential: {seq_time:.2f}s')
# Measure parallel time
start = time.time()
results = RunnableParallel(a=chain_a, b=chain_b, c=chain_c).invoke(input_data)
par_time = time.time() - start
print(f'Parallel: {par_time:.2f}s')
print(f'Speedup: {seq_time/par_time:.1f}x')异步并行链
要在异步应用中实现真正的非阻塞并行执行,请使用 RunnableParallel.ainvoke()。在底层,LCEL 使用 asyncio.gather() 在事件循环上并发运行各个分支。这在 FastAPI 服务中尤其重要,因为每个请求处理器都是协程——使用异步接口,即使同时发起多个 LLM 调用,也能避免阻塞事件循环。
import asyncio
async def analyze_document(text: str) -> dict:
parallel = RunnableParallel(
summary=summary_chain,
keywords=keyword_chain,
sentiment=sentiment_chain
)
# All three chains run concurrently with asyncio.gather internally
result = await parallel.ainvoke({'text': text})
return result
# In FastAPI:
from fastapi import FastAPI
app = FastAPI()
@app.post('/analyze')
async def analyze(request: dict):
return await analyze_document(request['text'])嵌套并行链与顺序链
复杂的处理流程通常会混合使用顺序步骤和并行步骤。您可以将 RunnableParallel 嵌套在顺序管道中,反过来也可以。例如:先顺序识别意图,然后并行执行检索和上下文格式化,最后再顺序生成最终响应。LangChain 能正确处理这种嵌套,让您无需编写杂乱的回调代码,也能构建清晰易读的复杂处理流程。
# Full pipeline: classify → parallel retrieval → generate
pipeline = (
RunnablePassthrough.assign(
intent=classify_chain # sequential: classify first
)
| RunnablePassthrough.assign(
context=RunnableParallel( # parallel: retrieve from both sources
faq=faq_retriever,
docs=doc_retriever
)
)
| format_context_chain # sequential: format merged context
| generate_answer_chain # sequential: call LLM
)处理分支中的错误
当 RunnableParallel 的某个分支失败时,默认情况下整个并行调用都会引发异常。您可以在各个分支的 Runnable 上使用 .with_fallbacks(),以优雅地处理分支级别的失败。对于 RunnableBranch,可以在 RunnableLambda 中为每个分支编写 try-except,或者使用默认的备用链来处理路由错误。
from langchain_core.runnables import RunnableParallel
# Wrap each branch with a fallback
safe_summary = summary_chain.with_fallbacks([
RunnableLambda(lambda x: 'Summary unavailable')
])
safe_sentiment = sentiment_chain.with_fallbacks([
RunnableLambda(lambda x: 'Sentiment unavailable')
])
robust_parallel = RunnableParallel(
summary=safe_summary,
sentiment=safe_sentiment
)
# Now one branch failing won't kill the entire parallel call快速检查
测试您对 LCEL 中分支链和并行链的理解。
课程回顾
在本课中,您学到了:RunnableParallel 会并发运行多个链,并返回一个结果字典,从而降低相互独立的 LLM 调用的延迟;RunnableBranch 会根据条件或 LLM 分类结果,将输入路由到不同的专用链;通过嵌套并行步骤和顺序步骤,您可以构建复杂的多路径处理流程,同时保持代码清晰易读。接下来,我们将学习 LangChain 中的流式输出。
常见问题解答
「分支链与并行链」课时是免费的吗?
是的 — 「分支链与并行链」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 AI Engineering Academy 课程的其余内容,请升级到 CoddyKit PRO。 AI Engineering Academy 课程共包含 4 节课。
「分支链与并行链」这节课中我会学到什么?
构建 RunnableParallel 和 RunnableBranch 结构,同时运行多条链,或根据动态条件将输入路由到不同的链。 你通过在浏览器中直接运行的动手代码来练习 AI Engineering Academy,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 AI Engineering Academy 需要有经验吗?
无需任何先前经验。CoddyKit 上的 AI Engineering Academy 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 3 节课,共 4 节。
「分支链与并行链」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 AI Engineering Academy 课中编写并运行代码吗?
能。每节 AI Engineering Academy 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。