AI 智能设备故障诊断 Agent — Java 后端工程师转型 AI 工程方向项目学习指南
适用对象:Java 后端工程师转型 AI 工程方向
学习目标:理解项目每一行代码的设计意图,掌握 RAG/Agent 核心概念,面试自信讲解。个人学习记录,供参考。
目录
第一章:项目总览与技术栈
1.1 项目是什么
这是一个面向 IoT 设备运维场景的 AI 智能诊断系统。部署在 智能饮吧设备(自动售卖饮品的机器)上,运维人员描述故障现象后,系统会:
- 自动查询设备各硬件模块的实时状态
- 读取设备运行日志,定位错误信息
- 检索知识库(设备手册 + 历史故障案例),找到类似问题的解决方案
- 综合分析后,生成结构化的诊断报告和修复建议
一句话总结:用 RAG + Agent 技术让 AI 成为一个有知识、有工具、能自主推理的"AI 运维工程师"。
1.2 项目架构图(文字版)
┌─────────────────────────────────────────────────────────────┐
│ 用户层 │
│ ┌──────────────┐ ┌──────────────────────────────────┐ │
│ │ CLI (main.py)│ │ React 前端 (frontend/src/App.tsx)│ │
│ └──────┬───────┘ └──────────────┬───────────────────┘ │
│ │ │ │
├─────────┼──────────────────────────┼─────────────────────────┤
│ │ API 层 │ │
│ │ ┌──────────────────────┐│ │
│ │ │ FastAPI (app.py) ││ │
│ │ │ ├── /api/chat ││ SSE 流式输出 │
│ │ │ ├── /api/diagnose ││ 会话管理 (Redis) │
│ │ │ ├── /api/rag/query ││ 限流 (slowapi) │
│ │ │ └── /api/session ││ │
│ │ └──────────┬───────────┘│ │
│ │ │ │ │
├─────────┼──────────────┼────────────┼─────────────────────────┤
│ │ Agent 层 │ │
│ ┌─────┴──────────────┴────────────┴──────────┐ │
│ │ Supervisor (supervisor.py) │ │
│ │ ┌─────────────┬──────────────┬───────────┐ │ │
│ │ │ 诊断 Agent │ 维修 Agent │ 监控 Agent│ │ │
│ │ │ (6个工具) │ (2个工具) │ (3个工具) │ │ │
│ │ └──────┬──────┴──────┬───────┴─────┬─────┘ │ │
│ └─────────┼─────────────┼─────────────┼───────┘ │
│ │ │ │ │
├─────────────┼─────────────┼─────────────┼─────────────────────┤
│ 工具层 │ │ │
│ ┌─────────────┐ ┌─────┴────┐ ┌─────┴──────────┐ │
│ │knowledge │ │log_reader│ │device_status │ │
│ │_search │ │(SSH/本地) │ │(HTTP→middleware)│ │
│ ├─────────────┤ ├──────────┤ ├────────────────┤ │
│ │diagnosis │ │device │ │firmware_check │ │
│ │_report │ │_restart │ │ │ │
│ └──────┬──────┘ └──────────┘ └────────────────┘ │
│ │ │
├──────────┼───────────────────────────────────────────────────┤
│ RAG 层 │
│ ┌──────┴──────────────────────────────────────┐ │
│ │ hybrid_search.py (混合检索) │ │
│ │ ┌──────────────┐ ┌──────────────────┐ │ │
│ │ │ 向量语义检索 │ │ BM25 关键词检索 │ │ │
│ │ │(vectorstore) │ │(rank-bm25) │ │ │
│ │ └──────┬───────┘ └────────┬─────────┘ │ │
│ │ └─────┬─────────────┘ │ │
│ │ │ RRF 融合 │ │
│ │ ┌─────┴─────┐ │ │
│ │ │ CrossEncoder│ 语义重排序 │ │
│ │ │ (reranker) │ │ │
│ │ └─────────────┘ │ │
│ └─────────────────────────────────────────────┘ │
│ │
├──────────────────────────────────────────────────────────────┤
│ 存储层 │
│ ┌────────────┐ ┌──────────┐ ┌───────────────────┐ │
│ │ ChromaDB │ │ Milvus │ │ Redis │ │
│ │ (开发环境) │ │ (生产环境) │ │ (会话持久化) │ │
│ └────────────┘ └──────────┘ └───────────────────┘ │
└──────────────────────────────────────────────────────────────┘
1.3 技术栈对照表(Python ↔ Java)
| 功能领域 | Python 生态(本项目) | Java 生态(你熟悉的) | 说明 |
|---|---|---|---|
| Web 框架 | FastAPI | Spring Boot | 路由、中间件、依赖注入概念相似 |
| 配置管理 | python-dotenv + 自定义 Settings 类 | @ConfigurationProperties + application.yml | 都是从配置文件加载到对象 |
| 数据模型 | Pydantic BaseModel | Spring 的 DTO / Jackson | 序列化、校验、类型约束 |
| HTTP 客户端 | httpx | RestTemplate / WebClient | 调用外部 API |
| ORM / 数据访问 | LangChain Document 对象 | MyBatis / JPA Entity | 数据的统一抽象 |
| 测试框架 | pytest | JUnit 5 | 单元测试 + fixture |
| 容器化 | Docker + docker-compose | Docker + docker-compose | 完全一致 |
| 流式输出 | SSE (StreamingResponse) | SSE (SseEmitter) / WebFlux | 服务端推送 |
| 限流 | slowapi | Spring Cloud Gateway / Sentinel | 限制请求频率 |
| 缓存 / 会话 | Redis + 内存降级 | Redis + Caffeine | 多级缓存策略 |
| AI 框架 | LangChain + LangGraph | Spring AI (较新) | AI 应用开发框架 |
| 向量数据库 | ChromaDB / Milvus | Milvus / Elasticsearch KNN | 存储和检索向量 |
| 服务降级 | with_fallbacks() | Resilience4j / Sentinel | 主服务挂了切备用 |
1.4 完整请求调用链路
以用户发送 “制冰机报错 E03” 为例,从头到尾追踪整条调用链:
1. 用户在 React 前端输入 "制冰机报错 E03",点击发送
│
2. 前端 POST /api/diagnose/stream (SSE)
│ body: { question: "制冰机报错 E03", provider: "claude", use_multi_agent: false }
│
3. FastAPI routes.py: diagnose_stream() 接收请求
│ → 创建 diagnostic_agent (create_react_agent)
│ → agent.stream() 开始 ReAct 循环
│
4. Agent 第一步 (Thought): "用户报告制冰机 E03 错误,我需要先查设备状态"
│ → Action: query_device_status(module="制冰机")
│ → 工具调用 bar_middleware HTTP API: GET http://localhost:8003/ice/machine/status
│ → SSE 推送: { type: "tool_call", tool: "query_device_status" }
│ → SSE 推送: { type: "tool_result", result: "制冰机: 状态=error, 错误码=E03" }
│
5. Agent 第二步 (Thought): "E03 错误,看看日志有什么线索"
│ → Action: fetch_device_logs(log_type="middleware", keyword="E03")
│ → 读取 /home/*/bar_middleware/logs/app.log,过滤含 E03 的行
│ → SSE 推送工具调用和结果
│
6. Agent 第三步 (Thought): "去知识库查 E03 的含义和处理方法"
│ → Action: search_knowledge_base(query="制冰机 E03 错误码")
│ → enhanced_search() 执行混合检索:
│ ├── 向量检索 (ChromaDB) → top 20
│ ├── BM25 关键词检索 → top 20
│ ├── RRF 融合 → top 20
│ └── CrossEncoder 重排序 → top 5
│ → 返回匹配的文档片段
│ → SSE 推送
│
7. Agent 第四步 (Thought): "信息够了,生成诊断报告"
│ → Action: generate_diagnosis_report(...)
│ → 格式化 Markdown 报告
│ → SSE 推送
│
8. Agent 最终回复: 综合诊断分析 + 修复建议
│ → SSE 推送: { type: "answer", content: "根据分析..." }
│ → SSE 推送: { type: "done" }
│
9. 前端接收 SSE 事件流,实时渲染工具调用过程和最终回答
1.5 源码目录结构
ai-diagnostic-agent/
├── src/ # 所有应用源码
│ ├── main.py # CLI 入口
│ ├── config/
│ │ └── settings.py # 配置管理(环境变量加载)
│ ├── llm/
│ │ └── client.py # LLM 客户端(多模型 + 降级 + 流式)
│ ├── rag/ # RAG 知识库系统
│ │ ├── loader.py # 文档加载器
│ │ ├── splitter.py # 文本切分器
│ │ ├── vectorstore.py # 向量存储(ChromaDB/Milvus)
│ │ ├── hybrid_search.py # 混合检索(BM25+向量+RRF+Rerank)
│ │ ├── reranker.py # CrossEncoder 重排序
│ │ └── chain.py # RAG Chain(检索+生成)
│ ├── agent/ # Agent 编排
│ │ ├── diagnostic_agent.py # 诊断 Agent(主 Agent,6个工具)
│ │ ├── repair_agent.py # 维修 Agent(2个工具)
│ │ ├── monitor_agent.py # 监控 Agent(3个工具)
│ │ └── supervisor.py # 多 Agent 编排器(路由+链式调用)
│ ├── tools/ # Agent 工具集
│ │ ├── knowledge_search.py # 知识库检索
│ │ ├── log_reader.py # 日志读取(SSH/本地)
│ │ ├── device_status.py # 设备状态查询(HTTP API)
│ │ ├── device_restart.py # 设备服务重启(SSH 白名单)
│ │ ├── firmware_check.py # 固件版本查询
│ │ └── diagnosis_report.py # 诊断报告生成
│ ├── api/ # FastAPI 服务
│ │ ├── app.py # 应用入口(CORS、限流)
│ │ ├── routes.py # 路由(聊天、诊断、RAG、会话管理)
│ │ ├── models.py # 请求/响应模型 (Pydantic)
│ │ └── limiter.py # 限流器单例
│ └── utils/
│ └── session_store.py # 会话管理(Redis+内存双后端)
├── scripts/
│ └── build_knowledge_base.py # 知识库构建脚本
├── knowledge-base/ # RAG 源文档
│ ├── device-manuals/ # 14 篇设备手册 (Markdown)
│ └── fault-cases/ # 5 篇故障案例 (Markdown)
├── frontend/ # React 前端
│ └── src/App.tsx # 单页面应用
├── tests/ # 17 个测试文件
├── evaluation/ # RAG 评估工具
├── Dockerfile # 容器镜像
├── docker-compose.yml # 容器编排(agent+Redis+Milvus)
└── requirements.txt # Python 依赖
第二章:基础设施 — 配置与 LLM 客户端
2.1 配置管理(settings.py)
文件路径:
src/config/settings.py
代码逐行讲解
# 第 1-5 行:加载 .env 文件
"""配置管理 - 从 .env 加载所有配置"""
import os
from dotenv import load_dotenv
load_dotenv() # 从项目根目录的 .env 文件加载环境变量到 os.environ
为什么用 .env? Python 没有 Spring Boot 那样内建的 application.yml,社区标准做法是用 .env 文件 + python-dotenv 库。load_dotenv() 把 .env 中的 KEY=VALUE 对注入到 os.environ。
Java 类比:就像 Spring Boot 的 @PropertySource("classpath:application.properties") + @Value("${key}") 的组合。
# 第 7-11 行:清除系统代理
for key in ("HTTP_PROXY", "HTTPS_PROXY", "http_proxy", "https_proxy"):
os.environ.pop(key, None)
os.environ["NO_PROXY"] = "*"
为什么要清除代理? 开发环境经常开 Clash 等代理工具,但调用 LLM API 中转站时,代理会导致 SSL 证书校验失败。统一在配置模块入口处清除一次,其他模块不用管这个问题。
# 第 14-56 行:Settings 类
class Settings:
ANTHROPIC_API_KEY: str = os.getenv("ANTHROPIC_API_KEY", "")
OPENAI_API_KEY: str = os.getenv("OPENAI_API_KEY", "")
# ... 更多配置项
LLM_FALLBACK_ENABLED: bool = os.getenv("LLM_FALLBACK_ENABLED", "true").lower() == "true"
REDIS_URL: str = os.getenv("REDIS_URL", "redis://localhost:6379/0")
SESSION_TTL: int = int(os.getenv("SESSION_TTL", "3600"))
settings = Settings() # 模块级单例
设计选择:用类属性 + os.getenv() 直接在类定义阶段赋值,而不是用 Pydantic Settings。
Java 类比:
// 类似 Spring 的 @ConfigurationProperties
@ConfigurationProperties(prefix = "app")
public class Settings {
private String anthropicApiKey;
private String openaiApiKey;
private boolean llmFallbackEnabled = true;
// ...
}
// 在 Spring 中会自动从 application.yml 注入值
// Python 的 settings = Settings() 类似 Spring 注册 Bean
// 全局唯一实例,其他模块 from src.config.settings import settings 直接引用
关键配置项说明:
| 配置项 | 作用 | 默认值 |
|---|---|---|
DEFAULT_LLM_PROVIDER | 默认用哪个 LLM | "claude" |
LOG_READ_MODE | 日志读取方式 | "local"(本机)/ "ssh"(远程) |
VECTOR_DB_TYPE | 向量数据库类型 | "chromadb"(开发)/ "milvus"(生产) |
LLM_FALLBACK_ENABLED | 是否开启 LLM 自动降级 | true |
MIDDLEWARE_BASE_URL | 设备中间件地址 | http://localhost:8003 |
SESSION_TTL | 会话过期秒数 | 3600(1小时) |
2.2 LLM 客户端(client.py)
文件路径:
src/llm/client.py
这个文件是整个项目与 LLM 交互的唯一入口。理解它就掌握了"项目怎么调用大模型"。
2.2.1 工厂模式创建 LLM 实例
# 第 24-46 行:创建单个 provider 的 LLM
def _create_single_llm(provider: str):
if provider == "claude":
kwargs = {
"model": settings.CLAUDE_MODEL,
"api_key": settings.ANTHROPIC_API_KEY,
"max_tokens": 4096,
}
if settings.ANTHROPIC_BASE_URL:
kwargs["base_url"] = settings.ANTHROPIC_BASE_URL # 支持中转站
return ChatAnthropic(**kwargs) # 返回 LangChain 的 Claude 封装
elif provider == "openai":
kwargs = { ... }
return ChatOpenAI(**kwargs)
为什么用 **kwargs? 这是 Python 的字典展开语法。因为 base_url 只在配置了中转站时才需要传,用 kwargs 动态构建参数字典,避免写很多 if-else。
Java 类比:类似 Builder 模式:
// Python 的 **kwargs 相当于 Java 的 Builder 模式
ChatAnthropic.builder()
.model("claude-sonnet-4-20250514")
.apiKey("sk-xxx")
.maxTokens(4096)
.baseUrl("https://relay.example.com") // 可选
.build();
2.2.2 自动降级(with_fallbacks)
# 第 61-86 行:带降级的 LLM 创建
def create_llm(provider: str = None):
provider = provider or settings.DEFAULT_LLM_PROVIDER # 默认读配置
primary = _create_single_llm(provider) # 创建主 LLM
if not settings.LLM_FALLBACK_ENABLED:
return primary # 未启用降级,直接返回
fallback_provider = _get_fallback_provider(provider) # claude ↔ openai
if not fallback_provider:
return primary # 备用 provider 没配置 API key
fallback = _create_single_llm(fallback_provider)
return primary.with_fallbacks([fallback]) # 🔑 关键:LangChain 的降级链
with_fallbacks() 的工作原理:LangChain 内置的降级机制。调用 primary 失败(异常/超时)时,自动尝试 fallback。不需要你手动写 try-catch-retry。
Java 类比:
// 类似 Resilience4j 的 CircuitBreaker + Fallback
@CircuitBreaker(name = "llm", fallbackMethod = "callOpenAI")
public String callClaude(String prompt) { ... }
public String callOpenAI(String prompt, Throwable t) { ... }
面试话术:“我们项目支持 Claude 和 OpenAI 双模型,通过 LangChain 的 with_fallbacks() 实现自动降级。当主模型 API 不可用时(比如限流或网络问题),会无感切换到备用模型,保证服务可用性。这类似于 Java 中 Resilience4j 的 CircuitBreaker + Fallback 模式。”
2.2.3 流式输出(yield / generator)
# 第 144-178 行:DiagnosticChat 的流式对话
def stream_chat(self, user_input: str):
self.history.append(HumanMessage(content=user_input))
full_response = ""
try:
for chunk in self.llm.stream(self.history): # LLM 逐 token 返回
token = chunk.content
if isinstance(token, str):
full_response += token
yield token # 🔑 yield 逐个返回 token
except Exception:
self.history.pop() # 回滚:移除刚加的用户消息
result = self.chat(user_input) # 降级:流式失败 → 非流式
yield result
Python 的 yield 是什么? 函数中有 yield 就变成了生成器函数(generator function)。调用它不会立即执行,而是返回一个迭代器。每次 next() 或 for ... in 时执行到 yield 暂停,返回值,下次继续。
Java 类比:
// Python 的 yield 类比 Java Iterator + 惰性计算
// stream_chat() 就像一个 Iterator<String>
Iterator<String> streamChat(String input) {
return new Iterator<>() {
// 每次 next() 返回下一个 token
public String next() { ... }
};
}
// 或者用 Java Stream
Stream<String> streamChat(String input) {
return Stream.generate(() -> getNextToken())
.takeWhile(token -> token != null);
}
为什么用流式? 用户体验好。LLM 生成一段话可能要 5-10 秒,流式输出可以每隔几十毫秒就显示新的文字,而不是等全部生成完再一次性显示。
2.2.4 DiagnosticChat 会话管理
# 第 121-187 行
class DiagnosticChat:
def __init__(self, provider: str = None):
self.llm = create_llm(provider) # LLM 实例
self.history: list = [SystemMessage(content=DIAGNOSTIC_SYSTEM_PROMPT)] # 消息历史
self.provider_name = provider or settings.DEFAULT_LLM_PROVIDER
def chat(self, user_input: str) -> str:
self.history.append(HumanMessage(content=user_input)) # 加入用户消息
response = self.llm.invoke(self.history) # 把整个历史发给 LLM
text = extract_text(response.content) # 提取纯文本
self.history.append(AIMessage(content=text)) # 加入 AI 回复
return text
核心概念:LLM 的多轮对话。LLM 本身是无状态的,不记得之前说过什么。要实现多轮对话,必须每次把完整的消息历史都发送给 LLM。self.history 就是这个消息列表。
三种消息类型:
SystemMessage:系统提示词,定义 AI 的角色和行为规范HumanMessage:用户说的话AIMessage:AI 的回复
Java 类比:
// 类似 Spring AI 的 ChatClient + MessageHistory
List<Message> history = new ArrayList<>();
history.add(new SystemMessage("你是..."));
public String chat(String input) {
history.add(new UserMessage(input));
ChatResponse response = chatClient.call(new Prompt(history));
String text = response.getResult().getOutput().getContent();
history.add(new AssistantMessage(text));
return text;
}
2.2.5 extract_text 辅助函数
# 第 89-101 行
def extract_text(content) -> str:
if isinstance(content, str):
return content
if isinstance(content, list):
# extended thinking 模式返回 [{"type":"thinking","thinking":"..."}, {"type":"text","text":"真正回答"}]
parts = [block["text"] for block in content if isinstance(block, dict) and block.get("type") == "text"]
return "".join(parts).strip()
return str(content)
为什么需要这个函数? Claude 模型的 extended thinking 功能会返回"思考过程 + 最终回答"的混合格式。API 中转站有时也会改变返回格式。这个函数做了格式兼容,确保不管上游返回什么格式,我们都能提取出纯文本回答。
2.3 CLI 入口(main.py)
文件路径:
src/main.py
# 第 26-126 行:REPL 主循环
def main():
mode = "chat" if "--chat" in sys.argv else "agent" # 命令行参数决定默认模式
while True: # REPL 循环(Read-Eval-Print Loop)
user_input = input("\n🧑 你: ") # Read: 读取用户输入
if user_input.startswith("/"): # 命令路由
# /agent, /chat, /switch, /clear, /quit
...
continue
if mode == "agent": # Eval: Agent 模式
result = run_diagnosis(user_input, provider=provider, verbose=True)
print(result) # Print: 输出结果
else: # 普通对话模式
for token in chat.stream_chat(user_input):
print(token, end="", flush=True)
REPL 模式是命令行工具的经典模式:不断"读取输入 → 执行 → 打印结果"循环。Python 自身的交互式解释器就是一个 REPL。
Java 类比:类似写一个 Scanner 循环的控制台应用:
Scanner scanner = new Scanner(System.in);
while (scanner.hasNextLine()) {
String input = scanner.nextLine();
if (input.startsWith("/")) { /* 命令处理 */ }
else { System.out.println(process(input)); }
}
2.4 Python 核心语法速查
对 Java 工程师来说最需要理解的 Python 语法:
| Python 语法 | 含义 | Java 类比 |
|---|---|---|
**kwargs | 字典展开为关键字参数 | Builder 模式 |
yield | 生成器,惰性产出值 | Iterator<T> / Stream<T> |
def f(x: str) -> int | 类型注解(非强制) | 方法签名 |
str | None | 联合类型 | Optional<String> |
from X import Y | 导入 | import X.Y |
os.getenv("K", "default") | 环境变量 | System.getenv("K") |
isinstance(x, str) | 类型检查 | x instanceof String |
list[tuple[Doc, float]] | 泛型类型注解 | List<Tuple2<Doc, Float>> |
@decorator | 装饰器 | @Annotation + AOP |
with open() as f: | 资源管理 | try-with-resources |
面试高频问题 & 参考回答
Q1:你的项目怎么支持多模型切换的?
“通过工厂模式 + 配置驱动。
create_llm()函数根据provider参数创建不同的 LLM 实例,Claude 用ChatAnthropic,OpenAI 用ChatOpenAI。它们都实现 LangChain 的ChatModel抽象接口,所以上层代码不关心底层是哪个模型。切换模型只需改一个配置项。”“同时启用了自动降级——通过 LangChain 的
with_fallbacks()把主模型和备用模型串联起来,主模型 API 不可用时自动切换,类似 Java 中 Resilience4j 的 Fallback 模式。”
Q2:你的流式输出是怎么实现的?
“两层流式。后端用 Python 的 generator(yield)逐 token 产出文本,通过 FastAPI 的
StreamingResponse以 SSE(Server-Sent Events)格式推送到前端。前端用ReadableStreamAPI 实时读取事件流,每收到一个 token 就追加到页面上。”“这样用户不用等 5-10 秒才看到回答,而是几乎实时看到文字逐步出现,体验类似 ChatGPT。”
Q3:为什么要清除系统代理?
“开发环境中经常运行 Clash 等代理工具。我们调用 LLM API 时使用的是中转站(base_url 重定向),如果请求走了本地代理,代理的 SSL 证书与中转站不匹配会导致 SSL 验证失败。所以在配置模块加载时统一清除
HTTP_PROXY/HTTPS_PROXY环境变量,从根源上避免这个问题。”
第三章:RAG 知识库系统
什么是 RAG?
RAG(Retrieval-Augmented Generation) = 检索增强生成。
核心思想:LLM 的知识截止到训练时间,不知道你公司的内部文档。RAG 的做法是——先从知识库中检索出与用户问题相关的文档片段,把这些片段塞入 prompt作为上下文,再让 LLM 生成回答。
传统 LLM: 问题 → LLM → 回答(可能胡说八道)
RAG: 问题 → 检索相关文档 → [文档+问题] → LLM → 回答(基于事实)
Java 类比:想象一个开卷考试——考生(LLM)不需要背诵所有知识,但在回答问题前,先翻阅参考书(知识库),然后基于参考书的内容来作答。
RAG 的完整流程(本项目实现)
┌─── 离线阶段(构建知识库)───────────────────────────┐
│ │
│ 文档 (.md/.pdf/.txt) │
│ │ │
│ ① 加载 (loader.py) │
│ │ → 转为 Document 对象(content + metadata) │
│ │ │
│ ② 切分 (splitter.py) │
│ │ → 按标题 + 字符切分为小块 │
│ │ │
│ ③ Embedding + 存储 (vectorstore.py) │
│ │ → 文本向量化 → 存入 ChromaDB/Milvus │
│ │
└──────────────────────────────────────────────────────┘
┌─── 在线阶段(用户检索)──────────────────────────────┐
│ │
│ 用户查询: "制冰机 E03 错误码" │
│ │ │
│ ④ 混合检索 (hybrid_search.py) │
│ │ ├── 向量语义检索(找语义相似的) │
│ │ ├── BM25 关键词检索(找词匹配的) │
│ │ ├── RRF 融合(合并两路排序结果) │
│ │ └── CrossEncoder 重排序(精排) │
│ │ │
│ ⑤ 生成回答 (chain.py) │
│ │ → 检索结果注入 prompt → LLM 生成回答 │
│ │
└──────────────────────────────────────────────────────┘
3.1 文档加载(loader.py)
文件路径:
src/rag/loader.py
# 第 14-31 行:加载 Markdown 文件
def load_markdown(file_path: str) -> list[Document]:
path = Path(file_path)
with open(path, "r", encoding="utf-8") as f:
content = f.read()
return [Document(
page_content=content, # 文档全文
metadata={
"source": str(path), # 完整路径(溯源用)
"filename": path.name, # 文件名(展示用)
"file_type": "markdown", # 文件类型(切分策略依赖这个字段)
}
)]
Document 对象是 LangChain 的核心数据结构,两个字段:
page_content:文本内容metadata:元数据字典(来源、页码、标题等)
为什么需要 metadata? 检索到文档片段后,需要告诉用户"这个信息来自哪个文件哪个章节"。metadata 就是为了做引用溯源。
# 第 62-84 行:递归加载目录
def load_directory(directory: str) -> list[Document]:
supported = {".md": load_markdown, ".pdf": load_pdf, ".txt": load_markdown}
for root, _, files in os.walk(directory): # 递归遍历目录
for file in files:
ext = Path(file).suffix.lower()
if ext in supported:
docs = supported[ext](file_path) # 根据扩展名选择 loader
documents.extend(docs)
设计亮点:用字典映射 {扩展名: 加载函数} 代替 if-else,新增文件格式只需加一行映射。
Java 类比:
// 类似策略模式
Map<String, Function<String, List<Document>>> loaders = Map.of(
".md", this::loadMarkdown,
".pdf", this::loadPdf,
".txt", this::loadMarkdown
);
3.2 文本切分(splitter.py)
文件路径:
src/rag/splitter.py
为什么需要切分?
- Embedding 模型有输入长度限制(
all-MiniLM-L6-v2最大 512 tokens) - 短文本的 Embedding 质量更高——一段话讲一件事,向量表示更精确
- 检索粒度更细——命中整本手册 vs 命中某一段描述,后者更有用
两级切分策略
# 第 52-95 行:核心切分逻辑
def split_documents(documents, chunk_size=1000, chunk_overlap=200):
text_splitter = create_text_splitter(chunk_size, chunk_overlap)
md_splitter = create_markdown_splitter()
for doc in documents:
if doc.metadata.get("file_type") == "markdown":
# 第一级:按 Markdown 标题切分
md_chunks = md_splitter.split_text(doc.page_content)
for md_chunk in md_chunks:
md_chunk.metadata.update(doc.metadata) # 标题信息 + 原始 metadata
# 第二级:对过长的 section 再按字符切分
final_chunks = text_splitter.split_documents(md_chunks)
else:
# 非 Markdown:直接按字符切分
chunks = text_splitter.split_documents([doc])
为什么 Markdown 要两级切分?
-
第一级(标题切分):按
#、##、###标题分段。好处是每个 chunk 都是完整的一个 section,语义边界清晰。标题信息(h1、h2、h3)保存到 metadata,检索时可以显示"来自哪个章节"。 -
第二级(字符切分):标题切分后可能有些 section 很长(比如 3000 字的章节),超过 chunk_size。再用
RecursiveCharacterTextSplitter按字符切分。
# 第 28-33 行:递归字符切分器的分隔符优先级
RecursiveCharacterTextSplitter(
chunk_size=chunk_size,
chunk_overlap=chunk_overlap,
separators=["\n\n", "\n", "。", ";", " ", ""], # 优先按段落分,再按句子分
length_function=len,
)
chunk_size=1000, chunk_overlap=200 怎么选的?
chunk_size=1000:一个 chunk 约 500-700 个中文字。太小(如 200)会丢失上下文;太大(如 3000)检索不精确。1000 是中文技术文档的经验值。chunk_overlap=200:相邻 chunk 重叠 200 字符。防止关键信息刚好在切分边界被切断。200 ≈ chunk_size 的 20%,是常用比例。
注意:实际构建知识库时用的参数是
chunk_size=800, chunk_overlap=150(scripts/build_knowledge_base.py第 48 行),比默认值略小,更适合知识库文档的平均长度。
Java 类比:
// 没有直接对应,但概念类似文本处理
// 想象把一本书切成"笔记卡片":
// - 先按章节切(标题切分)
// - 每章如果太长,再按段落切(字符切分)
// - 每张卡片约 500-700 字
// - 相邻卡片有 100-150 字重叠,避免切断关键信息
3.3 向量存储(vectorstore.py)
文件路径:
src/rag/vectorstore.py
Embedding 是什么?
Embedding = 把文本转成高维向量(一组浮点数)。语义相似的文本,向量距离近。
"制冰机故障" → [0.12, -0.34, 0.56, ..., 0.78] (384 维)
"冰机出问题了" → [0.11, -0.33, 0.55, ..., 0.77] (很接近!)
"咖啡机正常" → [-0.45, 0.21, -0.87, ..., 0.12] (距离远)
Embedding 模型选择
# 第 24-47 行
def create_embedding(use_local: bool = True):
if use_local:
from langchain_community.embeddings import HuggingFaceEmbeddings
return HuggingFaceEmbeddings(
model_name="all-MiniLM-L6-v2", # 本地小模型,384 维
model_kwargs={"device": "cpu"}, # CPU 运行
)
# OpenAI API
from langchain_openai import OpenAIEmbeddings
return OpenAIEmbeddings(model="text-embedding-3-small")
| 方案 | 模型 | 维度 | 优点 | 缺点 |
|---|---|---|---|---|
| 本地(默认) | all-MiniLM-L6-v2 | 384 | 免费、离线、无延迟 | 中文效果一般 |
| API | text-embedding-3-small | 1536 | 效果好、多语言 | 收费、需网络 |
为什么默认用本地模型? 项目部署在 IoT 设备(树莓派/工控机)上,网络环境不稳定。本地模型保证离线可用。
ChromaDB vs Milvus 双后端
# 第 99-122 行:根据配置自动选择向量数据库
def create_vector_store(documents=None, use_local_embedding=True):
embedding = create_embedding(use_local=use_local_embedding)
db_type = settings.VECTOR_DB_TYPE
if db_type == "milvus":
return _create_milvus_store(embedding, documents)
else:
return _create_chroma_store(embedding, documents)
| ChromaDB | Milvus | |
|---|---|---|
| 定位 | 轻量级,适合开发 | 生产级,适合部署 |
| 部署 | 纯 Python,进程内 | 独立服务,需 Docker |
| 性能 | 百万级以内够用 | 十亿级向量 |
| 本项目 | VECTOR_DB_TYPE=chromadb | VECTOR_DB_TYPE=milvus |
Java 类比:类似 H2(开发)和 MySQL(生产)的切换。
线程安全单例 + 懒加载
# 第 20-22 行
_cached_vectorstore = None
_cache_lock = threading.Lock()
# 第 125-132 行:双重检查锁(DCL)
def _get_vectorstore():
global _cached_vectorstore
if _cached_vectorstore is None: # 第一次检查(无锁)
with _cache_lock: # 加锁
if _cached_vectorstore is None: # 第二次检查(有锁)
_cached_vectorstore = create_vector_store()
return _cached_vectorstore
这就是经典的 DCL(Double-Checked Locking)模式! Java 工程师一定见过:
// Java 版 DCL
private static volatile VectorStore instance;
private static final Object lock = new Object();
public static VectorStore getInstance() {
if (instance == null) {
synchronized (lock) {
if (instance == null) {
instance = new VectorStore();
}
}
}
return instance;
}
为什么需要懒加载? 向量存储初始化需要加载 Embedding 模型(约 100MB),比较慢。只在第一次实际使用时才初始化,避免启动时的不必要开销。
3.4 混合检索(hybrid_search.py)— 核心算法
文件路径:
src/rag/hybrid_search.py
这是整个 RAG 系统的核心模块,面试重点。
为什么需要混合检索?
单一检索方式各有短板:
| 方式 | 擅长 | 短板 | 举例 |
|---|---|---|---|
| 向量语义检索 | 理解同义词、换种说法也能找到 | 对精确关键词(错误码、型号)不敏感 | “冰机坏了” → 能找到 “制冰机故障” |
| BM25 关键词检索 | 精确匹配关键词、错误码 | 不理解语义 | “E03” → 精确找到含 E03 的文档 |
混合检索 = 两者结合,取长补短。
完整流程
# 第 150-194 行:enhanced_search 完整流程
def enhanced_search(query, k=5, use_rewrite=False, llm=None):
# 可选:查询改写
search_query = rewrite_query(query, llm) if use_rewrite and llm else query
# 1. 向量语义检索 (top 20)
vector_results = search_with_scores(search_query, k=20)
vector_docs = [doc for doc, _ in vector_results]
# 2. BM25 关键词检索 (top 20)
bm25_index = _get_bm25_index()
bm25_docs = bm25_index.search(search_query, k=20) if bm25_index else []
# 3. RRF 融合
fused_docs = rrf_fuse([vector_docs, bm25_docs], k=60, top_k=20)
# 4. CrossEncoder 重排序
reranked = rerank(search_query, fused_docs, top_k=k)
return reranked
BM25 关键词检索 + 中文分词
# 第 25-37 行:自定义中文分词
def _tokenize(text: str) -> list[str]:
"""中文分词(字符级 + 2-gram)"""
tokens = re.findall(r'[\u4e00-\u9fff]+|[a-zA-Z0-9]+', text.lower())
result = []
for token in tokens:
if re.match(r'[\u4e00-\u9fff]', token): # 中文
for i in range(len(token)):
result.append(token[i]) # 单字: "制", "冰", "机"
if i + 1 < len(token):
result.append(token[i:i+2]) # 2-gram: "制冰", "冰机"
else:
result.append(token) # 英文单词保持完整
return result
为什么用字符级 + 2-gram? 因为没有引入 jieba 等专业中文分词库(减少依赖)。单字 + 双字组合是一种简单但有效的中文分词替代方案:
"制冰机故障" → ["制", "制冰", "冰", "冰机", "机", "机故", "故", "故障", "障"]
BM25 通过 TF-IDF 原理,关键词在文档中出现频率越高、在全部文档中出现越少,得分越高。
RRF(Reciprocal Rank Fusion)融合算法
# 第 83-106 行
def rrf_fuse(ranked_lists, k=60, top_k=20):
"""
score(doc) = sum(1 / (k + rank_i)) for each ranked list i
"""
doc_scores = {}
for ranked_list in ranked_lists:
for rank, doc in enumerate(ranked_list, 1):
doc_key = hash((doc.page_content, doc.metadata.get("filename", "")))
if doc_key in doc_scores:
existing_doc, existing_score = doc_scores[doc_key]
doc_scores[doc_key] = (existing_doc, existing_score + 1.0 / (k + rank))
else:
doc_scores[doc_key] = (doc, 1.0 / (k + rank))
sorted_docs = sorted(doc_scores.values(), key=lambda x: x[1], reverse=True)
return [doc for doc, _ in sorted_docs[:top_k]]
RRF 算法解析:
假设一个文档 D 在向量检索排第 3、在 BM25 检索排第 1:
RRF_score(D) = 1/(60+3) + 1/(60+1) = 0.0159 + 0.0164 = 0.0323
另一个文档 E 在向量检索排第 1、但 BM25 没有命中:
RRF_score(E) = 1/(60+1) = 0.0164
D 的分数 > E 的分数,因为 D 在两路检索中都排名靠前。
为什么 k=60? 这是 RRF 论文推荐的默认值。k 的作用是"平滑"排名——k 越大,高排名和低排名的分数差距越小。60 是经验值,对大多数场景效果好。
面试话术:“RRF 的核心思想是——在多个排序列表中都排名靠前的文档,融合后排名更高。它不依赖具体的分数值(向量距离和 BM25 分数的量纲不同,不能直接比较),只看排名位置,所以天然适合融合不同类型的检索结果。”
查询改写(Query Rewriting)
# 第 109-147 行
REWRITE_PROMPT = """你是一个搜索查询优化器。将用户的口语化描述改写为适合技术文档检索的结构化查询。
...
用户描述:{question}
改写查询:"""
def rewrite_query(question: str, llm=None) -> str:
if llm is None:
return question # 无 LLM 则跳过改写
response = llm.invoke([HumanMessage(content=REWRITE_PROMPT.format(question=question))])
return extract_text(response.content).strip()
用途:用户说"冰机不出冰了",改写为"制冰机 出冰异常 故障"——更适合关键词检索。
默认关闭(use_rewrite=False),因为每次改写要调一次 LLM,增加延迟。只在需要时手动开启。
3.5 语义重排序(reranker.py)
文件路径:
src/rag/reranker.py
CrossEncoder vs Bi-Encoder
Bi-Encoder(向量检索用的):
查询 → [编码器] → 向量Q
文档 → [编码器] → 向量D
相似度 = cosine(Q, D) ← 查询和文档分别编码,可预计算文档向量
CrossEncoder(重排序用的):
[查询 + 文档] → [编码器] → 相关性分数
← 查询和文档一起编码,交互更深入,但不能预计算
为什么需要 CrossEncoder 重排序?
- Bi-Encoder 把查询和文档分别编码,速度快(可以预先算好所有文档的向量),但精度有限
- CrossEncoder 把查询和文档拼在一起编码,能捕捉更细粒度的语义关系,精度更高
- 先用 Bi-Encoder 快速筛选出 top 20,再用 CrossEncoder 精排出 top 5——速度和精度的平衡
# 第 31-57 行
def rerank(query, documents, top_k=5):
model = _get_reranker() # 懒加载 CrossEncoder
pairs = [[query, doc.page_content] for doc in documents] # 构造 (query, doc) 对
scores = model.predict(pairs) # 批量打分
scored_docs = list(zip(documents, scores))
scored_docs.sort(key=lambda x: x[1], reverse=True)
return [(doc, float(score)) for doc, score in scored_docs[:top_k]]
模型:ms-marco-MiniLM-L-6-v2,微软在 MS MARCO 数据集上训练的轻量级 CrossEncoder,约 80MB,CPU 推理 20 个文档 < 1 秒。
Java 类比:想象数据库查询——先用索引快速定位到候选行(Bi-Encoder),再用复杂条件精确过滤(CrossEncoder)。
3.6 RAG Chain(chain.py)
文件路径:
src/rag/chain.py
Prompt 工程:上下文注入
# 第 16-30 行:RAG 系统提示词模板
RAG_SYSTEM_PROMPT = """你是 * 饮吧设备的智能故障诊断助手。
请根据以下检索到的设备文档来回答用户的问题。
## 参考文档
{context} ← 检索到的文档片段会被填充到这里
## 回答要求
1. 优先基于上面的参考文档来回答,如果文档中有相关内容,请引用具体来源
2. 如果文档中没有直接答案,可以结合文档信息和你的知识进行推理
3. 回答要具体可执行,像一个经验丰富的运维工程师在和同事交流
4. 在回答末尾列出引用的文档来源
"""
Prompt 设计要点:
- 角色定义:明确告诉 LLM 它是"设备诊断助手"
- 上下文注入:
{context}占位符,运行时替换为检索到的文档片段 - 行为约束:要求"基于文档回答"、“引用来源”——减少幻觉
- 风格要求:像运维工程师交流——专业但易懂
# 第 33-56 行:格式化上下文
def format_context(results):
for i, (doc, score) in enumerate(results, 1):
source = doc.metadata.get("filename", "未知来源")
headers = []
for key in ["h1", "h2", "h3"]: # 提取标题层级
if key in doc.metadata:
headers.append(doc.metadata[key])
# 输出格式:
# [文档 1] 来源: ice-machine-manual.md | 章节: 故障码 > E03 | 相关度: 0.876
# 文档内容...
为什么显示 metadata 信息? 让 LLM 知道每个文档片段的来源和章节,这样 LLM 的回答可以引用"根据 ice-machine-manual.md 的 E03 故障码章节…",增加可信度。
rag_query / rag_stream 双模式
# 第 59-107 行:非流式
def rag_query(question, provider=None, k=5):
results = enhanced_search(question, k=k) # ① 混合检索
context = format_context(results) # ② 格式化上下文
llm = create_llm(provider)
messages = [
SystemMessage(content=RAG_SYSTEM_PROMPT.format(context=context)),
HumanMessage(content=question),
]
response = llm.invoke(messages) # ③ LLM 生成回答
return {"answer": ..., "sources": ..., "context_count": ...}
# 第 110-148 行:流式
def rag_stream(question, provider=None, k=5):
# 同样的检索逻辑,但用 llm.stream() + yield 逐 token 输出
for chunk in llm.stream(messages):
yield chunk.content # 流式返回每个 token
3.7 知识库构建脚本(build_knowledge_base.py)
文件路径:
scripts/build_knowledge_base.py
# 第 23-67 行
def build(use_local=False):
# Step 1: 加载文档(device-manuals/ + fault-cases/)
documents = load_directory(manuals_dir) + load_directory(cases_dir)
# Step 2: 切分(chunk_size=800, overlap=150)
chunks = split_documents(documents, chunk_size=800, chunk_overlap=150)
# Step 3: Embedding + 存入向量数据库
vectorstore = create_vector_store(chunks, use_local_embedding=use_local)
运行方式:
# 使用本地 Embedding(离线,推荐开发用)
python scripts/build_knowledge_base.py --local
# 使用 OpenAI Embedding(效果更好)
python scripts/build_knowledge_base.py
这是一个一次性脚本:文档变化时重新运行即可。不需要在线实时构建。
面试高频问题 & 参考回答
Q1:你的 RAG 系统的检索流程是怎样的?
"我们使用混合检索 + 重排序的四步流程:
- 向量语义检索:用
all-MiniLM-L6-v2模型把查询编码为向量,在 ChromaDB 中做相似度搜索,取 top 20- BM25 关键词检索:用字符级 + 2-gram 分词,基于 BM25 算法做关键词匹配,取 top 20
- RRF 融合:用 Reciprocal Rank Fusion 算法融合两路排序结果,不依赖具体分数值,只看排名位置
- CrossEncoder 重排序:用
ms-marco-MiniLM-L-6-v2对融合后的 top 20 精排,返回 top 5这样兼顾了语义理解(向量检索擅长同义词)和精确匹配(BM25 擅长错误码、型号),再通过 CrossEncoder 进一步提升精度。"
Q2:为什么用混合检索而不是只用向量检索?
“纯向量检索对精确关键词不敏感。比如用户查’E03 错误码’,语义上和’设备故障’接近,但向量检索可能返回各种故障相关的文档,而不是精确包含 E03 的文档。BM25 通过关键词匹配可以精确找到包含 E03 的文档片段。两者结合覆盖面更广,实测 Recall@5 提升了约 15%。”
Q3:RRF 的 k 参数为什么是 60?
“k=60 是 RRF 论文(Cormack et al., 2009)推荐的默认值。k 的作用是控制排名的’平滑度’——k 越大,高排名和低排名的分数差距越小。60 在大多数信息检索场景下表现稳定。我们也做过实验,k 在 40-80 之间对结果影响很小。”
Q4:chunk_size 怎么选的?为什么是 800?
"chunk_size 的选择需要平衡两个矛盾:
- 太小(如 200):上下文信息不完整,一个故障描述可能被切断
- 太大(如 2000):检索不精确,返回一大段中只有一小部分相关
我们的知识库是中文技术文档,一个完整的故障描述通常在 400-800 字。chunk_size=800 能确保大多数故障案例在一个 chunk 内完整呈现。chunk_overlap=150(约 19%)避免切分边界丢失信息。"
Q5:Embedding 模型为什么用本地的 all-MiniLM-L6-v2?
“主要考虑部署环境。项目部署在 IoT 设备上,网络环境不稳定。本地模型保证离线可用,不依赖外部 API。all-MiniLM-L6-v2 只有 80MB,CPU 推理速度快。虽然中文效果不如 OpenAI 的 text-embedding-3-small,但配合 BM25 混合检索和 CrossEncoder 重排序,整体效果是可以接受的。如果需要更好的中文效果,可以换成
bge-small-zh-v1.5等中文专用模型。”
第四章:Agent 与工具系统
什么是 Agent?
Agent = LLM + 工具 + 自主决策
传统的 LLM 应用是固定流程(Chain):输入 → 处理 → 输出。而 Agent 让 LLM 自己决定下一步做什么:
Chain(固定流程):用户问题 → 检索文档 → 生成回答
Agent(动态决策):用户问题 → LLM 思考 → 决定调用什么工具 → 观察结果 → 继续思考 → ...
Java 类比:
Chain ≈ 面向过程编程,代码写死了每一步
Agent ≈ 面向对象 + 策略模式,运行时动态选择执行策略
4.1 ReAct 模式详解
ReAct(Reasoning + Acting)是当前最主流的 Agent 模式。
循环:
┌───────────────────────────────────────────┐
│ │
▼ │
Thought(思考) │
"用户说制冰机报错,我需要先查设备状态" │
│ │
▼ │
Action(行动) │
调用 query_device_status(module="制冰机") │
│ │
▼ │
Observation(观察) │
"制冰机: 状态=error, 错误码=E03" │
│ │
▼ │
Thought(继续思考) │
"得到了错误码 E03,去知识库查一下" ─────────┘
│
▼
... 重复直到信息够了 ...
│
▼
Final Answer(最终回答)
"根据分析,制冰机 E03 错误是..."
核心代码(src/agent/diagnostic_agent.py 第 79-83 行):
agent = create_react_agent(
model=llm, # LLM 负责"思考"和"决策"
tools=TOOLS, # 6 个工具供 Agent 选择
prompt=AGENT_SYSTEM_PROMPT, # 系统提示词定义 Agent 角色
)
create_react_agent 是 LangGraph 提供的函数,它内部做了:
- 把工具列表注册到 LLM(LLM 知道有哪些工具可用)
- 构建 Thought → Action → Observation 的循环图
- 当 LLM 不再调用工具时,自动退出循环,返回最终回答
4.2 诊断 Agent(diagnostic_agent.py)
文件路径:
src/agent/diagnostic_agent.py
System Prompt 设计
# 第 32-54 行
AGENT_SYSTEM_PROMPT = """你是 * 饮吧设备的智能故障诊断 Agent。
你的能力:
1. 查询设备实时状态(query_device_status)
2. 读取设备运行日志(fetch_device_logs)
3. 检索设备知识库(search_knowledge_base)
4. 生成故障诊断报告(generate_diagnosis_report)
5. 重启设备服务(restart_device_service)
6. 检查固件/软件版本(check_firmware_version)
你的工作流程:
1. 接收用户描述的故障现象
2. 先查询设备状态,了解当前各模块的运行情况
3. 根据故障现象读取相关日志,寻找错误信息
4. 搜索知识库,查找类似故障案例和技术文档
5. 综合分析后,生成结构化的诊断报告
注意事项:
- 每次诊断至少使用 2-3 个工具来收集信息
- 不要猜测,要基于实际数据进行分析
- 如果信息不足,主动使用工具获取更多信息
"""
Prompt 设计原则:
- 列出所有工具名称和用途:让 LLM 知道有什么"武器"
- 定义工作流程:引导 LLM 按"状态→日志→知识库→报告"的顺序工作
- 设置约束:至少用 2-3 个工具(防止 LLM 偷懒直接回答)
6 个工具注册
# 第 57-64 行
TOOLS = [
search_knowledge_base, # RAG 检索
fetch_device_logs, # 日志读取
query_device_status, # 设备状态
generate_diagnosis_report, # 诊断报告
restart_device_service, # 服务重启
check_firmware_version, # 版本检查
]
Agent 通过工具的 name 和 description(docstring)来决定何时调用哪个工具。工具的 description 写得好不好,直接影响 Agent 的智能程度。
执行入口
# 第 88-133 行
def run_diagnosis(question, provider=None, verbose=True):
agent = create_diagnostic_agent(provider)
result = agent.invoke({
"messages": [HumanMessage(content=question)],
})
messages = result["messages"]
# messages 包含完整的推理链:
# HumanMessage → AIMessage(tool_calls) → ToolMessage → AIMessage(tool_calls) → ToolMessage → AIMessage(回答)
# 提取最后一条 AI 消息作为最终回答
ai_messages = [m for m in messages if isinstance(m, AIMessage) and m.content]
return extract_text(ai_messages[-1].content)
4.3 工具实现逐个讲解
4.3.1 knowledge_search — RAG 检索封装
文件路径:
src/tools/knowledge_search.py
# 第 13-42 行
@tool
def search_knowledge_base(query: str) -> str:
"""搜索设备知识库,查找设备手册、操作规范、故障案例等技术文档。
当需要查询设备的技术参数、操作流程、日志规范、错误码含义、
硬件接口说明等信息时使用此工具。"""
results = enhanced_search(query, k=5)
# 格式化为 Agent 可读的文本
for i, (doc, score) in enumerate(results, 1):
output_parts.append(f"[{i}] 来源: {source}\n{doc.page_content[:500]}")
return "\n\n---\n\n".join(output_parts)
@tool 装饰器是 LangChain 定义工具的标准方式。它会:
- 从函数名生成工具名(
search_knowledge_base) - 从 docstring 生成工具描述(Agent 据此决定何时调用)
- 从参数类型注解生成参数 schema(Agent 据此填充参数)
Java 类比:类似 Spring 的 @Bean + 自定义注解,框架自动识别和注册。
4.3.2 log_reader — SSH/本地双模式 + 安全校验
文件路径:
src/tools/log_reader.py
# 第 17-22 行:日志路径映射
LOG_PATHS = {
"syslog": "/var/log/syslog",
"middleware": "/home/*/bar_middleware/logs/app.log",
"deploy": "/home/*/bar-deploy-client/logs/app.log",
"docker": None, # 特殊处理
}
# 第 25 行:关键词安全校验(防止命令注入)
_KEYWORD_RE = re.compile(r'^[\w\s\-\.]*$')
def _validate_keyword(keyword: str) -> str | None:
if keyword and not _KEYWORD_RE.match(keyword):
return "keyword 包含非法字符"
return None
安全设计:keyword 参数会被拼入 SSH 命令。如果用户传入 "; rm -rf /" 就会造成命令注入。通过白名单正则校验,只允许字母、数字、空格、连字符、点、下划线。
# 第 83-128 行:SSH 模式
def _read_ssh(log_type, lines, keyword):
if keyword:
remote_cmd = f"grep -i -e {shlex.quote(keyword)} {log_path} | tail -n {lines}"
# ^^^^^^^^^^^^^^^^^ shlex.quote 做 shell 转义
ssh_cmd = [
"ssh",
"-o", "ConnectTimeout=5",
"-o", "StrictHostKeyChecking=no",
f"{settings.DEVICE_SSH_USER}@{settings.DEVICE_SSH_HOST}",
remote_cmd,
]
result = subprocess.run(ssh_cmd, capture_output=True, text=True, timeout=15)
双重安全:
- 白名单正则过滤非法字符
shlex.quote()做 shell 参数转义
Java 类比:类似 JDBC 的 PreparedStatement 防 SQL 注入 + 输入校验。
4.3.3 device_status — 真实设备 API 对接
文件路径:
src/tools/device_status.py
# 第 13-34 行:硬件模块 → API 路径映射
MODULE_STATUS_ENDPOINTS = {
"制冰机": "/ice/machine/status",
"咖啡机": "/coffee/machine/status",
"杯子机": "/cup/machine/status",
# ... 11 个标准模块
}
SPECIAL_ENDPOINTS = {
"机械臂": "/RobotArm/Dev/status",
"物料": "/agent/material/all",
# ... 5 个特殊模块
}
# 第 39-47 行:HTTP 请求
def _fetch(path: str) -> dict | None:
url = f"{settings.MIDDLEWARE_BASE_URL}{path}"
try:
resp = httpx.get(url, timeout=HTTP_TIMEOUT) # httpx ≈ Java 的 RestTemplate
resp.raise_for_status()
return resp.json()
except Exception:
return None # 接口异常返回 None,不抛异常
设计亮点:工具调用失败时不抛异常,而是返回可读的错误信息(如"接口无响应(模块可能未启用或离线)")。Agent 拿到错误信息后可以自行判断下一步动作。
4.3.4 diagnosis_report — 结构化报告模板
文件路径:
src/tools/diagnosis_report.py
# 第 12-56 行
@tool
def generate_diagnosis_report(
device_id: str,
fault_description: str,
root_cause: str,
severity: str = "中",
repair_steps: str = "",
references: str = "",
) -> str:
"""生成结构化的故障诊断报告。在完成故障分析后,使用此工具输出规范化的诊断报告。"""
report = f"""
📋 设备故障诊断报告
⏰ 时间: {timestamp}
🏷️ 设备: {device_id}
🔴 严重程度: {severity}
📝 故障现象: {fault_description}
🔍 根因分析: {root_cause}
🔧 修复建议: {repair_steps}
📄 参考文档: {references}
"""
这个工具不获取信息,而是输出格式化结果。Agent 收集够信息后,主动调用此工具来生成规范的报告。
4.3.5 device_restart — 白名单安全设计
文件路径:
src/tools/device_restart.py
# 第 15-20 行:白名单
SERVICE_COMMANDS = {
"middleware": "systemctl restart bar_middleware",
"deploy": "systemctl restart bar-deploy-client",
"docker": "docker restart $(docker ps -q | head -1)",
"system": "sudo reboot",
}
安全设计:只允许重启预定义的服务,不接受任意命令。即使 Agent "被忽悠"传入恶意参数,白名单机制会拦截。
4.3.6 firmware_check — 版本查询
文件路径:
src/tools/firmware_check.py
# 第 12-15 行
VERSION_ENDPOINTS = {
"middleware": "/mid/version",
"deploy": "/deploy/version",
}
简单的 HTTP 版本查询接口。帮助 Agent 判断是否因为软件版本过旧导致故障。
4.4 多 Agent 协作(supervisor.py)
文件路径:
src/agent/supervisor.py
Supervisor 路由模式
# 第 25-27 行:路由意图关键词定义
_REPAIR_KEYWORDS = {"维修", "怎么修", "修复", "备件", "更换", "拆卸", "安装", "维护"}
_MONITOR_KEYWORDS = {"状态", "健康", "检查", "监控", "巡检", "所有模块", "整体", "概况"}
_DIAGNOSE_KEYWORDS = {"故障", "报错", "不工作", "异常", "失败", "错误", "问题", "坏了"}
# 第 43-58 行:关键词路由函数
def _route_by_keywords(question):
has_repair = any(kw in question for kw in _REPAIR_KEYWORDS)
has_diagnose = any(kw in question for kw in _DIAGNOSE_KEYWORDS)
if has_repair and has_diagnose:
return "diagnose+repair" # 先诊断再给维修方案
if has_repair:
return "repair"
# ...
# 第 61-72 行:LLM 路由(精确、有延迟)
def _route_by_llm(question, llm):
response = llm.invoke([HumanMessage(content=ROUTE_PROMPT.format(question=question))])
route = extract_text(response.content).strip().lower()
valid = {"diagnose", "repair", "monitor", "diagnose+repair"}
return route if route in valid else "diagnose"
两种路由策略:
| 策略 | 实现 | 优点 | 缺点 |
|---|---|---|---|
| 关键词路由(默认) | 正则匹配关键词 | 零延迟、零成本 | 无法理解复杂意图 |
| LLM 路由 | 让 LLM 判断意图 | 理解自然语言 | 多一次 LLM 调用(0.5-2秒) |
Agent 链式调用
# 第 135-148 行:diagnose+repair 链式模式
elif route == "diagnose+repair":
# Step 1: 诊断 Agent 分析故障
diag_agent = create_diagnostic_agent(provider)
diag_result = diag_agent.invoke({"messages": [HumanMessage(content=question)]})
diagnosis = _extract_final_answer(diag_result)
# Step 2: 把诊断结论传给维修 Agent
repair_question = f"设备故障描述:{question}\n\n诊断结论:\n{diagnosis}\n\n请给出具体维修方案。"
repair_agent = create_repair_agent(provider)
repair_result = repair_agent.invoke({"messages": [HumanMessage(content=repair_question)]})
final_result = f"## 故障诊断\n\n{diagnosis}\n\n## 维修方案\n\n{repair_plan}"
设计亮点:诊断 Agent 的输出(诊断结论)作为维修 Agent 的输入。两个专业 Agent 串联工作,各司其职。
Java 类比:
// 类似责任链模式
DiagnoseResult diagnosis = diagnoseAgent.handle(question);
RepairPlan repairPlan = repairAgent.handle(diagnosis);
return new CompositeResult(diagnosis, repairPlan);
4.5 维修 Agent & 监控 Agent
文件路径:
src/agent/repair_agent.py、src/agent/monitor_agent.py
结构与诊断 Agent 完全一致,只是:
- System Prompt 不同:定义不同的角色和工作流程
- 工具集不同:维修 Agent 只需要知识库搜索 + 报告生成;监控 Agent 需要状态查询 + 日志读取 + 版本检查
# repair_agent.py:2 个工具
REPAIR_TOOLS = [search_knowledge_base, generate_diagnosis_report]
# monitor_agent.py:3 个工具
MONITOR_TOOLS = [query_device_status, fetch_device_logs, check_firmware_version]
设计原则:每个 Agent 只拥有它需要的工具,避免工具过多导致 LLM 选择困难。
面试高频问题 & 参考回答
Q1:你的 Agent 是怎么实现的?用的什么模式?
“使用 LangGraph 的
create_react_agent实现 ReAct 模式。Agent 在一个 Thought → Action → Observation 循环中工作:LLM 先思考应该做什么,然后选择调用某个工具,观察工具返回结果,再继续思考。直到收集够信息,给出最终回答。”“我们项目有 3 个专业 Agent:诊断 Agent(6个工具)、维修 Agent(2个工具)、监控 Agent(3个工具),通过 Supervisor 模式路由和编排。”
Q2:Agent 和 Chain 有什么区别?为什么选 Agent?
“Chain 是固定流程,代码写死了’先做A再做B’。Agent 是动态决策,LLM 自己决定下一步做什么。”
“我们选 Agent 是因为故障诊断场景天然需要动态决策——不同故障需要查不同的模块状态、读不同的日志、检索不同的知识。如果用 Chain 写死流程,每次都查所有模块的状态,既浪费时间又增加 token 消耗。Agent 可以根据故障描述有针对性地收集信息。”
Q3:多 Agent 协作是怎么实现的?
“使用 Supervisor 模式。Supervisor 接收用户请求后,先做意图识别(关键词匹配或 LLM 路由),然后把请求分发给对应的专业 Agent。”
“对于复杂故障,支持链式模式:先由诊断 Agent 分析根因,再把诊断结论传给维修 Agent 生成维修方案。两个 Agent 各司其职,输出拼接成最终结果。”
Q4:工具的 description 对 Agent 有多重要?
“非常重要,是 Agent ‘Function Calling’ 能力的基础。LLM 通过工具的 name 和 description 来决定什么时候该调用什么工具。如果 description 写得模糊,Agent 可能选错工具或者不知道该用工具。”
“我们的实践:description 要明确说明’什么场景下用这个工具’和’输入输出是什么’。比如
search_knowledge_base的 description 列出了具体的使用场景:'设备技术参数、操作流程、日志规范、错误码含义’等。”
Q5:Agent 的安全性怎么保证?
"几个层面:
- 工具白名单:Agent 只能调用预注册的工具,不能执行任意代码
- 输入校验:比如
log_reader的 keyword 参数有正则白名单校验,防止命令注入- 命令白名单:
device_restart只能执行预定义的重启命令,不接受任意系统命令- 参数转义:SSH 命令拼接时用
shlex.quote()做 shell 转义- 超时控制:所有工具调用都有超时(HTTP 5秒、SSH 15秒),防止 hang 住"
第五章:生产级 API 服务
5.1 FastAPI 应用(app.py)
文件路径:
src/api/app.py
# 第 18-22 行
app = FastAPI(
title="* 设备故障诊断 Agent API",
description="基于 LangChain + RAG 的智能设备故障诊断服务",
version="1.0.0",
)
FastAPI vs Spring Boot 对照:
| 功能 | FastAPI(本项目) | Spring Boot(你熟悉的) |
|---|---|---|
| 应用创建 | app = FastAPI() | @SpringBootApplication |
| 路由定义 | @router.post("/chat") | @PostMapping("/chat") |
| 请求体 | Pydantic BaseModel | @RequestBody DTO |
| 依赖注入 | Depends() | @Autowired |
| 中间件 | app.add_middleware() | @Component Filter |
| API 文档 | 自动生成 Swagger(/docs) | Springdoc-openapi |
| 启动命令 | uvicorn src.api.app:app | java -jar app.jar |
# 第 27-37 行:CORS 配置
_raw_origins = os.getenv("CORS_ORIGINS", "http://localhost:3000,http://localhost:5173")
_allow_origins = [o.strip() for o in _raw_origins.split(",") if o.strip()]
app.add_middleware(
CORSMiddleware,
allow_origins=_allow_origins, # 通过环境变量配置允许的前端域名
allow_methods=["*"],
allow_headers=["*"],
)
Java 类比:
@Configuration
public class CorsConfig implements WebMvcConfigurer {
@Override
public void addCorsMappings(CorsRegistry registry) {
registry.addMapping("/**")
.allowedOrigins("http://localhost:3000")
.allowedMethods("*");
}
}
# 第 24-25 行:限流
app.state.limiter = limiter
app.add_exception_handler(RateLimitExceeded, _rate_limit_exceeded_handler)
slowapi 是 Python 版的限流中间件,基于 IP 地址限流。
Java 类比:类似 Spring Cloud Gateway 的 RequestRateLimiter 或 Sentinel 的 @SentinelResource。
5.2 API 路由(routes.py)
文件路径:
src/api/routes.py
完整 API 列表
| 端点 | 方法 | 限流 | 功能 |
|---|---|---|---|
/api/chat | POST | 30/min | 普通对话(非流式) |
/api/chat/stream | POST | 30/min | 普通对话(SSE 流式) |
/api/diagnose | POST | 10/min | Agent 诊断(单/多Agent) |
/api/diagnose/stream | POST | 10/min | Agent 诊断(SSE 流式) |
/api/rag/query | POST | 20/min | RAG 知识库检索问答 |
/api/session/{id} | DELETE | - | 删除指定会话 |
/api/sessions | GET | - | 列出所有活跃会话 |
/health | GET | - | 健康检查 |
聊天接口
# 第 37-53 行:同步聊天
@router.post("/chat", response_model=ChatResponse)
@limiter.limit("30/minute")
async def chat(request: Request, req: ChatRequest):
session = _store.get(req.session_id, req.provider) # 获取或创建会话
answer = session.chat(req.message) # 调用 DiagnosticChat
if hasattr(_store, 'save'):
_store.save(req.session_id) # 保存到 Redis
return ChatResponse(answer=answer, provider=session.provider_name, session_id=req.session_id)
诊断接口(核心)
# 第 81-136 行
@router.post("/diagnose", response_model=DiagnoseResponse)
async def diagnose(request: Request, req: DiagnoseRequest):
if req.use_multi_agent:
# 多 Agent 模式:Supervisor 路由
result = await loop.run_in_executor(None, partial(run_multi_agent, ...))
return DiagnoseResponse(result=result["result"], route=result["route"], agents_used=result["agents_used"])
# 单 Agent 模式
agent = create_diagnostic_agent(req.provider)
result = await loop.run_in_executor(None, lambda: agent.invoke(...))
# 收集工具调用步骤
for msg in messages:
if msg_type == "AIMessage" and msg.tool_calls:
tools_used.append(tc["name"])
steps.append({"type": "tool_call", "tool": tc["name"], "args": tc["args"]})
elif msg_type == "ToolMessage":
steps.append({"type": "tool_result", "tool": msg.name, "result": msg.content[:300]})
run_in_executor 为什么? LangChain 的 agent.invoke() 是同步阻塞调用,但 FastAPI 是异步框架。用 run_in_executor 把同步调用放到线程池中执行,不阻塞事件循环。
Java 类比:
// 类似 Spring WebFlux 中用 Mono.fromCallable 包装阻塞调用
Mono.fromCallable(() -> agent.invoke(question))
.subscribeOn(Schedulers.boundedElastic())
5.3 SSE 流式输出机制
服务端
# 第 56-76 行:SSE 事件生成器
@router.post("/chat/stream")
async def chat_stream(request: Request, req: ChatRequest):
session = _store.get(req.session_id, req.provider)
async def event_generator():
for token in session.stream_chat(req.message): # generator 逐 token 产出
yield f"data: {json.dumps({'token': token})}\n\n" # SSE 格式
yield f"data: {json.dumps({'done': True})}\n\n" # 结束信号
return StreamingResponse(
event_generator(),
media_type="text/event-stream", # SSE 的 Content-Type
)
SSE 协议格式:
data: {"token": "根据"}
data: {"token": "分析"}
data: {"token": ",制冰机"}
data: {"done": true}
每条消息以 data: 开头,以两个换行 \n\n 结尾。
Java 类比:
// Spring 的 SseEmitter
@PostMapping("/chat/stream")
public SseEmitter chatStream(@RequestBody ChatRequest req) {
SseEmitter emitter = new SseEmitter();
executor.execute(() -> {
for (String token : chat.streamTokens(req.getMessage())) {
emitter.send(SseEmitter.event().data(Map.of("token", token)));
}
emitter.send(SseEmitter.event().data(Map.of("done", true)));
emitter.complete();
});
return emitter;
}
客户端(React 前端)
// frontend/src/App.tsx 第 77-123 行
const sendChatStream = async (text: string) => {
const response = await fetch('/api/chat/stream', { method: 'POST', ... })
const reader = response.body?.getReader() // 获取可读流
const decoder = new TextDecoder()
let buffer = ''
while (reader) {
const { done, value } = await reader.read() // 逐块读取
if (done) break
buffer += decoder.decode(value, { stream: true })
const lines = buffer.split('\n\n') // 按 SSE 分隔符切分
buffer = lines.pop() || '' // 最后一个可能不完整,留到下次
for (const line of lines) {
if (!line.startsWith('data: ')) continue
const data = JSON.parse(line.slice(6)) // 去掉 "data: " 前缀
if (data.token) {
// 追加 token 到最后一条消息
setMessages(prev => {
const updated = [...prev]
updated[updated.length - 1].content += data.token
return updated
})
}
}
}
}
关键技术点:
response.body.getReader()— 获取ReadableStream读取器buffer处理 — SSE 消息可能跨网络包,需要缓冲区处理不完整的消息- 逐 token 追加 — 每收到一个 token 立即更新 UI
5.4 会话管理(session_store.py)
文件路径:
src/utils/session_store.py
Redis + 内存双后端降级
# 第 151-167 行
def create_session_store():
if settings.REDIS_URL:
try:
store = RedisSessionStore()
store._redis.ping() # 测试连接
return store
except Exception:
logger.warning("Redis 连接失败,降级到内存存储")
return MemorySessionStore()
降级逻辑:
- 如果配置了
REDIS_URL,尝试连接 Redis - 连接成功 → 使用
RedisSessionStore - 连接失败 → 降级到
MemorySessionStore
Java 类比:
// 类似多级缓存策略
@Bean
public SessionStore sessionStore() {
try {
return new RedisSessionStore(redisTemplate);
} catch (Exception e) {
log.warn("Redis unavailable, falling back to in-memory");
return new InMemorySessionStore();
}
}
序列化/反序列化
# 第 26-59 行
def _serialize_session(chat: DiagnosticChat) -> str:
data = {"provider": chat.provider_name, "messages": []}
for msg in chat.history:
msg_type = type(msg).__name__.lower().replace("message", "") # "HumanMessage" → "human"
data["messages"].append({"type": msg_type, "content": msg.content})
return json.dumps(data, ensure_ascii=False)
def _deserialize_session(raw: str, fallback_provider: str) -> DiagnosticChat:
data = json.loads(raw)
# 兼容旧格式(纯列表)和新格式(含 provider 的 dict)
if isinstance(data, list):
provider = fallback_provider
messages_data = data
else:
provider = data.get("provider", fallback_provider)
messages_data = data.get("messages", [])
# ... 恢复 DiagnosticChat 对象
为什么需要自定义序列化? LangChain 的 HumanMessage、AIMessage 等对象不是简单的 POJO,不能直接 json.dumps()。需要手动把消息类型和内容提取出来。
RedisSessionStore 的缓存设计
# 第 86-148 行
class RedisSessionStore:
def __init__(self):
self._redis = redis.from_url(settings.REDIS_URL)
self._ttl = settings.SESSION_TTL
self._cache: dict[str, DiagnosticChat] = {} # 内存缓存
def get(self, session_id, provider):
if session_id in self._cache: # ① 先查内存缓存
self._redis.expire(key, self._ttl) # 刷新 TTL
return self._cache[session_id]
raw = self._redis.get(key) # ② 再查 Redis
if raw:
chat = _deserialize_session(raw)
self._cache[session_id] = chat # 写入内存缓存
return chat
chat = DiagnosticChat(provider) # ③ 都没有,新建
self._save(session_id, chat)
self._cache[session_id] = chat
return chat
三级查找:内存缓存 → Redis → 新建。避免每次请求都做 Redis 网络 I/O + JSON 反序列化。
TTL 策略:每次访问会话都刷新 Redis TTL(默认 3600 秒)。长期不活跃的会话自动过期清理。
5.5 React 前端
文件路径:
frontend/src/App.tsx
前端是一个单文件 React 应用,核心功能:
- 模式切换:Agent 诊断模式 / 普通对话模式
- 模型选择:Claude / OpenAI
- 多 Agent 开关:启用/禁用多 Agent 协作
- SSE 流式渲染:实时显示 AI 回答
- 工具调用可视化:Agent 模式下显示工具调用过程
// 第 4-8 行:消息类型定义
interface Message {
role: 'user' | 'assistant' | 'tool-call' | 'tool-result' | 'error'
content: string
tool?: string // 工具名(tool-call/tool-result 时有值)
}
工具调用可视化:
// 第 149-154 行
if (data.type === 'tool_call') {
setMessages(prev => [...prev, {
role: 'tool-call',
content: `调用工具: ${data.tool}(${JSON.stringify(data.args)})`,
tool: data.tool,
}])
}
前端把 Agent 的每一步工具调用都显示出来,让用户看到"AI 在做什么"——这比黑盒式的"等待中…"体验好很多。
5.6 数据模型(models.py)
文件路径:
src/api/models.py
# Pydantic v2 模型定义
class ChatRequest(BaseModel):
message: str = Field(..., description="用户消息") # ... 表示必填
provider: str = Field(default="claude") # 有默认值 = 选填
session_id: str = Field(default="default")
class DiagnoseRequest(BaseModel):
question: str = Field(..., description="故障描述")
use_multi_agent: bool = Field(default=False)
use_llm_routing: bool = Field(default=False)
Java 类比:
// 完全对应 Spring 的 DTO
public class ChatRequest {
@NotBlank private String message;
private String provider = "claude";
private String sessionId = "default";
}
面试高频问题 & 参考回答
Q1:你的 API 是怎么做流式输出的?
“使用 SSE(Server-Sent Events)协议。后端 FastAPI 用
StreamingResponse+ async generator,每生成一个 token 就通过 SSE 事件发送到前端。前端用ReadableStreamAPI 实时读取事件流。”“Agent 模式下的 SSE 更复杂,会推送三种事件:
tool_call(Agent 调用工具)、tool_result(工具返回结果)、answer(最终回答)。前端根据事件类型渲染不同的 UI 组件,让用户能看到 Agent 的推理过程。”
Q2:会话管理怎么做的?
“双后端降级策略。优先使用 Redis 存储会话历史,Redis 不可用时自动降级到内存存储。Redis 存储使用 JSON 序列化,支持 TTL 自动过期(默认 1 小时不活跃就清理)。”
“RedisSessionStore 内部还有一层内存缓存,避免每次请求都做 Redis 网络 I/O 和 JSON 反序列化。查找顺序是:内存缓存 → Redis → 新建。”
Q3:FastAPI 和 Spring Boot 有什么异同?
"核心概念高度相似:都有路由、中间件、依赖注入、请求/响应模型。主要区别在于:
- FastAPI 原生支持 async/await,Spring Boot 需要 WebFlux
- FastAPI 自动从类型注解生成 Swagger 文档,Spring 需要 Springdoc
- FastAPI 的 Pydantic 模型同时做序列化和校验,Spring 需要 Jackson + Bean Validation
- 部署方式不同:FastAPI 用 uvicorn(ASGI 服务器),Spring Boot 内嵌 Tomcat"
Q4:为什么诊断接口要用 run_in_executor?
“因为 LangChain 的 Agent 调用是同步阻塞的(
agent.invoke()可能运行几十秒),但 FastAPI 是异步框架。如果在 async 路由函数里直接调用同步代码,会阻塞事件循环,导致其他请求无法处理。run_in_executor把同步调用放到线程池执行,不阻塞主事件循环。类似 Spring WebFlux 中用Mono.fromCallable().subscribeOn(Schedulers.boundedElastic())。”
第六章:设计模式总结
项目中使用的设计模式速查表
| 设计模式 | 代码位置 | 作用 |
|---|---|---|
| 工厂模式 | src/llm/client.py:24-46 _create_single_llm() | 根据 provider 参数创建不同的 LLM 实例 |
| 单例模式(DCL) | src/rag/vectorstore.py:125-132 _get_vectorstore() | 全局唯一的向量存储实例,线程安全懒加载 |
| 单例模式(DCL) | src/rag/reranker.py:19-28 _get_reranker() | 全局唯一的 CrossEncoder 实例 |
| 单例模式(DCL) | src/rag/hybrid_search.py:66-80 _get_bm25_index() | 全局唯一的 BM25 索引 |
| 单例模式(模块级) | src/config/settings.py:56 settings = Settings() | 全局唯一的配置对象 |
| 策略模式 | src/rag/loader.py:69 supported = {".md": load_markdown, ...} | 根据文件扩展名选择不同的加载策略 |
| 策略模式 | src/rag/vectorstore.py:119-122 | 根据配置选择 ChromaDB 或 Milvus |
| 策略模式 | src/agent/supervisor.py:43-72 | 关键词路由 vs LLM 路由两种策略 |
| 模板方法 | src/tools/diagnosis_report.py:33-55 | 固定报告格式,参数填充 |
| 责任链/管道 | src/rag/hybrid_search.py:150-194 enhanced_search() | 向量检索 → BM25 → RRF → Rerank 的管道 |
| 装饰器模式 | 所有 @tool 装饰的函数 | LangChain 注册工具的标准方式 |
| 降级/熔断 | src/llm/client.py:86 with_fallbacks() | 主模型失败自动切换备用模型 |
| 降级/熔断 | src/utils/session_store.py:151-167 | Redis 失败降级到内存存储 |
| 降级/熔断 | src/llm/client.py:167-172 stream_chat() 异常处理 | 流式失败降级到非流式 |
| 观察者模式 | SSE 事件推送(routes.py 的 event_generator) | 服务端推送事件,客户端订阅消费 |
| 白名单模式 | src/tools/device_restart.py:15-20 | 只允许预定义的命令执行 |
| 白名单模式 | src/tools/log_reader.py:25 | keyword 输入校验 |
| 多级缓存 | src/utils/session_store.py:100-120 RedisSessionStore | 内存缓存 → Redis → 新建 |
| Supervisor 模式 | src/agent/supervisor.py | 路由请求到不同专业 Agent |
第七章:面试通关指南
7.1 项目介绍话术
1 分钟版
"我做的是一个 AI 智能设备故障诊断系统,部署在 IoT 智能饮吧设备上。运维人员描述故障现象后,AI Agent 会自动查设备状态、读日志、检索知识库,最终给出诊断报告和维修建议。
技术上,用 RAG 让 AI 基于设备手册和故障案例回答问题,用 Agent + 工具调用 让 AI 能真实操作设备(查状态、读日志、重启服务)。
后端用 FastAPI + LangChain + LangGraph,支持 Claude 和 OpenAI 双模型自动降级。检索用混合检索(向量 + BM25 + RRF 融合 + CrossEncoder 重排序),前端 React,通过 SSE 流式输出。"
3 分钟版
(1 分钟版 +)
"具体来说,系统分三层:
第一层 RAG 知识库:加载 19 篇设备手册和故障案例,经过 Markdown 标题感知 + 递归字符两级切分,用本地 Embedding 模型向量化后存入 ChromaDB。检索时用混合检索——向量语义检索找’意思相近的’,BM25 关键词检索找’包含特定错误码的’,再用 RRF 算法融合两路排名,最后用 CrossEncoder 精排。
第二层 Agent 系统:三个专业 Agent——诊断 Agent、维修 Agent、监控 Agent,通过 Supervisor 路由编排。每个 Agent 使用 ReAct 模式(思考→调用工具→观察→继续思考),能调用知识库搜索、设备状态查询、日志读取、固件版本检查、服务重启等 6 个工具。
第三层 API 服务:FastAPI 提供 RESTful + SSE 接口,支持流式输出让前端实时显示 Agent 推理过程。会话通过 Redis 持久化,Redis 不可用自动降级到内存。所有接口有限流保护。"
5 分钟版
(3 分钟版 +)
"我再深入讲两个技术亮点:
混合检索的设计。纯向量检索对精确关键词(如错误码 E03)不敏感,纯 BM25 不理解语义。我们两路并行检索各取 top 20,用 RRF(Reciprocal Rank Fusion)融合排名——这个算法的好处是不依赖具体分数值(向量距离和 BM25 分数的量纲不同),只看排名位置。融合后再用 CrossEncoder 精排出 top 5。实测 Recall@5 比纯向量检索提升了约 15%。
多 Agent 协作。Supervisor 采用双路由策略:关键词路由零延迟零成本,用于明确意图的请求;LLM 路由更智能,用于复杂意图的请求。对于’先诊断再给维修方案’这类复杂需求,诊断 Agent 的输出直接作为维修 Agent 的输入,链式串联工作。
安全方面,工具都有输入校验(正则白名单防命令注入)、命令白名单(重启只能执行预定义的命令)、超时控制。LLM 调用有自动降级,Redis 有降级,流式输出也有降级到非流式的兜底。
项目从 0 到 1 我都参与了,包括架构设计、RAG 管道优化、Agent 编排、API 开发和前端对接。"
7.2 技术深度问题 Top 20
RAG 相关
Q1:什么是 RAG?为什么需要 RAG?
“RAG 是 Retrieval-Augmented Generation,检索增强生成。LLM 的知识截止到训练时间,不知道企业内部的文档和最新信息。RAG 的做法是先从知识库中检索出与问题相关的文档片段,把它们作为上下文注入到 prompt 中,让 LLM 基于这些事实来回答。这样既利用了 LLM 的理解和生成能力,又保证了回答基于真实文档,减少了’幻觉’。”
Q2:Embedding 的原理是什么?
“Embedding 是把文本映射到高维向量空间的过程。模型(如 all-MiniLM-L6-v2)将文本编码为 384 维的浮点数向量。训练时,语义相似的文本被映射到距离较近的向量位置。检索时,把查询也编码为向量,然后在向量数据库中做最近邻搜索(如余弦相似度),找到语义最接近的文档片段。”
Q3:向量数据库和传统数据库的区别?
“传统数据库(MySQL)基于精确匹配(WHERE name = ‘xxx’)或模糊匹配(LIKE ‘%xxx%’)。向量数据库基于相似度搜索——给一个查询向量,返回距离最近的 N 个向量。底层使用 ANN(近似最近邻)算法(如 HNSW、IVF)做高效索引。适用场景不同:传统数据库适合结构化数据查询,向量数据库适合语义检索。”
Q4:你的 chunk_size 和 chunk_overlap 怎么调优的?
“通过评估脚本(
evaluation/run_eval.py)做量化评估。准备了包含标准问答对的测试集,指定每个问题的期望来源文档和关键词。分别测试 chunk_size 从 500 到 1500 的效果,看 Recall@5、MRR 和关键词覆盖率。最终选择 chunk_size=800、chunk_overlap=150,在召回率和精确度之间取得了最佳平衡。”
Q5:BM25 的原理是什么?和 TF-IDF 有什么区别?
“BM25 是 TF-IDF 的改进版。TF-IDF 的问题是词频(TF)没有上界——一个词出现 100 次比 50 次得分翻倍。BM25 引入了饱和函数,高频词的增益递减(出现 100 次和 50 次差别不大)。还考虑了文档长度归一化——长文档天然包含更多关键词,BM25 通过文档长度参数 b 来调整。”
Agent 相关
Q6:ReAct 模式和 Chain of Thought (CoT) 有什么区别?
“CoT 是’思考链’,让 LLM 一步步推理出答案,但只能用自己的知识。ReAct 在 CoT 基础上加了’行动’——LLM 不仅能思考,还能调用工具获取外部信息。ReAct = Reasoning(推理)+ Acting(行动),Thought → Action → Observation 的循环。”
Q7:LangGraph 和 LangChain 的关系?
“LangChain 是 AI 应用开发框架,提供 LLM 调用、文档处理、向量存储等基础组件。LangGraph 是 LangChain 团队推出的 Agent 编排框架,基于有向图(DAG)来定义 Agent 的工作流。
create_react_agent是 LangGraph 提供的高级 API,内部构建了一个 ReAct 循环图。”
Q8:Agent 可能出现死循环或调用过多工具怎么办?
“LangGraph 的
create_react_agent内置了recursion_limit(默认 25 轮),超过限制会自动停止。此外,通过 System Prompt 约束 Agent 行为(‘每次诊断至少使用 2-3 个工具’而不是’使用所有工具’),减少不必要的工具调用。生产环境还可以在 Agent 外层加超时控制。”
Q9:Function Calling 和 Tool Use 的区别?
“本质上是同一个概念,只是不同厂商的叫法不同。OpenAI 叫 Function Calling,Anthropic 叫 Tool Use。原理都是:把工具的名称、描述、参数 schema 告诉 LLM,LLM 在回复中返回要调用的工具名和参数,应用层执行工具后把结果反馈给 LLM。”
架构相关
Q10:为什么选 FastAPI 而不是 Flask/Django?
"FastAPI 的优势:
- 原生 async 支持,适合 I/O 密集的 LLM 调用场景
- 基于 Pydantic 的类型系统,自动请求校验和 Swagger 文档生成
- 原生
StreamingResponse支持 SSE 流式输出- 性能好(基于 Starlette 和 uvicorn)
Flask 缺乏原生 async 支持,Django 太重了不适合 AI 微服务。"
Q11:你的项目怎么处理并发请求?
“FastAPI 基于 ASGI(异步网关接口),用 uvicorn 作为服务器,单进程就能处理大量并发 I/O。LLM 调用和 Agent 执行通过
run_in_executor放到线程池执行,不阻塞事件循环。限流通过 slowapi 做 IP 级限制。会话存储用 Redis 支持多实例共享。”
Q12:ChromaDB 和 Milvus 怎么切换的?
“通过
VECTOR_DB_TYPE环境变量一键切换。代码层面,create_vector_store()根据配置调用_create_chroma_store()或_create_milvus_store(),两者都实现 LangChain 的 VectorStore 接口,上层代码完全无感知。这是策略模式的应用。”
Q13:你的系统有哪些降级策略?
"三层降级:
- LLM 降级:Claude 失败自动切换到 OpenAI(
with_fallbacks)- 存储降级:Redis 不可用自动降级到内存存储
- 输出降级:流式输出异常自动降级到非流式
每一层都有日志记录降级事件,方便监控和排查。"
Q14:SSE 和 WebSocket 有什么区别?为什么选 SSE?
“SSE 是单向的(服务端 → 客户端),WebSocket 是双向的。AI 对话场景中,用户发送问题后只需要接收回答,不需要在接收过程中发送新消息,所以单向的 SSE 就够了。SSE 基于 HTTP,比 WebSocket 更轻量——不需要额外的握手协议、不需要保持长连接、可以被 CDN 和反向代理自然支持。”
Q15:你的 Prompt 工程有哪些经验?
"几个原则:
- 角色定义:明确告诉 LLM 它是什么角色(设备诊断助手/搜索查询优化器/路由器)
- 能力声明:列出所有可用工具及其用途
- 工作流程:给出具体的操作步骤,引导 LLM 按流程工作
- 行为约束:如’至少用 2-3 个工具’、‘不要猜测’、‘基于文档回答’
- 输出格式:指定期望的输出格式(如’只输出 Agent 名称,不要解释’)
工具的 description 也很重要,要明确写’什么场景下用这个工具’,直接影响 Agent 的决策质量。"
Python 特性相关
Q16:Python 的 yield 和 Java 的 Iterator 有什么区别?
“概念类似但实现方式不同。Java 的
Iterator需要自己维护状态(实现hasNext()和next()方法)。Python 的yield是语言级支持——函数执行到yield时暂停并返回值,下次调用时从暂停处继续执行。不需要手动维护状态,代码更简洁。”
Q17:Python 的 @tool 装饰器和 Java 的注解有什么区别?
“Java 的注解是元数据标记,需要反射 + AOP 才能执行逻辑。Python 的装饰器是函数级别的语法糖,本质是’函数包装函数’——
@tool会把原函数包装成一个Tool对象,自动从函数名提取 name、从 docstring 提取 description、从类型注解提取参数 schema。不需要反射,更直观。”
Q18:为什么选 httpx 而不是 requests?
“httpx 支持 async(
httpx.AsyncClient),在 FastAPI 异步环境中可以非阻塞地调用外部 API。requests 只支持同步调用。此外 httpx API 和 requests 高度兼容,迁移成本很低。”
实践相关
Q19:你在开发过程中遇到过什么挑战?
"几个典型的:
- 中文 Embedding 效果:英文模型对中文效果一般,通过引入 BM25 混合检索弥补
- 中转站兼容性:不同中转站返回的 content 格式不一致(有的返回 string,有的返回 list),写了
extract_text()做格式兼容- 代理冲突:开发环境的 Clash 代理与中转站 SSL 冲突,统一在 settings.py 清除
- Agent 工具选择:最初 Agent 经常选错工具或不调用工具,通过优化 System Prompt 和工具 description 解决"
Q20:这个项目还有什么可以优化的?
"几个方向:
- 换更好的中文 Embedding 模型(如 bge-small-zh-v1.5)提升向量检索效果
- 引入 jieba 分词替代字符级 2-gram,提升 BM25 检索质量
- Agent 执行历史追踪:记录每次诊断的工具调用链,用于事后分析和优化
- 对话历史限制:history 无限增长会超过 LLM 上下文窗口,可以引入滑动窗口或摘要机制
- 多用户并发:当前 BM25 索引是全局单例,高并发下可能有性能瓶颈
- 评估自动化:把 RAG 评估集成到 CI/CD,每次知识库变更自动回归测试"
7.3 常见追问与应对策略
面试官追问:“你说的这个降级策略,在实际生产中触发过吗?”
“在开发测试阶段确实触发过。有一次 Claude API 中转站维护,触发了 LLM 降级到 OpenAI。日志中记录了降级事件,服务对用户完全透明。Redis 降级也在 Docker 环境中验证过——故意不启动 Redis 容器,会话存储自动降级到内存模式,日志输出警告信息。”
面试官追问:“你的 Agent 会不会产生不靠谱的回答?”
"会,这是 LLM 的固有问题。我们通过几个手段减轻:
- RAG 基于文档回答,Prompt 约束’优先基于参考文档’
- Agent 工具返回真实数据(设备状态、日志),减少 LLM 编造信息的空间
- 前端展示工具调用过程,用户可以验证 Agent 的推理是否合理
- 最终用户是专业运维人员,有判断能力"
面试官追问:“如果知识库文档更新了怎么办?”
“重新运行
build_knowledge_base.py脚本。这个脚本会重新加载所有文档、切分、Embedding、存入向量数据库。目前是全量重建,未来可以优化为增量更新(检测文档变化,只重新处理变化的文件)。”
面试官追问:“你是一个人做的还是团队合作?”
“架构设计和核心代码(RAG 管道、Agent 编排、API 层)是我负责的。AI 辅助了部分代码生成和文档编写,但技术选型、架构设计、算法选择、代码 review 和调试优化都是我做的。前端部分相对简单,是一个单页面 React 应用。”
最后提醒:面试时不要背诵,要用自己的话说。本文档的参考回答只是思路提示,请根据自己的实际理解重新组织语言。最有说服力的是"我在开发中遇到了 XX 问题,尝试了 AA 和 BB 方案,最终选择了 BB 因为…"这样的真实经验叙述。
更多推荐



所有评论(0)