🧑 博主简介CSDN博客专家「历代文学网」(PC端可以访问:https://lidaiwenxue.com/#/?__c=1000,移动端可关注公众号 “ 心海云图 ” 微信小程序搜索“历代文学”)总架构师,首席架构师,也是联合创始人!16年工作经验,精通Java编程高并发设计分布式系统架构设计Springboot和微服务,熟悉LinuxESXI虚拟化以及云原生Docker和K8s,热衷于探索科技的边界,并将理论知识转化为实际应用。保持对新技术的好奇心,乐于分享所学,希望通过我的实践经历和见解,启发他人的创新思维。在这里,我希望能与志同道合的朋友交流探讨,共同进步,一起在技术的世界里不断学习成长。
🤝商务合作:请搜索或扫码关注微信公众号 “ 心海云图

在这里插入图片描述

在这里插入图片描述

从零构建 Java 智能体 RAG 系统:Milvus 向量数据库实战指南


一、引言

RAG(检索增强生成)通过为大语言模型提供外部知识库,解决了模型知识截止日期和“幻觉”问题。Milvus 作为 CNCF 毕业的向量数据库,提供专业的向量索引(HNSW、IVF_FLAT 等),可实现百亿数据的毫秒级检索。

本文档将带你使用 Java + Spring Boot + Milvus 从零构建一个企业级 RAG 智能问答系统。


二、环境准备

2.1 安装 Milvus(推荐 Docker 方式)

# 下载 docker-compose.yml
wget https://github.com/milvus-io/milvus/releases/download/v2.6.22/milvus-standalone-docker-compose.yml -O docker-compose.yml

# 启动 Milvus
sudo docker-compose up -d

# 验证服务(默认端口 19530)
docker ps

2.2 引入 Java SDK 依赖

Milvus Java SDK 要求 Java 8 或更高版本。根据 Milvus 版本选择对应的 SDK 版本:

Milvus 版本 Java SDK 版本
2.6.x 2.6.22
3.0.x 3.0.6

Maven 依赖:

<dependency>
    <groupId>io.milvus</groupId>
    <artifactId>milvus-sdk-java</artifactId>
    <version>2.6.22</version>
</dependency>

⚠️ 从 v2.5.2 起,SDK 拆分为两个包:milvus-sdk-javamilvus-sdk-java-bulkwriter。如不需要批量写入工具,仅引入第一个即可。


三、Milvus Java SDK 核心操作

💡 版本说明:新版 SDK 的核心客户端类是 MilvusClientV2,各语言 SDK 采用统一的 API 结构。

3.1 连接 Milvus 服务

import io.milvus.v2.client.ConnectConfig;
import io.milvus.v2.client.MilvusClientV2;
import io.milvus.v2.service.collection.response.ListCollectionsResp;
import io.milvus.v2.service.database.response.ListDatabasesResp;

public class MilvusConnection {
    
    private static MilvusClientV2 client;
    
    public static MilvusClientV2 connect() {
        ConnectConfig connectConfig = ConnectConfig.builder()
            .uri("http://localhost:19530")      // Milvus 服务地址
            .username("root")                    // 可选
            .password("Milvus")                  // 可选
            .dbName("default")                   // 可选,默认数据库
            .build();
        
        MilvusClientV2 client = new MilvusClientV2(connectConfig);
        
        // 验证连接:列出所有数据库
        ListDatabasesResp databases = client.listDatabases();
        System.out.println("已连接,数据库列表: " + databases.getDatabaseNames());
        
        // 列出所有集合
        ListCollectionsResp collections = client.listCollections();
        System.out.println("集合列表: " + collections.getCollectionNames());
        
        return client;
    }
}

3.2 创建集合(Collection)

集合类似于关系型数据库中的“表”。

import io.milvus.v2.common.DataType;
import io.milvus.v2.common.IndexParam;
import io.milvus.v2.service.collection.request.AddFieldReq;
import io.milvus.v2.service.collection.request.CreateCollectionReq;
import io.milvus.v2.service.collection.request.CreateIndexReq;

public class CollectionManager {
    
    public static void createRagCollection(MilvusClientV2 client, String collectionName) {
        // 1. 创建 Schema(模式)
        CreateCollectionReq.CollectionSchema schema = client.createSchema();
        
        // 2. 添加字段
        // 主键字段:id
        schema.addField(AddFieldReq.builder()
            .fieldName("id")
            .dataType(DataType.Int64)
            .isPrimaryKey(true)
            .autoID(true)        // 自动生成 ID
            .build());
        
        // 文本字段:存储原始文档内容
        schema.addField(AddFieldReq.builder()
            .fieldName("text")
            .dataType(DataType.VarChar)
            .maxLength(65535)
            .build());
        
        // 向量字段:存储文本的向量表示(维度需与 embedding 模型一致)
        schema.addField(AddFieldReq.builder()
            .fieldName("embedding")
            .dataType(DataType.FloatVector)
            .dimension(768)      // 以 nomic-embed-text 为例
            .build());
        
        // 元数据字段:文档来源
        schema.addField(AddFieldReq.builder()
            .fieldName("source")
            .dataType(DataType.VarChar)
            .maxLength(512)
            .build());
        
        // 3. 创建集合
        CreateCollectionReq createReq = CreateCollectionReq.builder()
            .collectionName(collectionName)
            .collectionSchema(schema)
            .build();
        
        client.createCollection(createReq);
        System.out.println("✅ 集合创建成功: " + collectionName);
        
        // 4. 创建索引(提升检索性能)
        IndexParam indexParam = IndexParam.builder()
            .fieldName("embedding")
            .indexType(IndexParam.IndexType.HNSW)   // HNSW 索引,适合高召回率场景
            .metricType(IndexParam.MetricType.COSINE)
            .extraParams(Map.of("M", 16, "efConstruction", 256))
            .build();
        
        CreateIndexReq indexReq = CreateIndexReq.builder()
            .collectionName(collectionName)
            .indexParams(Collections.singletonList(indexParam))
            .build();
        
        client.createIndex(indexReq);
        System.out.println("✅ 索引创建成功");
    }
}

索引类型选择建议

索引类型 适用场景 特点
HNSW 高召回率、低延迟 查询快,但构建索引耗时,内存占用较高
IVF_FLAT 大规模数据、平衡性能 需要调整 nlist/nprobe 参数
IVF_SQ8 存储受限场景 量化压缩,精度略降

3.3 插入向量数据

import com.google.gson.Gson;
import com.google.gson.JsonObject;
import io.milvus.v2.service.vector.request.UpsertReq;
import io.milvus.v2.service.vector.response.UpsertResp;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;

public class DataIngestion {
    
    // 模拟 embedding 调用(实际应调用 embedding 模型 API)
    private static float[] embedText(String text) {
        // 实际项目中调用 OpenAI、千帆等 embedding 接口
        // 此处返回随机向量作为示例
        float[] vector = new float[768];
        for (int i = 0; i < 768; i++) {
            vector[i] = (float) Math.random();
        }
        return vector;
    }
    
    public static void insertDocuments(MilvusClientV2 client, String collectionName) {
        // 准备文档数据
        List<Map<String, Object>> documents = List.of(
            Map.of("text", "Spring AI 是一个用于简化 AI 应用开发的 Spring 框架扩展。"),
            Map.of("text", "Milvus 是一款开源的向量数据库,专为海量向量数据的存储和相似性搜索而设计。"),
            Map.of("text", "RAG(检索增强生成)将信息检索与大语言模型生成相结合。")
        );
        
        Gson gson = new Gson();
        List<JsonObject> dataList = new ArrayList<>();
        
        for (Map<String, Object> doc : documents) {
            String text = (String) doc.get("text");
            float[] embedding = embedText(text);
            
            Map<String, Object> row = Map.of(
                "text", text,
                "embedding", embedding,
                "source", "demo_knowledge_base"
            );
            
            dataList.add(gson.toJsonTree(row).getAsJsonObject());
        }
        
        // 批量插入(Upsert:存在则更新,不存在则插入)
        UpsertReq upsertReq = UpsertReq.builder()
            .collectionName(collectionName)
            .data(dataList)
            .build();
        
        UpsertResp response = client.upsert(upsertReq);
        System.out.println("✅ 成功插入 " + response.getUpsertCnt() + " 条数据");
    }
}

3.4 向量检索

import io.milvus.v2.service.vector.request.SearchReq;
import io.milvus.v2.service.vector.response.SearchResp;

public class VectorSearch {
    
    public static List<SearchResp.SearchResult> search(
            MilvusClientV2 client, 
            String collectionName, 
            String queryText,
            int topK) {
        
        // 1. 将查询文本向量化
        float[] queryVector = embedText(queryText);  // 复用上面的 embedText 方法
        FloatVec queryVec = new FloatVec(queryVector);
        
        // 2. 构建搜索请求
        SearchReq searchReq = SearchReq.builder()
            .collectionName(collectionName)
            .data(Collections.singletonList(queryVec))
            .annsField("embedding")              // 指定向量字段
            .outputFields(Arrays.asList("text", "source"))  // 返回的字段
            .topK(topK)
            .build();
        
        // 3. 执行搜索
        SearchResp searchResp = client.search(searchReq);
        List<List<SearchResp.SearchResult>> results = searchResp.getSearchResults();
        
        // 4. 输出结果
        for (List<SearchResp.SearchResult> resultList : results) {
            System.out.println("Top " + topK + " 结果:");
            for (SearchResp.SearchResult result : resultList) {
                System.out.println("  - 相似度: " + result.getScore());
                System.out.println("    文本: " + result.getEntity().get("text"));
            }
        }
        
        return results.isEmpty() ? Collections.emptyList() : results.get(0);
    }
}

四、构建完整的 RAG 系统

4.1 技术栈选型

组件 推荐方案
后端框架 Spring Boot 3.x
AI 集成 Spring AI 1.0.0
向量数据库 Milvus 2.6.x
LLM OpenAI API 兼容接口(百度千帆、Ollama 等)

4.2 Spring Boot 配置

application.yml:

spring:
  ai:
    openai:
      base-url: http://localhost:11434    # Ollama 或千帆 API 地址
      chat:
        options:
          model: qwen2.5:7b
      embedding:
        options:
          model: nomic-embed-text
          dimensions: 768

milvus:
  uri: http://localhost:19530
  username: root
  password: Milvus
  database: default
  collection: rag_knowledge

4.3 文档 ETL 服务

将 PDF、Word 等文档切分并存入 Milvus:

import org.springframework.ai.document.Document;
import org.springframework.ai.reader.tika.TikaDocumentReader;
import org.springframework.ai.splitter.TokenTextSplitter;
import org.springframework.ai.vectorstore.VectorStore;
import org.springframework.core.io.Resource;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Flux;

@Service
public class DocumentIngestionService {
    
    private final VectorStore vectorStore;
    private final TokenTextSplitter textSplitter;
    
    public DocumentIngestionService(VectorStore vectorStore) {
        this.vectorStore = vectorStore;
        this.textSplitter = new TokenTextSplitter();
    }
    
    public void ingestDocument(Resource resource) {
        // 1. 读取文档(支持 PDF、Word、TXT 等)
        TikaDocumentReader reader = new TikaDocumentReader(resource);
        
        // 2. 文本分块(控制上下文长度)
        List<Document> documents = textSplitter.apply(reader.read());
        
        // 3. 存入 Milvus(自动向量化并存储)
        vectorStore.add(documents);
        
        System.out.println("✅ 成功加载 " + documents.size() + " 个文档片段到知识库");
    }
}

4.4 RAG 核心服务

实现“检索 + 生成”的完整流程:

import org.springframework.ai.chat.client.ChatClient;
import org.springframework.ai.chat.prompt.PromptTemplate;
import org.springframework.ai.document.Document;
import org.springframework.ai.vectorstore.SearchRequest;
import org.springframework.ai.vectorstore.VectorStore;
import org.springframework.stereotype.Service;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;

@Service
public class RAGService {
    
    private final ChatClient chatClient;
    private final VectorStore vectorStore;
    
    public RAGService(ChatClient.Builder chatClientBuilder, VectorStore vectorStore) {
        this.chatClient = chatClientBuilder.build();
        this.vectorStore = vectorStore;
    }
    
    public String ask(String question) {
        // 1. 从 Milvus 检索最相关的文档片段
        SearchRequest searchRequest = SearchRequest.builder()
            .query(question)
            .topK(5)
            .build();
        
        List<Document> documents = vectorStore.similaritySearch(searchRequest);
        
        // 2. 提取检索结果作为上下文
        List<String> context = documents.stream()
            .map(Document::getText)
            .collect(Collectors.toList());
        
        // 3. 构建 Prompt
        PromptTemplate promptTemplate = new PromptTemplate("""
            请根据以下背景信息回答用户的问题。
            如果背景信息不包含答案,请说明"依据现有资料无法回答"。
            
            ### 背景信息:
            {context}
            
            ### 用户问题:
            {question}
            
            ### 回答:
            """);
        
        // 4. 调用 LLM 生成答案
        return chatClient.prompt(
            promptTemplate.create(Map.of(
                "context", String.join("\n", context),
                "question", question
            ))
        ).call().content();
    }
}

4.5 REST API 接口

import org.springframework.web.bind.annotation.*;
import org.springframework.beans.factory.annotation.Autowired;

@RestController
@RequestMapping("/api/rag")
public class RAGController {
    
    @Autowired
    private RAGService ragService;
    
    @PostMapping("/ask")
    public String askQuestion(@RequestBody String question) {
        return ragService.ask(question);
    }
}

五、AI 智能体增强

5.1 从 RAG 到智能体记忆

传统 RAG 是“只读的外部知识库”,而 AI 智能体需要动态记忆能力——记住历史对话、积累上下文。

Milvus 可以存储两类记忆:

  • 短期记忆:当前会话的对话历史
  • 长期记忆:跨会话的用户偏好和历史知识

5.2 对话记忆存储

@Service
public class ConversationMemoryService {
    
    private final MilvusClientV2 client;
    private final String memoryCollection = "conversation_memory";
    
    public void saveConversation(String sessionId, String userMessage, String assistantMessage) {
        // 将对话内容向量化并存储
        float[] userVector = embedText(userMessage);
        float[] assistantVector = embedText(assistantMessage);
        
        JsonObject memory = new JsonObject();
        memory.addProperty("session_id", sessionId);
        memory.addProperty("user_message", userMessage);
        memory.addProperty("assistant_message", assistantMessage);
        memory.add("user_embedding", gson.toJsonTree(userVector));
        memory.add("assistant_embedding", gson.toJsonTree(assistantVector));
        memory.addProperty("timestamp", System.currentTimeMillis());
        
        UpsertReq req = UpsertReq.builder()
            .collectionName(memoryCollection)
            .data(Collections.singletonList(memory))
            .build();
        
        client.upsert(req);
    }
    
    public List<String> retrieveRelevantMemories(String sessionId, String currentQuery, int topK) {
        // 检索与当前问题相关的历史对话
        float[] queryVector = embedText(currentQuery);
        
        SearchReq searchReq = SearchReq.builder()
            .collectionName(memoryCollection)
            .data(Collections.singletonList(new FloatVec(queryVector)))
            .annsField("user_embedding")
            .filter("session_id == '" + sessionId + "'")  // 仅检索当前会话
            .outputFields(Arrays.asList("user_message", "assistant_message"))
            .topK(topK)
            .build();
        
        SearchResp resp = client.search(searchReq);
        // 解析并返回历史对话
        return extractMessages(resp);
    }
}

5.3 Hybrid RAG(混合检索)

Milvus 2.6 支持混合检索——结合语义搜索(密集向量)和关键词搜索(稀疏向量 / BM25),提升检索准确率。

// 混合检索示例(需要创建稀疏向量字段)
public SearchResp hybridSearch(MilvusClientV2 client, String collectionName, 
                                String query, float[] denseVector, float[] sparseVector) {
    // 同时执行密集向量检索和稀疏向量检索
    // 然后通过加权融合(如 RRF)合并结果
    
    SearchReq denseSearch = SearchReq.builder()
        .collectionName(collectionName)
        .data(Collections.singletonList(new FloatVec(denseVector)))
        .annsField("dense_embedding")
        .topK(10)
        .build();
    
    SearchReq sparseSearch = SearchReq.builder()
        .collectionName(collectionName)
        .data(Collections.singletonList(new SparseFloatVec(sparseVector)))
        .annsField("sparse_embedding")
        .topK(10)
        .build();
    
    // 分别执行后做结果融合
    // ...
}

六、性能优化与生产实践

6.1 批量插入优化

生产环境中,批量插入应遵循以下原则:

  • 每批次控制在 2-5MB 数据量
  • 使用多线程并发插入时保证批次有序
  • 异常处理包含重试机制
public void batchInsertWithRetry(MilvusClientV2 client, String collectionName, 
                                  List<JsonObject> dataList, int batchSize) {
    int maxRetries = 3;
    for (int i = 0; i < dataList.size(); i += batchSize) {
        int end = Math.min(i + batchSize, dataList.size());
        List<JsonObject> batch = dataList.subList(i, end);
        
        int retry = 0;
        while (retry < maxRetries) {
            try {
                UpsertReq req = UpsertReq.builder()
                    .collectionName(collectionName)
                    .data(batch)
                    .build();
                client.upsert(req);
                break;
            } catch (Exception e) {
                retry++;
                if (retry >= maxRetries) {
                    throw new RuntimeException("批量插入失败,已重试 " + maxRetries + " 次", e);
                }
                System.out.println("重试第 " + retry + " 次...");
                Thread.sleep(1000 * retry);
            }
        }
    }
}

6.2 监控与调优

  • 使用 Attu(Milvus 官方 GUI 管理工具)监控集群状态
  • 关注查询延迟(P99)、系统负载等指标
  • 根据数据增长定期重建索引或调整索引参数

6.3 安全实践

  • 不要在代码或配置文件中硬编码 API KEY,使用环境变量
  • 生产环境开启 Milvus 的认证与授权功能
  • 敏感数据考虑字段级加密

七、完整项目结构

spring-ai-rag/
├── src/main/java/com/example/rag/
│   ├── config/
│   │   └── MilvusConfig.java          # Milvus 连接配置
│   ├── service/
│   │   ├── DocumentIngestionService.java  # 文档 ETL
│   │   ├── RAGService.java               # RAG 核心服务
│   │   └── ConversationMemoryService.java # 对话记忆
│   ├── controller/
│   │   └── RAGController.java            # REST API
│   └── Application.java
├── src/main/resources/
│   ├── application.yml
│   └── data/                            # 测试文档目录
└── pom.xml

💡 完整的可运行示例代码可参考官方示例仓库:milvus-sdk-java/examples以及 Spring AI RAG 示例 spring-ai-rag


总结

本文从零开始,使用 Java + Milvus 构建了一个完整的 AI 智能体 RAG 系统,涵盖:

  1. Milvus 环境搭建与 Java SDK 配置
  2. 集合管理:Schema 设计、索引创建
  3. 数据操作:向量插入、相似性检索
  4. RAG 核心链路:文档 ETL → 向量检索 → LLM 生成
  5. 智能体增强:对话记忆存储与检索
  6. 生产实践:批量优化、监控、安全
Logo

这里是“一人公司”的成长家园。我们提供从产品曝光、技术变现到法律财税的全栈内容,并连接云服务、办公空间等稀缺资源,助你专注创造,无忧运营。

更多推荐