1. Runnable组合与链式调用概述
在LangChain 1.2.7版本中,Runnable作为核心执行单元,其组合能力是构建复杂工作流的基础。通过将多个Runnable组件按逻辑顺序串联或并行组合,我们可以创建结构化执行流程,实现输入输出的自动传递与数据格式适配。
1.1 核心设计理念
Runnable组合体系的设计主要围绕以下几个关键目标:
- 组件解耦与复用:每个Runnable组件只关注单一功能,通过组合实现复杂逻辑,降低维护成本
- 统一语法规范:提供标准化的组合方式,兼容同步和异步执行场景
- 动态配置传递:支持在整个执行链中传递配置参数和上下文信息
- 异常边界处理:提供统一的异常处理机制,确保流程稳定性
- 版本兼容性:完全适配1.2.7版本的Runnable核心方法(invoke/ainvoke、batch/abatch等)
1.2 主要组合方式
LangChain 1.2.7版本提供了三种核心组合方式:
- pipe(管道):线性串联多个组件,前一个组件的输出作为后一个组件的输入
- assign(并行赋值):通过RunnableParallel实现多个组件的并行执行
- RunnableLambda:将普通Python函数(同步/异步)封装为Runnable组件
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 环境准备与依赖管理
2.1 版本控制要点
为确保代码兼容性,必须严格匹配以下依赖版本:
bash复制pip install langchain==1.2.7
pip install langchain-core==1.2.7
pip install langchain-community==0.4.1
pip install langchain-openai==1.1.7
pip install langchain-classic==1.0.1
pip install openai==1.13.3 # 适配langchain-openai 1.1.7的兼容版本
注意:版本不匹配可能导致语法错误或运行时异常,特别是在使用异步功能时。
2.2 基础导入语句
以下导入语句是后续所有示例的基础:
python复制from langchain_core.runnables import (
Runnable,
RunnableParallel,
RunnableLambda,
RunnableSequence,
)
from langchain_core.output_parsers import StrOutputParser
from langchain_openai import ChatOpenAI
from langchain.prompts import ChatPromptTemplate
import os
# 环境配置(确保OpenAI API密钥有效)
os.environ["OPENAI_API_KEY"] = "your-openai-api-key"
3. pipe:线性串联组合
3.1 基本语法与工作原理
pipe是Runnable最基础的组合方式,实现多个组件的线性执行。在1.2.7版本中,可以使用|操作符或RunnableSequence来创建管道。
python复制# 使用|操作符
chain = component1 | component2 | component3
# 等价于使用RunnableSequence
chain = RunnableSequence.from_sequence([component1, component2, component3])
管道中的每个组件必须实现Runnable协议(即具备invoke方法),且前一个组件的输出格式必须与后一个组件的输入格式兼容。
3.2 完整示例:文本生成→翻译→摘要
python复制# 定义各个组件
prompt = ChatPromptTemplate.from_template("生成一篇关于 {topic} 的英文短文(100词左右)")
llm = ChatOpenAI(model="gpt-3.5-turbo", temperature=0.7)
translator_prompt = ChatPromptTemplate.from_template("将以下英文文本翻译成中文:{text}")
summarizer_prompt = ChatPromptTemplate.from_template("对以下中文文本生成30字以内的摘要:{text}")
# 初始化各个Runnable
generator = prompt | llm | StrOutputParser()
translator = translator_prompt | llm | StrOutputParser()
summarizer = summarizer_prompt | llm | StrOutputParser()
# 组合成完整流程
chain = generator | translator | summarizer
# 执行流程
input_data = {"topic": "人工智能在医疗领域的应用"}
result = chain.invoke(input_data)
print("最终结果:", result)
3.3 执行流程解析
- 生成阶段:
generator接收包含topic的输入,生成英文短文 - 翻译阶段:生成的英文文本自动传递给
translator,输出中文翻译 - 摘要阶段:中文翻译传递给
summarizer,生成最终摘要 - 输出:整个过程无需手动处理中间结果,输入输出自动传递
3.4 异步执行支持
1.2.7版本原生支持异步执行:
python复制import asyncio
async def run_chain():
result = await chain.ainvoke(input_data)
print("异步执行结果:", result)
asyncio.run(run_chain())
4. assign:并行数据补充组合
4.1 基本概念与语法
assign(通过RunnableParallel实现)用于在执行流程中并行补充多个数据字段,其核心特点是:
- 所有组件并行执行(基于asyncio实现)
- 输入数据会同时传递给所有并行组件
- 输出为包含所有结果的字典
python复制# 基础语法
parallel_runnable = RunnableParallel(
main=main_component,
extra1=extra_component1,
extra2=extra_component2
)
# 更简洁的写法(推荐)
parallel_runnable = main_component | RunnableParallel(
extra1=extra_component1,
extra2=extra_component2
)
4.2 实战案例:用户评价分析
python复制# 定义各个组件
main_prompt = ChatPromptTemplate.from_template("写一段关于 {product} 的用户评价(50字左右)")
keyword_prompt = ChatPromptTemplate.from_template("提取以下文本的3个核心关键词:{text}")
sentiment_prompt = ChatPromptTemplate.from_template("分析以下文本的情感倾向(正面/负面/中性):{text}")
main_chain = main_prompt | llm | StrOutputParser()
keyword_chain = keyword_prompt | llm | StrOutputParser()
sentiment_chain = sentiment_prompt | llm | StrOutputParser()
# 组合并行流程
parallel_chain = RunnableParallel(
original_text=main_chain,
keywords=keyword_chain,
sentiment=sentiment_chain
)
# 执行并获取结果
input_data = {"product": "无线蓝牙耳机"}
result = parallel_chain.invoke(input_data)
print("并行组合结果:", result)
# 输出示例:{"original_text": "...", "keywords": "...", "sentiment": "..."}
4.3 进阶用法:并行结果作为后续输入
python复制final_prompt = ChatPromptTemplate.from_template(
"用户评价:{original_text}\n"
"关键词:{keywords}\n"
"情感倾向:{sentiment}\n"
"请基于以上信息,生成一份产品改进建议(30字左右)"
)
final_chain = parallel_chain | final_prompt | llm | StrOutputParser()
result = final_chain.invoke(input_data)
print("最终改进建议:", result)
5. RunnableLambda:自定义逻辑嵌入
5.1 基本语法与特性
RunnableLambda用于将普通Python函数封装为Runnable组件:
python复制# 同步函数封装
def sync_func(input_data):
# 处理逻辑
return processed_data
runnable_sync = RunnableLambda(sync_func)
# 异步函数封装
async def async_func(input_data):
# 异步处理逻辑
return processed_data
runnable_async = RunnableLambda(async_func)
# 匿名函数简化
runnable_lambda = RunnableLambda(lambda x: x.upper())
5.2 数据清洗实战案例
python复制def clean_text(text):
"""去除特殊字符,统一换行符"""
import re
text = re.sub(r"[^\u4e00-\u9fa5a-zA-Z0-9\s]", "", text)
text = re.sub(r"\n+", "\n", text).strip()
return {"cleaned_text": text}
# 创建清洗组件
cleaner = RunnableLambda(clean_text)
# 组合流程:生成→清洗→摘要
generate_chain = ChatPromptTemplate.from_template("写一段关于 {theme} 的短文") | llm | StrOutputParser()
summarize_chain = ChatPromptTemplate.from_template("摘要:{cleaned_text}") | llm | StrOutputParser()
final_chain = generate_chain | cleaner | summarize_chain
# 执行
result = final_chain.invoke({"theme": "环境保护"})
print("清洗后摘要:", result)
5.3 异步API调用案例
python复制async def fetch_external_data(query):
"""模拟异步API调用"""
import aiohttp
async with aiohttp.ClientSession() as session:
async with session.get(f"https://api.example.com/search?q={query}") as resp:
return await resp.json()
async_fetcher = RunnableLambda(fetch_external_data)
# 组合到流程中
chain = (
ChatPromptTemplate.from_template("生成关于{topic}的搜索查询")
| llm
| StrOutputParser()
| async_fetcher
)
# 异步执行
result = await chain.ainvoke({"topic": "最新AI技术"})
6. 复杂流程构建与最佳实践
6.1 混合组合案例:问答系统
python复制# 1. 定义各环节组件
def retrieve_relevant_documents(query):
"""模拟检索相关文档"""
docs = [
"Langchain 1.2.7 版本中,Runnable 协议支持多组件组合",
"RunnableParallel 可实现组件并行执行,提升效率",
"RunnableLambda 支持同步/异步函数封装"
]
return {"query": query, "relevant_docs": "\n".join(docs)}
retriever = RunnableLambda(retrieve_relevant_documents)
generate_prompt = ChatPromptTemplate.from_template(
"基于以下参考文档回答问题:\n{relevant_docs}\n问题:{query}\n回答:"
)
generator = generate_prompt | llm | StrOutputParser()
def validate_answer(data):
"""校验回答质量"""
answer = data["answer"]
return {
"answer": answer,
"is_valid": len(answer) > 50,
"reason": "长度符合要求" if len(answer) > 50 else "长度不足50字"
}
validator = RunnableLambda(validate_answer)
optimize_prompt = ChatPromptTemplate.from_template(
"原回答:{answer}\n问题:{query}\n参考文档:{relevant_docs}\n"
"要求:回答长度需大于50字,请重新生成"
)
optimizer = optimize_prompt | llm | StrOutputParser()
# 2. 组合复杂流程
chain = (
retriever
| RunnableParallel(
original_answer=generator,
doc_length=RunnableLambda(lambda x: len(x["relevant_docs"]))
)
| RunnableLambda(lambda x: {**x, "validation": validator.invoke({"answer": x["original_answer"]})})
| RunnableLambda(
lambda x: optimizer.invoke({
"answer": x["validation"]["answer"],
"query": x["query"],
"relevant_docs": x["relevant_docs"]
}) if not x["validation"]["is_valid"] else x["validation"]["answer"]
)
)
# 3. 执行
result = chain.invoke({"query": "Langchain 1.2.7中Runnable有哪些组合方式?"})
print("最终回答:", result)
6.2 性能优化技巧
-
并行执行:对无依赖关系的组件使用
RunnableParallelpython复制# 不推荐 - 顺序执行 chain = component1 | component2 | component3 # 推荐 - 并行执行component2和component3 chain = component1 | RunnableParallel(result2=component2, result3=component3) -
异步适配:I/O密集型操作封装为异步函数
python复制async def async_io_operation(input): # 异步I/O操作 return result async_chain = RunnableLambda(async_io_operation) -
缓存策略:使用
RunnableWithCache缓存重复计算结果python复制from langchain_core.runnables import RunnableWithCache cached_chain = RunnableWithCache( generator, cache_key=lambda x: f"prompt:{x['topic']}" )
6.3 错误处理机制
-
全局异常处理
python复制def exception_handler(exc) -> str: return f"执行出错:{str(exc)}" safe_chain = chain.with_config( exception_handlers={Exception: exception_handler} ) -
组件级异常处理
python复制safe_generator = generator.with_config( exception_handlers={Exception: lambda exc: "生成失败,请重试"} )
6.4 扩展性设计模式
-
工厂模式创建组件
python复制def create_llm_component(model="gpt-3.5-turbo", temperature=0.7) -> Runnable: return ChatOpenAI(model=model, temperature=temperature) | StrOutputParser() # 使用时 llm_component = create_llm_component(model="gpt-4") -
动态配置注入
python复制from langchain_core.runnables import configurable @configurable(fields={"temperature": float}) def generate_with_config(input_data, temperature=0.7): prompt = ChatPromptTemplate.from_template("生成关于 {topic} 的文本") return (prompt | ChatOpenAI(temperature=temperature) | StrOutputParser()).invoke(input_data) configurable_chain = RunnableLambda(generate_with_config) # 执行时动态配置 result = configurable_chain.invoke( {"topic": "AI"}, config={"configurable": {"temperature": 0.3}} )
7. 版本迁移与兼容性
7.1 与低版本的主要差异
| 低版本功能 | 1.2.7版本替代方案 |
|---|---|
SimpleSequentialChain |
pipe或RunnableSequence |
ParallelChain |
RunnableParallel |
LambdaChain |
RunnableLambda |
AsyncLambdaChain |
RunnableLambda(支持异步) |
7.2 常见迁移问题解决
-
链式调用语法变化
python复制# 旧版 from langchain.chains import SimpleSequentialChain chain = SimpleSequentialChain(chains=[c1, c2, c3]) # 新版 chain = c1 | c2 | c3 # 或 chain = RunnableSequence.from_sequence([c1, c2, c3]) -
并行处理语法变化
python复制# 旧版 from langchain.chains import ParallelChain chain = ParallelChain(chains=[c1, c2], output_keys=["k1", "k2"]) # 新版 chain = RunnableParallel(k1=c1, k2=c2) -
Lambda函数封装变化
python复制# 旧版 from langchain.chains import LambdaChain chain = LambdaChain(func=lambda x: x.upper()) # 新版 chain = RunnableLambda(lambda x: x.upper())
8. 实际应用中的经验分享
8.1 组件设计原则
-
单一职责:每个Runnable组件应该只做一件事,并做好这件事。例如:
- 一个组件负责生成文本
- 一个组件负责解析输出
- 一个组件负责数据清洗
-
合理粒度:组件不宜过大或过小。过大的组件难以复用,过小的组件会增加组合复杂度。
-
格式统一:尽量使用字典作为输入输出格式,便于参数传递和并行组合。
8.2 调试技巧
-
逐步构建:先测试单个组件,再逐步组合,最后构建完整流程。
-
中间结果检查:使用
RunnableLambda插入调试打印:python复制def debug_print(data): print("当前数据:", data) return data debugger = RunnableLambda(debug_print) chain = component1 | debugger | component2 -
异常定位:为关键组件单独配置异常处理器,精确定位问题来源。
8.3 性能调优
-
并行化:识别可以并行执行的独立任务,使用
RunnableParallel。 -
异步化:将I/O密集型操作封装为异步函数,使用
ainvoke执行。 -
缓存:对重复计算的结果使用
RunnableWithCache缓存。 -
批处理:对于批量任务,使用
batch或abatch方法而非循环调用。
8.4 可维护性建议
-
文档注释:为每个自定义组件添加清晰的文档字符串。
-
配置集中管理:将模型参数、API密钥等配置集中管理。
-
版本控制:严格记录依赖版本,使用requirements.txt或pyproject.toml。
-
单元测试:为关键组件编写单元测试,确保组合后的行为符合预期。
