浅谈LangChain\LangGraph(上)
前言:
简单聊一下使用LangChain的一些心得,个人感觉就是操作大模型调用自己的向量数据库,然后触发自定义工具包的一个生产框架,主要包括RAG(向量数据库)和Agent(工具包)两部分核心组成。
RAG一般来源于企业生产中的文档资料,需要转换成大模型识别的向量数据库,一般要先把Word\Excel\Pdf这种文件资料,先进行整理,转成这种文件内容:(举例汽车客服问答内容)
faq_data = [
# ------------------ 一、 常规保养与维护 (10条) ------------------
{"id": 1, "category": "常规保养", "question": "我这个车轮胎不是很好了,一般像咱家的车多久要更换轮胎?", "answer": "先生/女士,按我们的建议,轮胎一般行驶5年或5万公里左右需要更换。如果平时磨损较重或出现裂纹、鼓包,建议提前来店检测。"},
{"id": 2, "category": "常规保养", "question": "首保是免费的吗?首保需要做哪些项目?", "answer": "是的,新车首保完全免费。首保项目主要包括更换机油、机滤,以及对底盘、制动系统、灯光和全车电子系统进行全面检查。"},
{"id": 3, "category": "常规保养", "question": "新车多久做一次首保?", "answer": "一般建议在购车后行驶到 5000 公里或 6 个月(以先到者为准)时来店做首保。"},
{"id": 4, "category": "常规保养", "question": "如果不按时保养,会影响整车质保吗?", "answer": "如果长期不按厂家标准进行必要保养导致零件损坏,相关损坏部件可能无法享受质保,建议您按时保养以保障您的权益。"},
{"id": 5, "category": "常规保养", "question": "全合成机油多长时间换一次?", "answer": "全合成机油建议每行驶 10000 公里或 1 年更换一次。"},
{"id": 6, "category": "常规保养", "question": "半合成机油和全合成机油有什么区别?", "answer": "全合成机油润滑和抗氧化性能更好,更换周期长(10000公里);半合成机油性价比高,但更换周期稍短(5000-7500公里)。"},
{"id": 7, "category": "常规保养", "question": "空调滤芯多久更换一次?", "answer": "空调滤芯建议每 15000 公里或每年更换一次。如果经常在空气质量较差的地区行驶,建议每半年检查或更换。"},
{"id": 8, "category": "常规保养", "question": "发动机空气滤芯多久换一次?", "answer": "空气滤芯建议每 20000 公里更换一次,定期更换可以保护发动机并降低油耗。"},
{"id": 9, "category": "常规保养", "question": "火花塞多久需要更换?", "answer": "根据火花塞材质不同,普通火花塞 3-4 万公里更换,铱金/铂金火花塞建议 6-8 万公里更换。"},
{"id": 10, "category": "常规保养", "question": "汽车防冻液多久更换一次?", "answer": "一般防冻液建议每 2 年或 4 万公里更换一次,长效防冻液可延长至 5 年或 10 万公里。"}
# 100 这下面还有很多就不列举了
]拿到这种文件后,要把他们变成json:
[
{
"id": 1,
"category": "常规保养",
"question": "我这个车轮胎不是很好了,一般像咱家的车多久要更换轮胎?",
"answer": "先生/女士,按我们的建议,轮胎一般行驶5年或5万公里左右需要更换。如果平时磨损较重或出现裂纹、鼓包,建议提前来店检测。"
},
{
"id": 2,
"category": "常规保养",
"question": "首保是免费的吗?首保需要做哪些项目?",
"answer": "是的,新车首保完全免费。首保项目主要包括更换机油、机滤,以及对底盘、制动系统、灯光和全车电子系统进行全面检查。"
},
{
"id": 3,
"category": "常规保养",
"question": "新车多久做一次首保?",
"answer": "一般建议在购车后行驶到 5000 公里或 6 个月(以先到者为准)时来店做首保。"
},
{
"id": 4,
"category": "常规保养",
"question": "如果不按时保养,会影响整车质保吗?",
"answer": "如果长期不按厂家标准进行必要保养导致零件损坏,相关损坏部件可能无法享受质保,建议您按时保养以保障您的权益。"
},
{
"id": 5,
"category": "常规保养",
"question": "全合成机油多长时间换一次?",
"answer": "全合成机油建议每行驶 10000 公里或 1 年更换一次。"
},
{
"id": 6,
"category": "常规保养",
"question": "半合成机油和全合成机油有什么区别?",
"answer": "全合成机油润滑和抗氧化性能更好,更换周期长(10000公里);半合成机油性价比高,但更换周期稍短(5000-7500公里)。"
},
{
"id": 7,
"category": "常规保养",
"question": "空调滤芯多久更换一次?",
"answer": "空调滤芯建议每 15000 公里或每年更换一次。如果经常在空气质量较差的地区行驶,建议每半年检查或更换。"
},
{
"id": 8,
"category": "常规保养",
"question": "发动机空气滤芯多久换一次?",
"answer": "空气滤芯建议每 20000 公里更换一次,定期更换可以保护发动机并降低油耗。"
},
{
"id": 9,
"category": "常规保养",
"question": "火花塞多久需要更换?",
"answer": "根据火花塞材质不同,普通火花塞 3-4 万公里更换,铱金/铂金火花塞建议 6-8 万公里更换。"
},
{
"id": 10,
"category": "常规保养",
"question": "汽车防冻液多久更换一次?",
"answer": "一般防冻液建议每 2 年或 4 万公里更换一次,长效防冻液可延长至 5 年或 10 万公里。"
}
# 100 这下面还有很多就不列举了
]然后拿到这种格式的json后,我们就可以把他变成向量数据库了:
import json
from langchain_core.documents import Document
from langchain_community.vectorstores import Chroma
from langchain_huggingface import HuggingFaceEmbeddings
import os
os.environ["HF_ENDPOINT"] = "https://hf-mirror.com"
def build_vector_store():
# faq.json 就是上面刚刚那个json文件了
with open("faq.json", "r", encoding="utf-8") as f:
faq_list = json.load(f)
documents = []
for item in faq_list:
page_content = f"问题:{item['question']}\n客服答:{item['answer']}"
metadata = {
"id": item["id"],
"category": item["category"]
}
doc = Document(page_content=page_content, metadata=metadata)
documents.append(doc)
embedding_model = HuggingFaceEmbeddings(
model_name="BAAI/bge-small-zh-v1.5"
)
vectorstore = Chroma.from_documents(
documents=documents,
embedding=embedding_model,
persist_directory="./chroma_db"
)
print("Done")
if __name__ == "__main__":
build_vector_store()经过这个操作后,我们拿到了chroma_db,然后就可以写向量工具类了:
from typing import List
from langchain_chroma import Chroma
from langchain_core.tools import tool
from langchain_huggingface import HuggingFaceEmbeddings
from app.core.reranker import LocalReranker
# 1. 向量数据库初始化
embedding_model = HuggingFaceEmbeddings(
model_name="./models/BAAI--bge-small-zh-v1.5"
)
vectorstore = Chroma(
persist_directory="./chroma_db",
embedding_function=embedding_model
)
# 检索 Top 10
retriever = vectorstore.as_retriever(search_kwargs={"k": 10})
# 2. 我使用了本地离线模型(主要是这比较快)
reranker = LocalReranker(model_name="./models/BAAI--bge-small-zh-v1.5")
def format_docs(docs: List) -> str:
formatted = []
for i, doc in enumerate(docs, 1):
score = doc.metadata.get("rerank_score", 0.0)
formatted.append(f"[参考资料 {i}] (置信度得分: {score:.4f}):\n{doc.page_content}")
return "\n\n".join(formatted)
@tool
async def search_knowledge_base(query: str) -> str:
raw_docs = await retriever.ainvoke(query)
if not raw_docs:
return "知识库中未检索到相关内容。"
reranked_docs = reranker.rerank(query, raw_docs, top_n=2)
if reranked_docs[0].metadata["rerank_score"] < -2.0: # bge-reranker 得分区间
return "知识库检索结果置信度过低,未匹配到有效文档。"
return format_docs(reranked_docs)上面我拿到了RAG(向量数据库),接下来,我又写了一个工具类,获取天气情况(就是Agent的用法):
import httpx
from langchain_core.tools import tool
@tool
async def get_weather(city: str) -> str:
# 使用免费无需 API Key 的 wttr.in API(支持中文与 JSON 返回)
url = f"https://wttr.in/{city}?format=j1&lang=zh"
headers = {
"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36"
}
try:
# 设置 8 秒超时
async with httpx.AsyncClient(timeout=8.0) as client:
response = await client.get(url, headers=headers)
if response.status_code != 200:
return f"天气服务接口异常,HTTP 状态码: {response.status_code}"
data = response.json()
# 解析 API 返回的 JSON 字段
current = data.get("current_condition", [{}])[0]
if not current:
return f"未查询到 {city} 的有效天气数据,请检查城市名称输入是否正确。"
temp_c = current.get("temp_C", "未知")
humidity = current.get("humidity", "未知")
wind_speed = current.get("windspeedKmph", "未知")
lang_zh = current.get("lang_zh", [{}])
weather_desc = lang_zh[0].get("value") if lang_zh else None
if not weather_desc:
weather_desc = current.get("weatherDesc", [{}])[0].get("value", "未知")
return f"【{city}】实时天气信息:{weather_desc},当前气温 {temp_c}°C,相对湿度 {humidity}%,风速 {wind_speed} km/h。"
except httpx.TimeoutException:
return f"查询失败:请求 {city} 的远程天气接口超时,请稍后再试。"
except httpx.RequestError as e:
return f"网络异常:无法连接到远程天气 API 服务 ({str(e)})。"
except Exception as e:
return f"解析天气数据时发生未知错误: {str(e)}"准备好了,这两个工具后,就可以开始结合整体的fastapi完成下面的操作了:
1、环境准备:
本地大模型环境和缓存环境
Ollma qwen 2.5:3b +Redis
2、代码准备:
接口层:
发送请求并获取流式输出结果:
import json
from fastapi import APIRouter, HTTPException, status
from fastapi.responses import StreamingResponse
from pydantic import BaseModel, Field
from app.services.agent_service import AgentService
# 定义路由组
router = APIRouter(prefix="/agent", tags=["Agent 智能助手"])
# 初始化单例 Agent 服务
agent_service = AgentService(redis_url="redis://localhost:6379/0")
# 显式校验 message 和 session_id
class AgentRequest(BaseModel):
message: str = Field(..., min_length=1, description="用户发送的问题文本", example="今天天气怎么样?")
session_id: str = Field(..., min_length=1, description="唯一会话ID,用于隔离与持久化历史记忆", example="sess_123456")
@router.post("/chat/stream", summary="Agent SSE 流式对话接口")
async def agent_chat_stream_endpoint(request: AgentRequest):
try:
return StreamingResponse(
agent_service.get_stream_response(request.message, request.session_id),
media_type="text/event-stream"
)
except Exception as e:
raise HTTPException(
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
detail=f"Agent 流式服务内部异常: {str(e)}"
)
@router.post("/chat", summary="Agent 普通同步对话接口")
async def agent_chat_sync_endpoint(request: AgentRequest):
try:
full_content = ""
# 复用异步流生成器,在内部拼接所有 chunk 结果
async for chunk in agent_service.get_stream_response(request.message, request.session_id):
if chunk.startswith("data: "):
data_str = chunk.replace("data: ", "").strip()
if data_str and data_str != "[DONE]":
try:
payload = json.loads(data_str)
if "content" in payload:
full_content += payload["content"]
except json.JSONDecodeError:
continue
return {"response": full_content, "session_id": request.session_id}
except Exception as e:
raise HTTPException(
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
detail=f"Agent 同步服务内部异常: {str(e)}"
)服务层:
(获取工具列表,判断是否调用工具,或直接使用大模型)
import json
from typing import AsyncGenerator, List
from langchain_community.chat_message_histories import RedisChatMessageHistory
from langchain_core.messages import SystemMessage, HumanMessage, ToolMessage, AIMessage, BaseMessage
from app.core.llm import get_llm
from app.tools import ALL_TOOLS, TOOLS_MAP
class AgentService:
def __init__(self, redis_url: str = "redis://localhost:6379/0"):
self.llm = get_llm()
self.tools_map = TOOLS_MAP
self.redis_url = redis_url
self.llm_with_tools = self.llm.bind_tools(ALL_TOOLS)
def _get_history(self, session_id: str) -> RedisChatMessageHistory:
return RedisChatMessageHistory(
session_id=session_id,
url=self.redis_url,
ttl=86400
)
async def get_stream_response(self, question: str, session_id: str) -> AsyncGenerator[str, None]:
# 1.Redis
history_store = self._get_history(session_id)
# 2.Prompt
system_prompt = SystemMessage(
content=(
"你是一个严谨且高效的智能助手。\n"
"规则要求:\n"
"1. 如果调用的工具返回了错误或‘未收录’等提示,必须如实告知用户无法提供该信息,绝对严禁凭空捏造数据!\n"
"2. 严格基于工具返回的真实结果进行回答。"
)
)
saved_messages: List[BaseMessage] = history_store.messages
current_human_msg = HumanMessage(content=question)
messages = [system_prompt] + saved_messages + [current_human_msg]
new_messages_to_save: List[BaseMessage] = [current_human_msg]
first_response = await self.llm_with_tools.ainvoke(messages)
messages.append(first_response)
new_messages_to_save.append(first_response)
# 触发Tool
if first_response.tool_calls:
for tool_call in first_response.tool_calls:
tool_name = tool_call["name"]
tool_args = tool_call["args"]
tool_func = self.tools_map.get(tool_name)
if tool_func:
# 使用异步,消除阻塞
tool_result = await tool_func.ainvoke(tool_args)
print(f"[Agent Log] 异步执行工具 [{tool_name}],输入: {tool_args},输出: {tool_result}")
tool_msg = ToolMessage(
content=str(tool_result),
tool_call_id=tool_call["id"]
)
messages.append(tool_msg)
new_messages_to_save.append(tool_msg)
full_ai_content = ""
async for chunk in self.llm_with_tools.astream(messages):
if chunk.content:
full_ai_content += chunk.content
payload = json.dumps({"content": chunk.content}, ensure_ascii=False)
yield f"data: {payload}\n\n"
if full_ai_content:
new_messages_to_save.append(AIMessage(content=full_ai_content))
else:
# 普通场景
if first_response.content:
payload = json.dumps({"content": first_response.content}, ensure_ascii=False)
yield f"data: {payload}\n\n"
# 持久化至 Redis
history_store.add_messages(new_messages_to_save)Reranker 工具类
from typing import List
from langchain_core.documents import Document
from sentence_transformers import CrossEncoder
class LocalReranker:
def __init__(self, model_name: str = "BAAI/bge-reranker-base"):
# 加载模型
self.model = CrossEncoder(model_name)
def rerank(self, query: str, docs: List[Document], top_n: int = 2) -> List[Document]:
if not docs:
return []
pairs = [[query, doc.page_content] for doc in docs]
scores = self.model.predict(pairs)
for doc, score in zip(docs, scores):
doc.metadata["rerank_score"] = float(score)
sorted_docs = sorted(docs, key=lambda x: x.metadata["rerank_score"], reverse=True)
return sorted_docs[:top_n]工具层:
from app.tools.rag_tool import search_knowledge_base
from app.tools.weather_tool import get_weather
# 工具列表
ALL_TOOLS = [
search_knowledge_base,
get_weather,
]
# 导出工具字典
TOOLS_MAP = {t.name: t for t in ALL_TOOLS}这里调用的rag_tool 和weather_tool都在文章前言章里写了,这里直接调用就行。
许可协议:
CC BY 4.0