尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

基于LangChain4j与pgvector的AI网盘:RAG问答与SSE实时通知实现

基于LangChain4j与pgvector的AI网盘:RAG问答与SSE实时通知实现 CloudVault 是一个偏工程向的 AI 网盘项目它在传统文件上传、下载、文件夹管理之外额外引入了 LangChain4j RAG让用户上传的文档可以直接被检索和提问而不是只能当成附件下载。这类系统在面试、毕设和内部知识库建设中都很常见但多数资料只写了文件 CRUD 或只写了 RAG Demo很少有人把 PostgreSQL、pgvector、Redis 和实时通知串成一条完整链路。下面把 CloudVault 的实现拆成文件服务、知识库索引、问答检索、缓存与通知四条线按可以落地跑通的顺序展开。最终你会得到一个能支持文档上传、自动切分、向量化、相似度检索、生成回答并通过 SSE 向用户推送处理进度的项目骨架。1. 先拆清楚 CloudVault 的两条主线文件服务与知识库问答1.1 仿网盘系统要解决的工程问题仿网盘系统的最小内核并不是“把文件存到磁盘”而是文件元数据管理。文件上传后至少要记录文件名、父目录、大小、MIME 类型、存储路径、上传人、上传时间。用户浏览文件列表时读的是这些元数据用户下载时系统通过存储路径找到真实文件。复杂网盘还会涉及分片上传、秒传、断点续传、权限控制、分享链接、文件预览。CloudVault 作为一个 AI 增强型网盘先要把最基础的文件表设计好否则后面做文档解析、向量入库、问答检索时会经常在“这个片段属于哪个文件”这件事上返工。文件服务这层还需要处理一个问题上传文件后系统不能立刻假设文件可以用于问答。PDF 可能是扫描件Word 可能有图片和表格MD 可能有代码块。所以文件状态要区分“上传完成”“解析中”“可问答”“解析失败”几个阶段。1.2 LangChain4j RAG 在 CloudVault 中的定位RAG 的全称是 Retrieval-Augmented Generation中文常称为检索增强生成。它的思路是用户提问时不把整份文档或全部知识库塞给大模型而是先从数据库中检索出和问题最相关的几个片段再把片段作为上下文交给大模型生成回答。LangChain4j 是 Java 生态中面向大模型应用的开源框架。它解决的问题是“Java 项目如何用一套统一 API 接入大模型、做文档切分、做向量存储、做检索和生成”。CloudVault 选择 LangChain4j是因为项目主体是 Spring Boot如果为了接大模型再引入一套 Python 服务会增加部署成本和调试成本。CloudVault 的具体 RAG 链路是用户上传文档 - 文档解析器抽文本 - DocumentSplitter 切分片段 - EmbeddingModel 生成向量 - pgvector 存储向量和文本 - 用户提问 - 问题向量化 - pgvector 相似度检索 - 拼入 Prompt - ChatModel 生成回答这一条链路定义了“文档问答”不是简单调用聊天接口而是先做知识召回再做生成。1.3 为什么配 PostgreSQL pgvector 和 Redispgvector 是 PostgreSQL 的向量类型扩展。它让 PostgreSQL 既能存业务表又能存向量数据不用在项目初期单独引入 Milvus、Weaviate 或 Qdrant。CloudVault 的文档片段和文件元数据天然有强关联用同一套 PostgreSQL 管理可以减少运维节点也方便按文件 ID 做权限过滤。Redis 在 CloudVault 里承担的是三件事缓存热点数据、做分布式锁、做跨实例通知分发。文件列表和文档问答结果都适合缓存文档解析任务需要防止同一文件被重复索引实时通知在单实例下可以直接用内存里的 SSE Emitter多实例时必须通过 Redis Pub/Sub 广播。如果后续知识库达到千万级向量、需要复杂的分区策略和多租户隔离再考虑迁移到 Milvus 这类专用向量库也是合理的。但 CloudVault 这类单体项目第一版更适合从 pgvector 起步。2. 技术选型和环境准备版本、依赖、数据库初始化2.1 推荐技术栈和版本基线下面表格是示例版本实际落地前要确认官方发布页的最新版本特别是 LangChain4j 的 API 在不同版本之间差异较大。技术组件示例版本用途备注JDK17运行 Spring Boot 和 LangChain4j生产建议使用 LTSSpring Boot3.2.xWeb、JPA、Redis、异步任务需要大版本一致LangChain4j0.33.x文档解析、切分、Embedding、AiServicesAPI 会演进以发布说明为准PostgreSQL16.x业务数据、向量数据需支持 pgvectorpgvector0.7.x向量类型和 HNSW 索引Windows 安装需要额外注意Redis7.x缓存、分布式锁、Pub/Sub单节点即可跑通 Demo学习环境和生产环境可以共用这套技术栈但生产环境建议使用 PostgreSQL 专用镜像pgvector/pgvector或者通过系统包管理器安装 pgvector避免只在本机能运行而部署到服务器时扩展装不上。2.2 初始化 PostgreSQL 与 pgvector先创建一个数据库CREATE DATABASE cloudvault;进入数据库后开启扩展\c cloudvault CREATE EXTENSION IF NOT EXISTS vector; SELECT extname, extversion FROM pg_extension;如果SELECT能查到vector说明 pgvector 已生效。然后创建文件表和向量表CREATE TABLE file_meta ( id BIGSERIAL PRIMARY KEY, parent_id BIGINT, name VARCHAR(255) NOT NULL, storage_path VARCHAR(512) NOT NULL, size_bytes BIGINT NOT NULL DEFAULT 0, mime_type VARCHAR(128), status SMALLINT NOT NULL DEFAULT 0, chunk_count INT NOT NULL DEFAULT 0, created_by VARCHAR(64), created_at TIMESTAMP NOT NULL DEFAULT now() ); CREATE TABLE doc_chunks ( id BIGSERIAL PRIMARY KEY, file_id BIGINT NOT NULL, chunk_index INT NOT NULL, text TEXT NOT NULL, embedding vector(1024) NOT NULL ); CREATE INDEX ON doc_chunks USING hnsw (embedding vector_cosine_ops);这里vector(1024)的维度必须和 EmbeddingModel 的输出维度一致。如果使用 qwen embedding 或 OpenAI 兼容接口通常是 1024如果使用本地 all-MiniLM-L6-v2 模型则可能是 384。维度不一致时入库和检索都会报错。HNSW索引适合中小规模数据查询性能好但构建时占内存。数据量很小时也可以先不建索引便于快速验证整条链路。2.3 Maven 依赖配置创建一个 Spring Boot 项目后引入以下核心依赖properties java.version17/java.version langchain4j.version0.33.0/langchain4j.version /properties dependencies dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-jpa/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-websocket/artifactId /dependency dependency groupIdorg.postgresql/groupId artifactIdpostgresql/artifactId scoperuntime/scope /dependency dependency groupIddev.langchain4j/groupId artifactIdlangchain4j/artifactId version${langchain4j.version}/version /dependency dependency groupIddev.langchain4j/groupId artifactIdlangchain4j-pgvector/artifactId version${langchain4j.version}/version /dependency dependency groupIddev.langchain4j/groupId artifactIdlangchain4j-open-ai/artifactId version${langchain4j.version}/version /dependency dependency groupIddev.langchain4j/groupId artifactIdlangchain4j-document-parser-apache-pdfbox/artifactId version${langchain4j.version}/version /dependency /dependencies如果项目的嵌入模型不走 OpenAI 兼容接口就不要引入langchain4j-open-ai改为实现EmbeddingModel接口。示例采用 OpenAI 兼容接口是为了方便对接 qwen embedding、Ollama 等本地或云端模型。使用哪个模型应该由配置决定而不是由代码写死。2.4 application.yml 核心配置在src/main/resources/application.yml中至少配置数据源、Redis、文件存储根目录和 AI 模型地址spring: datasource: url: jdbc:postgresql://localhost:5432/cloudvault username: postgres password: postgres jpa: hibernate: ddl-auto: validate open-in-view: false data: redis: host: localhost port: 6379 cloudvault: storage: root: ${CLOUDVAULT_STORAGE_ROOT:./data/files} ai: embedding: base-url: ${EMBEDDING_BASE_URL:http://localhost:8000/v1} api-key: ${EMBEDDING_API_KEY:local} model: ${EMBEDDING_MODEL:text-embedding-v3} dimension: ${EMBEDDING_DIMENSION:1024} chat: base-url: ${CHAT_BASE_URL:http://localhost:8000/v1} api-key: ${CHAT_API_KEY:local} model: ${CHAT_MODEL:qwen-plus}这里用环境变量覆盖默认值是为了避免把 API Key 提交到代码仓库。ddl-auto: validate表示 JPA 不自动改表结构数据库结构交给 SQL 迁移管理比如 Flyway。这样生产环境不会因为实体字段和表结构不一致而悄悄出错。3. 实现文件管理与 RAG 问答的闭环3.1 文件实体与上传接口文件表的实体可以这样设计Entity Table(name file_meta) public class FileMeta { Id GeneratedValue(strategy GenerationType.IDENTITY) private Long id; private Long parentId; private String name; private String storagePath; private Long sizeBytes; private String mimeType; private Integer status; private Integer chunkCount; private String createdBy; private LocalDateTime createdAt; // getter / setter }上传接口要做的不是立刻解析而是先保存文件元数据再异步触发解析。也就是“上传成功”和“可问答”是两个状态。RestController RequestMapping(/api/files) public class FileController { private final FileService fileService; private final AiDocumentIndexService aiDocumentIndexService; private final NotificationService notificationService; public FileController(FileService fileService, AiDocumentIndexService aiDocumentIndexService, NotificationService notificationService) { this.fileService fileService; this.aiDocumentIndexService aiDocumentIndexService; this.notificationService notificationService; } PostMapping(/upload) public ResultFileMeta upload(RequestParam(file) MultipartFile file, RequestParam(parentId) Long parentId, RequestParam(userId) String userId) throws IOException { FileMeta meta fileService.saveFile(file, parentId, userId); notificationService.publish(userId, file.uploaded, meta.getId().toString()); aiDocumentIndexService.indexAsync(meta.getId(), file.getBytes(), file.getOriginalFilename(), userId); return Result.ok(meta); } }这里选择先返回文件元数据再异步做 AI 索引。用户在上传大文件时不会因为向量化耗时而卡在请求里。notificationService会通知前端“文件已上传”“开始解析”“解析进度”“解析完成”。3.2 文档解析与文本切分文档问答的前提是抽取文本。PDF 示例使用ApachePdfBoxDocumentParser如果是 Markdown、TXT可以直接读取字符串如果需要处理 DOCX还要引入对应的 POI 或 Tika 解析器。public Document loadDocument(String fileName, byte[] bytes) throws IOException { if (fileName.endsWith(.pdf)) { DocumentParser parser new ApachePdfBoxDocumentParser(); return parser.parse(new ByteArrayInputStream(bytes)); } if (fileName.endsWith(.txt) || fileName.endsWith(.md)) { return Document.from(new String(bytes, StandardCharsets.UTF_8)); } throw new UnsupportedOperationException(暂不支持该文件类型: fileName); }文档抽取后不能直接整篇入库因为大模型上下文长度有限检索粒度太大会降低相关度。LangChain4j 提供DocumentSplitters.recursive可以按段落和句子递归切分。DocumentSplitter splitter DocumentSplitters.recursive(800, 200); ListTextSegment segments splitter.split(document);800是每个片段的字符数200是相邻片段之间的重叠字符数。重叠的作用是避免一个完整语义刚好被切断让检索时能覆盖到边界位置。对于技术文档推荐 500 到 1000 字之间对于代码为主的文档可以降到 300 到 600 字。3.3 使用 LangChain4j 完成向量化与存储先构建EmbeddingModel和PgVectorEmbeddingStoreEmbeddingModel embeddingModel OpenAiEmbeddingModel.builder() .baseUrl(baseUrl) .apiKey(apiKey) .modelName(modelName) .build(); EmbeddingStoreTextSegment embeddingStore PgVectorEmbeddingStore.builder() .host(localhost) .port(5432) .database(cloudvault) .user(postgres) .password(postgres) .table(doc_chunks) .dimension(embeddingDimension) .build();PgVectorEmbeddingStore会把TextSegment的文本和向量写入doc_chunks表。file_id和chunk_index需要额外放到TextSegment.metadata()中否则后续无法知道某条向量来自哪个文件。异步索引服务的核心逻辑Async(aiTaskExecutor) public void indexAsync(Long fileId, byte[] bytes, String fileName, String userId) { try { notificationService.publish(userId, index.start, fileId.toString()); Document document loadDocument(fileName, bytes); DocumentSplitter splitter DocumentSplitters.recursive(800, 200); ListTextSegment segments splitter.split(document); for (int i 0; i segments.size(); i) { TextSegment segment segments.get(i); segment.metadata().put(fileId, String.valueOf(fileId)); segment.metadata().put(chunkIndex, String.valueOf(i)); Embedding embedding embeddingModel.embed(segment).content(); embeddingStore.add(embedding, segment); notificationService.publish(userId, index.progress, (i 1) / segments.size()); } fileService.markReady(fileId, segments.size()); notificationService.publish(userId, index.finish, fileId.toString()); } catch (Exception e) { fileService.markFailed(fileId); notificationService.publish(userId, index.failed, e.getMessage()); } }演示代码为了进度准确对每个片段单独做 embedding。生产环境更推荐embeddingModel.embedAll(segments)批量处理再按批统计进度否则大量小片段会浪费请求时间。3.4 文档问答接口检索增强生成问答接口面向用户输入是问题输出是回答。为了把 RAG 串起来定义一个Assistant接口public interface CloudVaultAssistant { String chat(UserMessage String question); }在配置类中构建ChatLanguageModel、EmbeddingStoreRetriever和AiServicesChatLanguageModel chatModel OpenAiChatModel.builder() .baseUrl(chatBaseUrl) .apiKey(chatApiKey) .modelName(chatModelName) .build(); RetrieverTextSegment retriever EmbeddingStoreRetriever.from(embeddingStore, embeddingModel, 4); CloudVaultAssistant assistant AiServices.builder(CloudVaultAssistant.class) .chatLanguageModel(chatModel) .contentRetriever(retriever) .build();EmbeddingStoreRetriever的关键参数是4表示从向量库召回 4 个相关片段。这个参数不能拍脑袋需要根据文档切分大小和问题复杂度调整。问答接口RestController RequestMapping(/api/ask) public class AskController { private final CloudVaultAssistant assistant; public AskController(CloudVaultAssistant assistant) { this.assistant assistant; } PostMapping public String ask(RequestBody AskRequest request) { return assistant.chat(request.getQuestion()); } }到这里CloudVault 已经完成从上传文件到自动建索引再到问答的闭环。如果只是本地验证可以暂时不做多轮对话后续需要多轮时再按 userId 维护ChatMemory。4. 用 Redis 处理缓存、分布式锁与实时通知4.1 Redis 在 CloudVault 里负责什么Redis 不能替代 PostgreSQL 做持久化但在 CloudVault 里它有明确价值场景具体操作为什么用 Redis热点缓存文件列表、检索结果、登录态减少 PostgreSQL 压力分布式锁同一文件并发解析时只允许一个任务防止索引重复发布订阅多实例之间同步 SSE 通知单个节点内存 Emitter 无法跨实例通信计数与限流上传频率、问答频率避免资源被单个用户占满第一版可以只用StringRedisTemplate不需要引入太重的 Redis 客户端。4.2 文件索引的分布式锁文档解析是很典型的“不能让多个消费者同时执行”的任务。如果同一文件被重复上传或者重试任务和原始任务同时运行doc_chunks会出现重复向量问答结果也会受影响。用 Redis 做锁的示例public boolean tryLock(String key, String value, long timeoutSeconds) { return Boolean.TRUE.equals( redisTemplate.opsForValue().setIfAbsent(key, value, Duration.ofSeconds(timeoutSeconds)) ); } public void unlock(String key, String value) { String current redisTemplate.opsForValue().get(key); if (value.equals(current)) { redisTemplate.delete(key); } }在indexAsync方法开头获取锁String lockKey cloudvault:index:lock: fileId; String lockValue UUID.randomUUID().toString(); if (!tryLock(lockKey, lockValue, 600)) { notificationService.publish(userId, index.duplicate, fileId.toString()); return; } try { // 执行解析和向量化 } finally { unlock(lockKey, lockValue); }锁的 timeout 要明显大于一次解析的预期时间。如果 timeout 太小任务没跑完锁就过期了仍然会出现重复执行。生产环境建议把超时时间、重试次数做成配置项。4.3 基于 SSE 的实时通知实时通知不一定要用 WebSocketCloudVault 这类“服务端主动推送状态”的场景使用 SSE 更简单。SSE 是浏览器向服务端发一次 HTTP 请求服务端保持连接不断推送事件。通知服务示例Component public class NotificationService { private final MapString, CopyOnWriteArrayListSseEmitter userEmitters new ConcurrentHashMap(); private final StringRedisTemplate redisTemplate; public void subscribe(String userId, SseEmitter emitter) { userEmitters.computeIfAbsent(userId, k - new CopyOnWriteArrayList()) .add(emitter); } public void unsubscribe(String userId, SseEmitter emitter) { ListSseEmitter emitters userEmitters.get(userId); if (emitters ! null) { emitters.remove(emitter); } } public void publish(String userId, String event, String data) { String payload event: event \ndata: data \n\n; sendToLocal(userId, payload); redisTemplate.convertAndSend(cloudvault:notify: userId, payload); } private void sendToLocal(String userId, String payload) { ListSseEmitter emitters userEmitters.get(userId); if (emitters null) { return; } for (SseEmitter emitter : emitters) { try { emitter.send(payload); } catch (Exception e) { emitter.completeWithError(e); unsubscribe(userId, emitter); } } } }SSE 接入接口RestController RequestMapping(/api/sse) public class SseController { private final NotificationService notificationService; public SseController(NotificationService notificationService) { this.notificationService notificationService; } GetMapping(/{userId}) public SseEmitter connect(PathVariable String userId) {
返回列表