
一、系列回顾与本篇定位① 篇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_nameiot_diagnosis, dimension384, # bge-small-zh 的维度 metric_typeCOSINE, auto_idTrue, ) # 加标量字段支持只检索同设备型号的历史 client.alter_collection( collection_nameiot_diagnosis, properties{enable_dynamic_field: True} )Embedding 模型选bge-small-zh384 维在设备日志/告警文本上效果足够体积小、推理快CPU 即可 1200 句/s。四、Flink 侧写入向量 检索增强4.1 异常 → 向量并检索历史public class RAGDiagnosisFunction extends RichAsyncFunctionAlert, DiagnosisReport { private transient MilvusServiceClient milvus; private transient EmbeddingClient embed; // bge-small 本地服务 Override public void asyncInvoke(Alert a, ResultFutureDiagnosisReport 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 接入主流水线SingleOutputStreamOperatorDiagnosisReport 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平均诊断字数5268更具体—诊断延迟 P50380 ms420 ms检索 40ms可接受峰值吞吐5 万 msg/s5 万 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 全栈1ESP32 传感器采集到云端大模型推理的完整链路Edge AI 全栈 2ESP32 MQTT Flink 实时 IoT 数据管道Edge AI 全栈实战 3Flink 异常检测联动云端大模型生成诊断报告RAG 架构设计 7 个关键决策从 Chunk 策略到 Reranker 的生产级方案