0Pricing
AI Engineering Academy · 课时

分支链与并行链

构建 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 反馈 — 无需本地设置。

此课程中的所有课时

  1. LangChain 架构与核心抽象
  2. 使用 LCEL 构建链
  3. 分支链与并行链
  4. 在 LangChain 中流式输出
← 返回 AI Engineering Academy