Deep Read

LangChain - LCEL 高级特性与组件

 

LCEL(LangChain Expression Language)用 | 把各个 Runnable 串成链。除了 prompt | llm | parser 这种基础用法,还有几个常用组件能把自定义函数、并行分支、数据透传等能力接进链里。

组件作用
RunnableLambda把普通 Python 函数包装成 Runnable,接进链里
RunnableParallel多个分支接收同一份输入,并行执行,结果汇总成 dict
RunnablePassthrough把输入原样传递下去,或用 assign 增强后再传

RunnableLambda —— 把自定义函数加入链

 

RunnableLambda 把任意普通 Python 函数(自定义函数)包装成一个 Runnable,从而能用 | 接进链里。

案例:在 LCEL 链中用 RunnableLambda 对 LLM 的输出统计字数。

from langchain_core.runnables import RunnableLambda

llm = ChatOpenAI(api_key=api_key, base_url=base_url, model=model, streaming=True)

def count_and_wrap(text: str) -> str:
    word_count = len(text)
    return f"【AI回复共 {word_count} 个字符】\n\n{text}"

prompt = ChatPromptTemplate.from_template(
    '请用中文写一段关于"{topic}"的简短自我介绍,100字以内。'
)

chain = prompt | llm | StrOutputParser() | RunnableLambda(count_and_wrap)

if __name__ == "__main__":
    result = chain.invoke({"topic": "Python编程"})
    print(result)

 

另一种把自定义函数包装成 Runnable 的方法是使用装饰器 @chain

from langchain_core.runnables import chain

@chain
def count_and_wrap(text: str) -> str:
    word_count = len(text)
    return f"【AI回复共 {word_count} 个字符】\n\n{text}"

prompt = ChatPromptTemplate.from_template(
    '请用中文写一段关于"{topic}"的简短自我介绍,100字以内。')

my_chain = prompt | llm | StrOutputParser() | count_and_wrap

if __name__ == "__main__":
    result = my_chain.invoke({"topic": "Python编程"})
    print(result)

 

RunnableParallel —— 并行执行多个任务

RunnableParallel(也叫 RunnableMap)是 LCEL 里的一个核心组件,作用:

让多个 Runnable 接收同一份输入,并行执行,然后把各自的结果汇总成一个字典返回

核心特点

  • 同一输入,多路处理:所有分支拿到的是完全相同的输入。
  • 并行执行:各分支同时运行(同步模式下用线程池,异步模式下用 asyncio),整体耗时约等于最慢的那个分支,而不是所有分支耗时之和。
  • 输出是字典:key 是你定义的名字,value 是对应分支的输出。

三种等价写法

from langchain_core.runnables import RunnableParallel

# 写法 1:显式构造
chain = RunnableParallel(joke=joke_chain, poem=poem_chain)

# 写法 2:传字典
chain = RunnableParallel({"joke": joke_chain, "poem": poem_chain})

# 写法 3:在 LCEL 中直接用字典字面量(最常见)
# LangChain 会自动把 dict 转成 RunnableParallel
chain = {"joke": joke_chain, "poem": poem_chain}

第 3 种最常用——只要在管道 | 中出现一个普通的 dict,LangChain 会自动包装成 RunnableParallel

基本示例

from langchain_openai import ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser
from langchain_core.runnables import RunnableParallel

from models import DEEPSEEK_URL, DEEPSEEK_V4_PRO
import os
from dotenv import load_dotenv

load_dotenv()
api_key = os.getenv("MY_DEEPSEEK_API_KEY")
base_url = DEEPSEEK_URL
model = DEEPSEEK_V4_PRO
llm = ChatOpenAI(api_key=api_key, base_url=base_url, model=model)

joke_temp = ChatPromptTemplate.from_template("讲一个关于{topic}的笑话")
poem_temp = ChatPromptTemplate.from_template("写一首关于{topic}的诗")

joke_chain = joke_temp | llm | StrOutputParser()
poem_chain = poem_temp | llm | StrOutputParser()

parallel_chain = RunnableParallel({
    "joke": joke_chain,
    "poem": poem_chain,
})

result = parallel_chain.invoke({"topic": "熊"})
print("笑话:", result["joke"])
print("诗:", result["poem"])

joke_chainpoem_chain同时调用模型,而不是一个接一个。执行流程如下:

flowchart TD
    IN["invoke({topic: 熊})"] --> RP{{"RunnableParallel<br/>同一份输入分发给两个分支"}}

    RP -->|并行| J1["joke_temp<br/>填入 topic"]
    RP -->|并行| P1["poem_temp<br/>填入 topic"]

    subgraph joke_chain
        J1 --> J2["llm 调用模型"] --> J3["StrOutputParser 取纯文本"]
    end

    subgraph poem_chain
        P1 --> P2["llm 调用模型"] --> P3["StrOutputParser 取纯文本"]
    end

    J3 --> OUT["合并成 dict<br/>{joke: ..., poem: ...}"]
    P3 --> OUT

两个分支同时启动,整体耗时约等于较慢的那条链,最后按 key 合并成字典返回。

最常见用途:在 RAG 中传递多个值

这是 RunnableParallel 出场率最高的场景——一边检索上下文,一边把原始问题透传下去:

from langchain_core.runnables import RunnablePassthrough

retrieval_chain = (
    {
        "context": retriever,              # 用问题去检索
        "question": RunnablePassthrough(),  # 原样透传问题
    }
    | prompt
    | model
)
retrieval_chain.invoke("什么是 LangChain?")

这里 contextquestion 并行准备好后,合成字典交给后面的 prompt

常配合使用的伙伴

组件作用
RunnablePassthrough()把输入原样传递下去
RunnablePassthrough.assign(...)在保留原输入的基础上,追加新字段(底层就是用 RunnableParallel 实现的)
RunnableLambda把普通函数包装成 Runnable,放进某个分支
# assign:在原 dict 基础上增加 context 字段
chain = RunnablePassthrough.assign(context=retriever)

小结

  • 本质:多分支并行 + 结果合并为 dict。
  • 价值:① 提升性能(并行 I/O);② 在链路中同时构造多个下游需要的字段。
  • 记忆点:LCEL 里只要看到 { } 字典,基本就是 RunnableParallel 在工作。

 

RunnablePassthrough —— 传递 / 增强数据

RunnablePassthrough 解决的是一个很具体的问题:在链路中途,怎么把"原始输入"保留下来,传给后面真正需要它的环节。

LCEL 的链是单向流水线,前一步的输出就是后一步的输入。但很多时候后面的 prompt 既需要"加工过的数据"(比如检索到的上下文),又需要"原始输入"(比如用户最初的问题)。RunnablePassthrough 就是那条"什么都不做、原样往后递"的支线。

两种用法

用法行为输出
RunnablePassthrough()把输入原封不动地传下去输入是什么,输出就是什么
RunnablePassthrough.assign(key=...)保留原输入的基础上,追加计算出的新字段原 dict + 新字段(要求输入是 dict)

二者的关键区别:() 是"透传",assign() 是"透传 + 增量"。assign 底层就是用 RunnableParallel 实现的——一路透传原值,一路计算新字段,再合并。

示例:先生成答案,再据此出题

下面的链分两步:

  • 第一步用 assign 生成对概念的解释并追加为 answer
  • 第二步的 quiz_template 同时需要 conceptanswer——这正是必须用 assign 而非普通 chain 的原因(普通链只会输出解释字符串,丢掉原始的 concept)。
template = ChatPromptTemplate.from_template("用一句话介绍{concept}的概念")
chain = template | llm | StrOutputParser()

# assign 跑 chain 得到解释,以 answer 为键追加,原 concept 保留:
#   {"concept": "递归"} → {"concept": "递归", "answer": "递归是……"}

# quiz_template 同时需要 {concept} 和 {answer}
quiz_template = ChatPromptTemplate.from_template(
    "关于{concept}的概念:{answer}\n\n请根据以上内容出一道选择题(给出4个选项和正确答案)"
)
quiz_chain = quiz_template | llm | StrOutputParser()

full_chain = RunnablePassthrough.assign(answer=chain) | quiz_chain
print(full_chain.invoke({"concept": "递归"}))

数据流向:

flowchart TD
    IN["invoke({concept: 递归})"] --> RP{{"RunnablePassthrough.assign(answer=chain)"}}
    RP -->|原值透传| A["concept: 递归"]
    RP -->|跑 chain 生成解释| B["answer: 递归是……"]
    A --> M["合并:{concept: 递归, answer: 递归是……}"]
    B --> M
    M --> QC["quiz_chain<br/>用 concept + answer 出题"]
    QC --> OUT["一道关于递归的选择题<br/>(4 个选项 + 正确答案)"]

典型场景:RAG 中保留原始问题

最常见的写法就是 RunnableParallel 里的检索链——一路用问题去检索上下文,另一路用 RunnablePassthrough() 把问题原样透传,让后面的 prompt 同时拿到 contextquestion

retrieval_chain = (
    {
        "context": retriever,               # 用问题去检索
        "question": RunnablePassthrough(),  # 原样透传问题
    }
    | prompt
    | model
)
retrieval_chain.invoke("什么是 LangChain?")

如果上游已经是一个 dict、只想往里补一个 context 字段,用 assign 更简洁:

# 输入 {"question": "..."} → 输出 {"question": "...", "context": <检索结果>}
chain = RunnablePassthrough.assign(context=lambda x: retriever.invoke(x["question"]))

小结

  • RunnablePassthrough():原样透传,常放在链首或并行分支里"留住"原始输入。
  • RunnablePassthrough.assign(...):透传 + 追加字段,不丢原数据,是构造 RAG 输入字典的利器。
  • 记忆点:当后面的环节"既要新数据、又要原始输入"时,就该想到它。

综合示例

把上面几个组件组合起来。这里还用到两个工具:

  • itemgetter("foo"):从输入 dict 里取出 foo 字段,等价于 lambda d: d["foo"]
  • @chain 装饰器:和 RunnableLambda 作用相同,直接把一个函数声明成 Runnable。
from operator import itemgetter
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.runnables import RunnableLambda, chain
from models import get_lc_model_client

client = get_lc_model_client()
chat_template = ChatPromptTemplate.from_template("{a} + {b} 是多少?")

def length_function(text):
    return len(text)

@chain
def multiple_length_function(_dict):
    return len(_dict["text1"]) * len(_dict["text2"])

chain1 = chat_template | client

# 最外层的 dict 会被自动转成 RunnableParallel:并行算出 a、b 再喂给 chain1
full_chain = (
    {
        "a": itemgetter("foo") | RunnableLambda(length_function),
        "b": {"text1": itemgetter("foo"), "text2": itemgetter("bar")} | multiple_length_function,
    }
    | chain1
)
print(full_chain.invoke({"foo": "abc", "bar": "abcd"}))