Java 程序员第 44 阶段07:大模型微服务拆分,独立服务解耦便于扩容维护,知识库管理服务:文档上传、切片、索引的独立服务化
- 知识库管理为何需要独立服务
- 独立知识库服务的核心职责
- 文档上传与存储设计
- 文本切片策略与最佳实践
- 向量化与索引构建流水线
- 版本控制与增量更新机制
- 多租户与权限隔离
- Spring Boot 实现详解
- 一致性、幂等性与异常处理
- 总结与展望
1.1 从“代码写死”到“知识运营”
在 RAG 应用的早期阶段,很多团队会把知识库相关的逻辑直接耦合在业务服务里。例如,上传一个 PDF 后就地解析、切片、调用 Embedding 模型,然后写入向量库。这种方式在 demo 阶段非常省事,但当知识库规模扩大、运营人员需要频繁更新文档、不同业务线需要独立知识库时,问题就会接踵而来:
- **上传阻塞主流程**:文档解析和向量化可能耗时数秒到数分钟,同步调用会拖慢业务接口。
- **更新不可控**:没有统一的文档版本管理,文档更新后旧索引残留,导致检索结果混杂。
- **格式处理复杂**:PDF、Word、Excel、Markdown、HTML、图片 OCR 等解析逻辑越来越重,污染业务代码。
- **多租户隔离困难**:不同租户的知识库数据如果没有独立管理,容易混用,带来安全和合规风险。
- **运营能力缺失**:没有可视化的知识库管理后台,运营人员无法自主维护内容。
因此,把知识库管理拆分为独立服务,让文档上传、切片、索引构建、版本控制、增量更新等能力平台化,是 RAG 架构演进的必然选择。
1.2 独立服务带来的收益
|
维度 |
耦合在业务服务 |
独立知识库管理服务 |
|
--- |
--- |
--- |
|
上传体验 |
同步阻塞,超时风险高 |
异步处理,上传即返回任务 ID |
|
更新控制 |
无版本管理,索引残留 |
版本化、原子切换、灰度发布 |
|
格式扩展 |
业务代码膨胀 |
独立演进,支持插件化解析 |
|
多租户 |
隔离逻辑散落各处 |
命名空间/租户维度统一隔离 |
|
运营能力 |
几乎无 |
可建设完整管理后台 |
|
稳定性 |
重索引影响业务 |
独立资源,不影响在线对话 |
2.1 功能边界
知识库管理服务应聚焦于“知识的运营”,其核心职责包括:
- **文档管理**:上传、下载、删除、查看、分类、标签、元数据维护。
- **文档解析**:把不同格式文件转换为原始文本,包括 PDF、Word、Excel、PPT、HTML、Markdown、图片 OCR 等。
- **文本切片**:按照固定长度、语义段落、递归、滑动窗口等策略把长文本切分为片段(Chunk)。
- **向量化**:调用 Embedding 模型把每个 Chunk 映射为向量。
- **索引构建**:将向量与元数据写入向量数据库,建立 ANN 索引。
- **版本管理**:对知识库版本进行快照管理,支持版本切换、回滚、灰度。
- **增量更新**:仅对已变更的文档或 Chunk 重新索引,避免全量重建。
- **任务调度**:对耗时操作进行异步任务调度,提供进度查询与失败重试。
2.2 在微服务架构中的位置
┌─────────────────────────────────────────────────────────┐
│ 业务服务层 │
│ 客服机器人 │ 内部知识助手 │ 智能报表 │ 编程助手 │
└──────────┬──────────────────────────────────┬───────────┘
│ │
▼ ▼
┌──────────────┐ ┌──────────────┐
│ RAG 检索服务 │◄──────────────────│ 知识库管理服务 │
└──────────────┘ └──────┬───────┘
│
┌───────────────────────────────┼──────────────┐
▼ ▼ ▼
┌──────────────┐ ┌──────────────┐ ┌──────────────┐
│ 对象存储 │ │ 向量数据库 │ │ Embedding │
│ MinIO/OSS │ │ Milvus/Qdrant│ │ 模型服务 │
└──────────────┘ └──────────────┘ └──────────────┘
知识库管理服务不直接面向终端用户回答问题,而是为业务服务和 RAG 检索服务提供“知识原料”和“索引能力”。
3.1 上传接口设计
文档上传应该支持同步预检与异步处理。典型的上传流程是:
- 客户端调用上传接口,服务端接收文件流并校验格式、大小、病毒等。
- 文件写入对象存储(MinIO、OSS、S3 等),同时生成唯一 document_id。
- 服务端创建索引任务,返回任务 ID 给客户端。
- 异步工作线程或任务队列消费任务,完成解析、切片、向量化、索引写入。
- 客户端可通过任务 ID 轮询处理进度。
@RestController
@RequestMapping("/api/v1/knowledge-base")
public class DocumentController {
private final DocumentUploadService uploadService;
@PostMapping("/{kbId}/documents")
public ResponseEntity<UploadTaskResponse> upload(
@PathVariable String kbId,
@RequestParam("file") MultipartFile file) {
UploadTask task = uploadService.submit(kbId, file);
return ResponseEntity.accepted().body(
UploadTaskResponse.builder()
.taskId(task.getId())
.status(task.getStatus())
.build());
}
@GetMapping("/tasks/{taskId}")
public ResponseEntity<UploadTask> getTask(@PathVariable String taskId) {
return ResponseEntity.ok(uploadService.getTask(taskId));
}
}
3.2 对象存储与元数据
上传的文件通常不直接存在本地磁盘,而是写入对象存储,这样便于共享、备份和扩展。元数据则存入关系型数据库(如 PostgreSQL)或文档数据库:
@Entity
@Table(name = "kb_document")
public class KnowledgeDocument {
@Id
private String id;
private String kbId;
private String filename;
private String storagePath;
private long size;
private String mimeType;
private String status; // PENDING, PARSING, INDEXING, READY, FAILED
private String version;
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
@Convert(converter = JsonMapConverter.class)
private Map<String, Object> metadata;
}
3.3 大文件与分片上传
对于大文件(如几百 MB 的 PDF 或视频转写文本),建议支持分片上传:
- 客户端先申请 upload session,获得 session_id 和分片大小。
- 分片上传后服务端暂存,所有分片上传完成后合并为完整文件。
- 合并成功后触发解析任务。
这种方式可以避免单次上传超时,也便于断点续传。
4.1 为什么切片质量决定 RAG 效果
RAG 检索的基本单元是 Chunk。Chunk 太大,会导致无关内容混入上下文,稀释关键信息;Chunk 太小,又会丢失上下文,导致语义不完整。切片策略直接影响检索精度和生成质量。
4.2 常见切片策略
|
策略 |
说明 |
适用场景 |
|
--- |
--- |
--- |
|
固定长度切片 |
按字符数或 token 数固定切分,可设置重叠区 |
通用场景,实现简单 |
|
段落切片 |
按自然段落或标题切分 |
结构清晰的文档 |
|
递归切片 |
先按大段落切,再对超长段落递归细分 |
兼顾语义与长度 |
|
语义切片 |
使用 Embedding 相似度判断语义边界 |
对效果要求高 |
|
文档结构切片 |
按章节、条款、表格、代码块等结构切分 |
法律、技术文档 |
4.3 固定长度切片示例
@Component
public class FixedLengthChunker implements DocumentChunker {
private final int chunkSize;
private final int overlap;
public FixedLengthChunker(@Value("${kb.chunk.size:512}") int chunkSize,
@Value("${kb.chunk.overlap:50}") int overlap) {
this.chunkSize = chunkSize;
this.overlap = overlap;
}
@Override
public List<Chunk> chunk(String documentId, String text) {
List<Chunk> chunks = new ArrayList<>();
int start = 0;
int idx = 0;
while (start < text.length()) {
int end = Math.min(start + chunkSize, text.length());
chunks.add(Chunk.builder()
.id(documentId + "_chunk_" + idx)
.documentId(documentId)
.index(idx)
.content(text.substring(start, end))
.build());
start += chunkSize - overlap;
idx++;
}
return chunks;
}
}
4.4 递归切片示例
@Component
public class RecursiveChunker implements DocumentChunker {
private final List<Integer> separators = List.of("\n\n", "\n", ".", " ");
@Override
public List<Chunk> chunk(String documentId, String text) {
return split(documentId, text, 0, 0);
}
private List<Chunk> split(String documentId, String text, int sepIndex, int chunkIndex) {
if (text.length() <= MAX_CHUNK_SIZE || sepIndex >= separators.size()) {
return List.of(Chunk.builder()
.id(documentId + "_chunk_" + chunkIndex)
.documentId(documentId)
.index(chunkIndex)
.content(text)
.build());
}
String sep = separators.get(sepIndex);
String[] parts = text.split(sep);
List<Chunk> result = new ArrayList<>();
int idx = chunkIndex;
for (String part : parts) {
if (part.length() > MAX_CHUNK_SIZE) {
result.addAll(split(documentId, part, sepIndex + 1, idx));
} else if (!part.isBlank()) {
result.add(Chunk.builder()
.id(documentId + "_chunk_" + idx)
.documentId(documentId)
.index(idx)
.content(part)
.build());
}
idx = result.size();
}
return result;
}
}
4.5 切片元数据设计
每个 Chunk 应保留足够的元数据,方便后续检索过滤和结果展示:
public class Chunk {
private String id;
private String documentId;
private String kbId;
private int index;
private String content;
private int startOffset;
private int endOffset;
private String title;
private String pageNumber;
private Map<String, Object> metadata;
}
5.1 流水线阶段

文档到索引的完整流水线通常包括:
- 下载文件:从对象存储下载原始文件。
- 格式解析:提取纯文本和结构信息。
- 文本清洗:去除多余空格、页眉页脚、乱码、无意义符号。
- 文本切片:生成 Chunk 列表。
- 批量向量化:调用 Embedding 模型获取向量。
- 元数据组装:把 Chunk 内容、来源、页码等写入 Payload。
- 向量写入:批量 upsert 到向量数据库。
- 状态更新:更新文档状态为 READY,记录版本信息。
5.2 异步任务调度
耗时步骤应通过任务队列异步执行。Spring 中可以使用 `@Async` 或集成 RabbitMQ/Kafka:
@Service
public class IndexingPipelineService {
private final DocumentParser parser;
private final DocumentChunker chunker;
private final EmbeddingService embeddingService;
private final VectorStore vectorStore;
@Async("indexingExecutor")
public CompletableFuture<Void> process(String taskId, String documentId) {
try {
taskService.updateStatus(taskId, IndexTaskStatus.PARSING);
Document document = documentService.getById(documentId);
String rawText = parser.parse(document.getStoragePath());
String cleanedText = cleanText(rawText);
taskService.updateStatus(taskId, IndexTaskStatus.CHUNKING);
List<Chunk> chunks = chunker.chunk(documentId, cleanedText);
taskService.updateStatus(taskId, IndexTaskStatus.EMBEDDING);
List<VectorRecord> records = embedChunks(document, chunks);
taskService.updateStatus(taskId, IndexTaskStatus.INDEXING);
vectorStore.upsert(document.getKbId(), records);
documentService.updateStatus(documentId, DocumentStatus.READY);
taskService.updateStatus(taskId, IndexTaskStatus.COMPLETED);
} catch (Exception e) {
taskService.updateStatus(taskId, IndexTaskStatus.FAILED, e.getMessage());
throw e;
}
return CompletableFuture.completedFuture(null);
}
private List<VectorRecord> embedChunks(Document doc, List<Chunk> chunks) {
List<String> texts = chunks.stream().map(Chunk::getContent).collect(Collectors.toList());
Map<String, List<Float>> embeddings = embeddingService.embedBatch(texts);
return chunks.stream().map(chunk -> VectorRecord.builder()
.id(chunk.getId())
.vector(embeddings.get(chunk.getContent()))
.content(chunk.getContent())
.documentId(doc.getId())
.chunkIdx(chunk.getIndex())
.metadata(Map.of(
"title", doc.getFilename(),
"page", chunk.getPageNumber(),
"kbId", doc.getKbId()))
.build()).collect(Collectors.toList());
}
}
5.3 批量与限流
Embedding 模型和向量库通常有并发限制,批量写入虽好,但过大会触发限流或 OOM。建议:
- 批量大小可配置,默认 32 或 64。
- 批量之间加入短暂休眠,避免瞬间打满模型服务。
- 使用 Sentinel 或 Guava RateLimiter 对 Embedding 调用限流。
- 向量库写入失败时,按指数退避重试。
6.1 知识库版本模型
生产环境中,知识库内容需要不断更新。如果每次更新都直接覆盖现有索引,可能出现检索结果不一致、更新失败难以回滚等问题。因此,引入版本控制机制非常关键。
public class KnowledgeBaseVersion {
private String versionId;
private String kbId;
private String status; // BUILDING, ACTIVE, ARCHIVED, FAILED
private List<String> documentIds;
private String vectorNamespace;
private LocalDateTime createdAt;
private LocalDateTime activatedAt;
}
每次发布新版本时:
- 新建一个版本记录,状态为 BUILDING。
- 在新命名空间(例如 `kb_v2`)中构建索引。
- 构建完成后切换状态为 ACTIVE,RAG 检索服务开始读取新命名空间。
- 旧版本进入 ARCHIVED 状态,保留一段时间后可清理。
6.2 增量更新策略
对于小幅度更新,不需要全量重建版本。可以按文档粒度进行增量更新:
- 文档新增:直接解析切片并写入当前命名空间。
- 文档修改:删除旧文档关联的所有 Chunk,重新写入新 Chunk。
- 文档删除:删除该文档关联的所有 Chunk。
public void updateDocument(String kbId, String documentId, MultipartFile newFile) {
// 1. 删除旧索引
List<String> oldChunkIds = chunkService.findIdsByDocumentId(documentId);
vectorStore.deleteByIds(kbId, oldChunkIds);
// 2. 重新上传与解析
UploadTask task = uploadService.submit(kbId, newFile);
documentService.bindNewVersion(documentId, task.getId());
}
6.3 版本切换与灰度
版本切换时,可以通过配置中心或数据库中的活跃版本字段控制:
@Service
public class KnowledgeBaseVersionService {
public String getActiveNamespace(String kbId) {
KnowledgeBaseVersion active = versionRepository
.findActiveByKbId(kbId)
.orElseThrow(() -> new KnowledgeBaseNotFoundException(kbId));
return active.getVectorNamespace();
}
@Transactional
public void activateVersion(String versionId) {
KnowledgeBaseVersion version = versionRepository.findById(versionId).orElseThrow();
versionRepository.deactivateAllByKbId(version.getKbId());
version.setStatus("ACTIVE");
version.setActivatedAt(LocalDateTime.now());
versionRepository.save(version);
}
}
RAG 检索服务每次查询前都读取当前活跃命名空间,这样就能实现版本切换的零停机。
7.1 命名空间隔离

多租户知识库最常用的是命名空间隔离。每个租户或每个知识库拥有独立的 Collection/Namespace:
tenant_a_kb_default
tenant_a_kb_manual
tenant_b_kb_default
这种方式隔离性最强,便于按租户做配额和计费。但管理成本较高,需要统一维护命名空间映射。
7.2 字段级隔离
在共享 Collection 的场景下,可以通过 `tenant_id` 字段过滤实现逻辑隔离:
public SearchQuery withTenantFilter(String tenantId, SearchQuery query) {
Map<String, Object> filter = new HashMap<>(query.getFilter());
filter.put("tenant_id", tenantId);
return query.toBuilder().filter(filter).build();
}
字段级隔离适合中小规模、希望降低运维成本的场景。需要注意的是,必须确保所有查询都强制带租户过滤,避免越权。
7.3 权限模型
知识库管理服务应支持细粒度权限:
- 读权限:可检索和查看文档。
- 写权限:可上传、修改、删除文档。
- 管理权限:可创建知识库、切换版本、配置切片策略。
- 审计权限:可查看操作日志和索引任务。
权限可以通过 RBAC 模型或 ABAC 模型实现,与统一认证服务(如 OAuth2、SSO)集成。
8.1 领域模型
public class KnowledgeBase {
private String id;
private String name;
private String tenantId;
private String activeVersionId;
private String chunkStrategy; // FIXED, RECURSIVE, SEMANTIC
private Map<String, Object> config;
private LocalDateTime createdAt;
}
public class Document {
private String id;
private String kbId;
private String filename;
private String storagePath;
private String status;
private String versionId;
private long size;
private Map<String, Object> metadata;
}
public class Chunk {
private String id;
private String documentId;
private String kbId;
private int index;
private String content;
private Map<String, Object> metadata;
}
8.2 解析器插件化
为了支持多种文档格式,可以设计一个解析器 SPI:
public interface DocumentParser {
boolean supports(String mimeType);
String parse(InputStream inputStream, String filename);
}
@Component
public class PdfDocumentParser implements DocumentParser {
@Override
public boolean supports(String mimeType) {
return "application/pdf".equals(mimeType);
}
@Override
public String parse(InputStream inputStream, String filename) {
try (PDDocument doc = PDDocument.load(inputStream)) {
PDFTextStripper stripper = new PDFTextStripper();
return stripper.getText(doc);
} catch (IOException e) {
throw new DocumentParseException("Failed to parse PDF: " + filename, e);
}
}
}
Spring 启动时自动收集所有 `DocumentParser` 实现,按 mimeType 路由:
@Service
public class DocumentParserRouter {
private final List<DocumentParser> parsers;
public DocumentParserRouter(List<DocumentParser> parsers) {
this.parsers = parsers;
}
public String parse(InputStream inputStream, String mimeType, String filename) {
return parsers.stream()
.filter(p -> p.supports(mimeType))
.findFirst()
.orElseThrow(() -> new UnsupportedDocumentTypeException(mimeType))
.parse(inputStream, filename);
}
}
8.3 任务表设计
@Entity
@Table(name = "kb_index_task")
public class IndexTask {
@Id
private String id;
private String documentId;
private String kbId;
@Enumerated(EnumType.STRING)
private IndexTaskStatus status;
private String errorMessage;
private int progress; // 0-100
private LocalDateTime createdAt;
private LocalDateTime completedAt;
}
任务状态流转:PENDING -> PARSING -> CHUNKING -> EMBEDDING -> INDEXING -> COMPLETED / FAILED。
8.4 管理接口示例
@RestController
@RequestMapping("/api/v1/knowledge-base")
public class KnowledgeBaseAdminController {
private final KnowledgeBaseService kbService;
@PostMapping
public ResponseEntity<KnowledgeBase> create(@RequestBody @Valid CreateKbRequest request) {
return ResponseEntity.status(HttpStatus.CREATED).body(kbService.create(request));
}
@PostMapping("/{kbId}/versions")
public ResponseEntity<KnowledgeBaseVersion> createVersion(@PathVariable String kbId) {
return ResponseEntity.status(HttpStatus.CREATED).body(kbService.createVersion(kbId));
}
@PostMapping("/{kbId}/versions/{versionId}/activate")
public ResponseEntity<Void> activateVersion(@PathVariable String kbId,
@PathVariable String versionId) {
kbService.activateVersion(versionId);
return ResponseEntity.ok().build();
}
@GetMapping("/{kbId}/documents")
public ResponseEntity<Page<Document>> listDocuments(@PathVariable String kbId,
Pageable pageable) {
return ResponseEntity.ok(kbService.listDocuments(kbId, pageable));
}
}
9.1 幂等上传
同一文件重复上传时,应避免重复解析和索引。可以通过文件内容哈希或业务侧唯一标识去重:
public String computeFileHash(MultipartFile file) throws IOException {
return DigestUtils.sha256Hex(file.getInputStream());
}
public Optional<Document> findByHash(String kbId, String hash) {
return documentRepository.findByKbIdAndContentHash(kbId, hash);
}
9.2 任务幂等执行
索引任务执行前应检查状态,避免同一任务被多个 worker 重复处理:
@Transactional
public boolean acquireTask(String taskId) {
int updated = taskRepository.updateStatus(
taskId, IndexTaskStatus.PENDING, IndexTaskStatus.PARSING);
return updated == 1;
}
9.3 失败重试与死信
任务失败后可按固定间隔或指数退避重试。超过最大重试次数后进入死信队列或人工处理队列:
@Retryable(retryFor = {EmbeddingUnavailableException.class},
backoff = @Backoff(delay = 1000, multiplier = 2, maxDelay = 30000))
public List<VectorRecord> embedWithRetry(List<String> texts) {
return embeddingService.embedBatch(texts);
}
9.4 最终一致性
文档上传和索引构建是异步过程,因此业务侧需要接受“最终一致性”。可以通过任务状态接口、Webhook 回调或事件总线通知业务方。
知识库管理服务是大模型 RAG 体系中不可或缺的“内容中台”。把它从业务服务中独立出来后,我们能够:
- 让文档运营流程标准化、可视化、可审计。
- 通过版本管理实现知识库的安全更新与灰度切换。
- 通过增量更新降低索引维护成本。
- 通过插件化解码器支持多格式文档的长期演进。
- 通过多租户隔离满足不同客户或业务线的数据安全需求。
Java 程序员在实现这一服务时,应重点关注:异步任务调度、幂等性、版本控制、多租户隔离和可观测性。与 RAG 检索服务类似,知识库管理服务也是无状态的,可以独立水平扩展。
下一篇文章,我们将讨论**对话会话服务的独立化**,重点解决多轮上下文状态管理、会话持久化、上下文压缩与 Token 控制等问题。敬请期待。
更多推荐


所有评论(0)