一、系列回顾与本篇定位
① 篇:ESP32 → MQTT → 云端(能通)
② 篇:ESP32 → MQTT → Flink → 实时管道 + 异常检测(生产级)
③ 篇:② + 异常时调 LLM 生成诊断(会思考)
④ 篇:③ + 向量库存历史诊断,检索增强(有记忆) ← 本篇
③ 篇的瓶颈很明显:LLM 每次都从零开始推理。
[温度漂移异常] → LLM:基于通用知识猜原因 → 诊断报告
问题是:工厂里"3 号机组温度漂移"上周刚处理过,根因是冷却泵滤网堵塞。LLM 不知道这件事。④ 篇解决的就是"让系统记住经验"。
二、整体架构
在 ③ 篇基础上加一个向量库和检索环节:
ESP32 ──MQTT──▶ Flink ──▶ 异常检测 ──▶ ┌─────────────────────┐ │ RAG 诊断算子 │ │ 1. 异常 → embedding │ │ 2. 检索 Top-K 历史 │ │ 3. 拼 prompt 调 LLM │ └─────────┬───────────┘ │ ┌────────────────────▼────────────┐ │ Milvus 向量库(历史诊断记忆) │ │ (异常向量, 根因, 处置, 时间戳) │ └─────────────────────────────────┘ │ ┌────────────────────▼────────────┐ │ LLM 诊断报告(含参考案例)→ 库/钉钉│ └─────────────────────────────────┘关键:每次诊断完,把本次异常向量 + 最终根因写回 Milvus,形成正向循环。
三、向量库选型与建表
为什么用 Milvus(而不是简单的内存向量检索):历史诊断会累积到百万级,需要持久化 + 高效 ANN 检索 + 可按设备/时间过滤。
from pymilvus import MilvusClient client = MilvusClient("http://milvus:19530") client.create_collection( collection_name="iot_diagnosis", dimension=384, # bge-small-zh 的维度 metric_type="COSINE", auto_id=True, ) # 加标量字段,支持"只检索同设备型号的历史" client.alter_collection( collection_name="iot_diagnosis", properties={"enable_dynamic_field": True} )Embedding 模型选bge-small-zh(384 维):在设备日志/告警文本上效果足够,体积小、推理快(CPU 即可 1200 句/s)。
四、Flink 侧:写入向量 + 检索增强
4.1 异常 → 向量,并检索历史
public class RAGDiagnosisFunction extends RichAsyncFunction<Alert, DiagnosisReport> { private transient MilvusServiceClient milvus; private transient EmbeddingClient embed; // bge-small 本地服务 @Override public void asyncInvoke(Alert a, ResultFuture<DiagnosisReport> f) { // 1. 异常文本 → 向量 float[] vec = embed.embed(buildAlertText(a)); // 2. 检索同型号设备的 Top-3 历史案例 SearchParam sp = SearchParam.newBuilder() .withCollectionName("iot_diagnosis") .withVector(vec) .withTopK(3) .withExpr("device_type == \"" + a.deviceType + "\"") .build(); SearchResults res = milvus.search(sp); // 3. 拼检索到的历史案例进 prompt String hist = formatHistory(res); String prompt = buildPrompt(a, hist); // 4. 调 LLM 生成(同 ③ 篇的 vLLM 调用) callLLM(prompt).thenAccept(report -> { // 5. 诊断完写回向量库,形成记忆 milvus.insert(InsertParam.newBuilder() .withCollectionName("iot_diagnosis") .withFields(buildFields(a, vec, report)) .build()); f.complete(Collections.singletonList(report)); }); } private String buildPrompt(Alert a, String hist) { return "你是工业设备运维助手。当前异常:" + a.type + ",数值=" + a.value + ",设备型号=" + a.deviceType + "。\n历史相似案例及根因:\n" + hist + "\n请参考历史案例判断本次根因并给处置建议,不超过 80 字。"; } }4.2 接入主流水线
SingleOutputStreamOperator<DiagnosisReport> reports = AsyncDataStream .unorderedWait(alerts, new RAGDiagnosisFunction(), 30, TimeUnit.SECONDS, 20) .name("rag-diagnosis"); reports.addSink(new DiagnosisSink());五、Embedding 服务(轻量)
# embedding_service.py — 本地 bge-small,被 Flink 通过 HTTP 调 from sentence_transformers import SentenceTransformer from fastapi import FastAPI model = SentenceTransformer("BAAI/bge-small-zh") app = FastAPI() @app.post("/embed") async def embed(req: dict): v = model.encode(req["text"]).tolist() return {"vector": v}CPU 上单条告警 embedding 约 8ms,对实时管道无感。
六、实测数据
环境:ESP32×3 + Flink(2 并行) + Milvus(单机) + vLLM(Llama-3-8B, 1×A100)。 先运行 2 周累积 1500 条历史诊断,再测新异常。
| 指标 | ③ 篇(无记忆) | ④ 篇(RAG 增强) | 提升 |
|---|---|---|---|
| 根因命中率(人工复核) | 71% | 93% | +22pt |
| 平均诊断字数 | 52 | 68(更具体) | — |
| 诊断延迟 P50 | 380 ms | 420 ms(+检索 40ms) | 可接受 |
| 峰值吞吐 | 5 万 msg/s | 5 万 msg/s | 不变 |
| 向量库规模 | — | 1500→持续增长 | 记忆累积 |
RAG 让 LLM 站在"历史经验"的肩膀上,根因命中率提升显著,且延迟仅多 40ms。
七、踩坑记录
| 问题 | 现象 | 解决 |
|---|---|---|
| 检索到不相关案例 | 根因被带偏 | 加device_type标量过滤,只查同型号 |
| 向量维度不匹配 | 写入失败 | 确认 embedding 模型维度=建表 dimension |
| 历史太少检索无意义 | 前期命中低 | 冷启动先用 ③ 篇规则兜底,积累 500 条后开 RAG |
| 写入和检索争抢 | 延迟抖动 | Milvus 读写分离 + 批量插入 |
| 旧案例误导(根因已修正) | 错误传承 | 加verified字段,只检索人工确认过的案例 |
| 向量库膨胀 | 存储暴涨 | 按时间 TTL,保留近 90 天 |
八、总结
本篇让 Edge AI 全栈真正"有记忆":
硬件采集(ESP32)→ 大数据流处理(Flink)→ AI 推理(vLLM)+ 检索增强(Milvus/RAG)四线闭环
历史诊断持续累积成"经验库",LLM 诊断根因命中率 71% → 93%
标量过滤保证检索相关性,冷启动用规则兜底
硬件成本仍 ¥66,云端 Milvus + vLLM 可复用
下篇预告:⑤ 篇我们做"端侧轻量化"——把诊断模型本身下沉到边缘(ESP32-S3 + TinyML),让简单异常在设备端就地判断、只把疑难杂症上云,引出"云边协同"的工程范式,给整个系列收尾。
往期回顾:
Edge AI 全栈1:ESP32 传感器采集到云端大模型推理的完整链路
Edge AI 全栈 2:ESP32 + MQTT + Flink 实时 IoT 数据管道
Edge AI 全栈实战 3:Flink 异常检测联动云端大模型生成诊断报告
RAG 架构设计 7 个关键决策:从 Chunk 策略到 Reranker 的生产级方案