⚙️ LCEL 核心原理 ★
这是 LangChain 的灵魂章节。我们已经集齐了消息、提示、模型、解析器四块积木——LCEL(LangChain Expression Language)用 | 管道把它们串联起来,并自动获得流式、并发、异步、重试等能力。掌握 LCEL,就掌握了 LangChain 的设计哲学。
本章目标
- 理解
Runnable抽象——LangChain 的"统一接口" - 掌握
|管道(RunnableSequence)和字典并发(RunnableParallel)两大组合原语 - 会用
RunnablePassthrough/assign/pick处理数据流 - 能独立写出 RAG 检索链(LCEL 的标志性应用)
- 理解"声明式组合"为何自动获得 batch/stream/async 能力
前 4 章你学了 4 个独立的积木。从这章开始,它们能组合了。而"组合"正是 LLM 应用的核心——RAG 是"检索+生成"的组合,Agent 是"思考+行动"的组合。LCEL 让组合变得像搭积木一样自然。
Runnable:统一接口
前面每章都提到"XX 是 Runnable"——消息模板是、模型是、解析器是。那 Runnable 到底是什么?看官方定义(源码 docstring 极其精炼):
class Runnable(ABC, Generic[Input, Output]):
"""一个"工作单元",可以被 invoke、batch、stream、transform 和组合。"""
# 核心方法(所有 Runnable 都有):
# - invoke/ainvoke : 把单个输入转成输出
# - batch/abatch : 高效地把多个输入转成输出(默认用线程池并行)
# - stream/astream : 流式输出单个输入的结果
# - astream_log : 流式输出 + 选定的中间结果
# 内建优化:
# - Batch: 默认用线程池并行跑多个 invoke(),子类可重写优化
# - Async: 'a' 前缀的方法是异步版,默认降级用线程池跑同步版
# 所有方法都接受可选的 config 参数(tags/metadata/callbacks/并发数)
# 暴露 schema 信息:input_schema / output_schema / config_schema()
用大白话:Runnable 就是一个"能被统一调用的盒子",有进有出,并且天生支持同步/异步/批量/流式四种调用方式。
哪些东西是 Runnable?
| Runnable | 输入 | 输出 | 从哪章来 |
|---|---|---|---|
ChatPromptTemplate | dict | 消息列表 | LC2 |
ChatModel | 消息 | AIMessage | LC3 |
OutputParser | AIMessage | 目标类型 | LC4 |
RunnableLambda | 任意 | 任意 | 本章(包普通函数) |
RunnablePassthrough | 任意 | 原样输出 | 本章 |
BaseTool | dict | 任意 | LC6 |
管道 | :RunnableSequence
这是 LCEL 最标志性的语法。用 | 把多个 Runnable 串联,前一个的输出自动成为后一个的输入:
class Runnable:
def pipe(self, *others, name=None):
"""把这个 Runnable 和 others 串联,返回 RunnableSequence。
`a | b` 等价于 `a.pipe(b)`。"""
return RunnableSequence(self, *others, name=name)
def __or__(self, other):
"""重载 | 运算符,让 a | b 语法成立。"""
return self.pipe(other)
# 所以这三行等价:
chain = prompt.pipe(model).pipe(parser) # 显式
chain = prompt | model | parser # 语法糖(最常用)
chain = RunnableSequence(prompt, model, parser) # 最底层
管道能成立的前提是类型衔接:前一个的输出类型 = 后一个的输入类型。prompt 输出消息 → model 接收消息 ✓。如果塞一个不匹配的,运行时会报错。Runnable 的 Generic[Input, Output] 泛型就是为类型安全设计的。
为什么 | 自动获得 batch/stream/async
这是 LCEL 最让人惊艳的地方。你只写了 prompt | model | parser,没写任何并发/流式代码,但它自动支持:
chain = prompt | model | StrOutputParser()
# 1. 单次同步
chain.invoke({"question": "你好"})
# 2. 并发批量(自动用线程池并行!)
chain.batch([{"question": "你好"}, {"question": "再见"}, {"question": "谢谢"}])
# 3. 流式(逐字输出,自动透传 model 的 stream)
for chunk in chain.stream({"question": "讲个故事"}):
print(chunk, end="", flush=True)
# 4. 异步(自动有 ainvoke/abatch/astream)
result = await chain.ainvoke({"question": "你好"})
原理:RunnableSequence 的 batch 实现,会把每个输入并行跑 invoke(用线程池);stream 实现,会逐级透传每个 Runnable 的流式输出。你声明了"做什么"(结构),框架自动给你"怎么做"(并发/流式/异步)——这就是"声明式"的威力。
class RunnableSequence(RunnableSerializable):
"""管道组合。把多个 Runnable 串联,前一个输出 = 后一个输入。
自动实现:
- invoke: 依次调用每一步
- batch: 并发处理多个输入(每步内部并行)
- stream: 逐级透传流式输出
- ainvoke/abatch/astream: 异步版(自动降级到线程池)
"""
first: Runnable # 管道第一步
middle: list[Runnable] # 中间步骤
last: Runnable # 最后一步
字典并发:RunnableParallel
第二个组合原语。用一个字典字面量,让多个 Runnable 并发执行,结果合并成 dict:
# 字典字面量 = 并发执行
from langchain_core.runnables import RunnableLambda
# 三个函数并发跑,结果合并成 {"a": ..., "b": ..., "c": ...}
parallel = {
"a": RunnableLambda(lambda x: x + 1),
"b": RunnableLambda(lambda x: x * 2),
"c": RunnableLambda(lambda x: x ** 2),
}
result = RunnableLambda(lambda x: x).pipe(parallel).invoke(5)
# 等价写法:
from langchain_core.runnables import RunnableParallel
result = RunnableParallel(a=..., b=..., c=...).invoke(5)
# → {"a": 6, "b": 10, "c": 25}
在管道里用字典,它会自动变成 RunnableParallel:
# retrieve 和 passthrough 并发,结果合并成 {"context": ..., "question": ...}
chain = (
RunnableParallel({
"context": retriever, # 检索相关文档
"question": RunnablePassthrough() # 原样透传问题
})
| prompt
| model
| parser
)
RunnablePassthrough / assign / pick
处理数据流的三个常用原语,源码在 runnables/passthrough.py:
| 原语 | 作用 | 示例 |
|---|---|---|
RunnablePassthrough() |
原样透传输入(什么也不做,常用于"保留原值") | RunnablePassthrough().invoke(5) → 5 |
.assign(key=fn) |
在原 dict 基础上新增字段(原字段保留) | .assign(extra=lambda d: d["x"]+1) |
.pick("key") |
从 dict 中选取某个/某些字段 | .pick("answer") |
from langchain_core.runnables import RunnablePassthrough
# RunnablePassthrough.assign 保留原 dict,并新增字段
chain = RunnablePassthrough.assign(
summary=lambda d: f"问题长度={len(d['question'])}"
)
chain.invoke({"question": "你好吗"}) # 注意是 dict 输入
# → {"question": "你好吗", "summary": "问题长度=3"} 原字段还在!
# 对比 RunnableLambda:会替换整个输出
from langchain_core.runnables import RunnableLambda
RunnableLambda(lambda d: f"长度{len(d['question'])}").invoke({"question": "你好吗"})
# → "长度3" (原 dict 没了)
实战:用 LCEL 写 RAG 链
RAG(检索增强生成)是 LCEL 最经典的应用。完整流程:"问题 → 并发(检索+保留问题) → 拼 prompt → 调模型 → 解析"。用 LCEL 一行管道搞定:
# pip install langchain langchain-openai chromadb
# export OPENAI_API_KEY="sk-..."
from langchain_openai import ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser
from langchain_core.runnables import RunnablePassthrough, RunnableLambda
model = ChatOpenAI(model="gpt-4o-mini")
# ============ ① 最简单的 chain ============
prompt = ChatPromptTemplate.from_messages([
("system", "你是有帮助的助手。"),
("human", "{question}"),
])
chain = prompt | model | StrOutputParser()
print(chain.invoke({"question": "1+1=?"}))
# ============ ② 用 RunnableLambda 包普通函数 ============
# 任何函数都能变成 Runnable,从而加入管道
def word_count(text: str) -> str:
return f"回答字数: {len(text)}"
chain2 = prompt | model | StrOutputParser() | RunnableLambda(word_count)
print(chain2.invoke({"question": "介绍Python"}))
# ============ ③ 经典 RAG 链 ============
# 假装有个 retriever(LC8 之后会讲真正的向量检索)
def fake_retrieve(question: str) -> str:
return "(这是检索到的相关文档内容...)"
retriever = RunnableLambda(fake_retrieve)
rag_prompt = ChatPromptTemplate.from_messages([
("system", "根据以下上下文回答问题。\n上下文:{context}"),
("human", "{question}"),
])
# LCEL 的精髓在这一行!
rag_chain = (
{ # ← 字典 = RunnableParallel(并发)
"context": retriever, # 检索
"question": RunnablePassthrough(), # 保留原问题
}
| rag_prompt # ← 拼 prompt
| model # ← 调模型
| StrOutputParser() # ← 取文本
)
# 调用时只传问题字符串,retriever 和 passthrough 都接收它
print(rag_chain.invoke("Python是什么?"))
# 流式也自动可用:
# for chunk in rag_chain.stream("Python是什么?"): print(chunk, end="")
- 输入字符串
"Python是什么?" - 进入字典(RunnableParallel):
retriever收到它 → 输出文档;RunnablePassthrough收到它 → 原样输出 - 合并成
{"context": "文档", "question": "Python是什么?"} - 进入
rag_prompt→ 填充变量 → 消息列表 - 进入
model→ AIMessage - 进入
StrOutputParser→ 最终字符串
@chain 装饰器:自定义 Runnable
除了 RunnableLambda,还可以用 @chain 装饰器把函数直接变成 Runnable(更优雅):
from langchain_core.runnables import chain # 注意是小写的装饰器
@chain
def my_step(x: dict) -> str:
"""函数自动变成 Runnable,可放入管道。"""
return f"处理了: {x['question']}"
# 等价于 RunnableLambda(my_step),但写法更简洁
chain = prompt | my_step | model | StrOutputParser()
config:追踪与配置
所有 Runnable 方法都接受 config 参数,用于打标签、追踪、控制并发。源码在 runnables/config.py:
result = chain.invoke(
{"question": "你好"},
config={
"tags": ["production", "v2"], # 追踪标签
"metadata": {"user_id": "u123"}, # 自定义元数据
"max_concurrency": 5, # batch 并发上限
# "callbacks": [...] # 回调(日志/metrics)
}
)
# 这些信息会出现在 LangSmith 追踪里,便于调试和监控
与生产实践对照
| LCEL 概念 | OpenCode 对应 | |
|---|---|---|
| Runnable 统一接口 | → | Effect 框架的 Stream/Effect(统一可组合单元) |
| prompt | model | parser 管道 | → | runLoop 内组装 system prompt → streamText → processor 处理 |
| RunnableParallel 并发 | → | Effect 的 Stream并发处理工具结果 |
| config 追踪标签 | → | OpenTelemetry span(每个工具调用一个 span) |
| stream() 流式透传 | → | LLMEvent 流(text-delta 逐字透传) |
小结
- Runnable 是 LangChain 的统一接口:一个有进有出、天生支持 invoke/batch/stream/async 的"工作单元"。
- 两大组合原语:
|管道(RunnableSequence,串联)和字典(RunnableParallel,并发)。 RunnablePassthrough(透传)/assign(新增字段)/pick(选取)是数据流处理三件套。- "声明式组合"自动获得 batch/stream/async 能力——这是 LCEL 的核心价值。
- RAG 链
{context, question} | prompt | model | parser是 LCEL 的标志性应用。
你已经能组合"提示+模型+解析"了。但 Agent 还差一块——工具。下一章 LC6 · 工具体系:学 BaseTool 和 @tool 装饰器,把普通 Python 函数变成 LLM 可调用的工具。这是 Agent "行动力"的另一半。