面试准备:KnowFlow项目深度解析

KnowFlow项目深度解析:企业级RAG知识库系统面试全攻略

目录


一、项目概述

KnowFlow 是一个企业级RAG(Retrieval-Augmented Generation)知识库系统,旨在帮助企业构建智能化的文档检索和知识问答平台。

1.1 项目基本信息

  • 项目名称: KnowFlow
  • 项目类型: 企业级RAG知识库系统
  • 开发周期: X个月
  • 团队规模: X人
  • 我的角色: 后端核心开发/架构设计

1.2 核心技术栈

1
2
3
4
5
6
7
8
9
后端框架:SpringBoot 3.x
数据存储:MySQL 8.0、Redis 7.0
搜索引擎:Elasticsearch 8.x
向量数据库:Milvus 2.3
对象存储:MinIO
监控系统:Prometheus + Grafana
消息队列:RabbitMQ
文档解析:Apache Tika、PDFBox
AI模型:OpenAI API / 本地LLM

1.3 系统特点

  • 混合检索: 结合全文检索和向量检索,提升检索准确率
  • 高性能: 支持百万级文档存储,毫秒级检索响应
  • 高可用: 分布式架构,服务可用性99.9%
  • 安全可控: 完善的权限管理和数据安全机制
  • 可观测性: 全链路监控和性能分析

二、项目背景和业务价值

2.1 为什么做这个项目?

业务痛点:

  1. 信息孤岛严重: 企业文档散落在各个系统,查找困难
  2. 检索效率低: 传统关键词检索无法理解语义,召回率低
  3. 知识利用率低: 大量历史文档和经验无法被有效利用
  4. 人工成本高: 员工需要花费大量时间查找和整理信息

市场需求:

  • ChatGPT等大模型兴起,企业希望构建私有知识库
  • RAG技术成为主流方案,避免模型幻觉问题
  • 企业对数据安全和隐私保护要求高

2.2 解决了什么问题?

  1. 智能检索: 支持自然语言查询,理解用户意图
  2. 准确召回: 混合检索策略,召回率提升40%
  3. 知识沉淀: 统一的知识管理平台
  4. 降本增效: 减少70%的信息查找时间

2.3 业务价值

  • 提升效率: 员工查找信息时间从平均30分钟降至5分钟以内
  • 降低成本: 减少重复性咨询,客服工作量降低50%
  • 知识复用: 历史经验被有效利用,新员工培训周期缩短60%
  • 决策支持: 快速获取相关信息,提升决策质量

三、技术架构详解

3.1 整体架构图

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
┌─────────────────────────────────────────────────────────────┐
│ 客户端层 │
│ Web / Mobile / API │
└──────────────────────┬──────────────────────────────────────┘

┌──────────────────────▼──────────────────────────────────────┐
│ 网关层(Gateway) │
│ Nginx + Spring Cloud Gateway │
│ 认证、鉴权、限流、日志、负载均衡 │
└──────────────────────┬──────────────────────────────────────┘

┌──────────────────────▼──────────────────────────────────────┐
│ 应用服务层 │
├─────────────────────────────────────────────────────────────┤
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ 文档服务 │ │ 检索服务 │ │ 用户服务 │ │
│ │ DocumentSvc │ │ SearchSvc │ │ UserSvc │ │
│ └──────────────┘ └──────────────┘ └──────────────┘ │
│ │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ 解析服务 │ │ 向量服务 │ │ 问答服务 │ │
│ │ ParseSvc │ │ VectorSvc │ │ QASvc │ │
│ └──────────────┘ └──────────────┘ └──────────────┘ │
└──────────────────────┬──────────────────────────────────────┘

┌──────────────────────▼──────────────────────────────────────┐
│ 数据存储层 │
├─────────────────────────────────────────────────────────────┤
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ MySQL │ │ Redis │ │Elasticsearch││ Milvus │ │
│ │结构化数据 │ │ 缓存 │ │ 全文检索 │ │向量检索 │ │
│ └──────────┘ └──────────┘ └──────────┘ └──────────┘ │
│ │
│ ┌──────────┐ ┌──────────┐ │
│ │ MinIO │ │ RabbitMQ │ │
│ │对象存储 │ │ 消息队列 │ │
│ └──────────┘ └──────────┘ │
└──────────────────────┬──────────────────────────────────────┘

┌──────────────────────▼──────────────────────────────────────┐
│ 基础设施层 │
│ Prometheus + Grafana + ELK + Jaeger │
│ 监控、日志、链路追踪 │
└─────────────────────────────────────────────────────────────┘

3.2 核心处理流程

文档上传处理流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
用户上传文档


网关层验证(大小、格式、权限)


文档服务接收

├─► 保存原文件到MinIO

├─► 发送MQ消息到解析队列

└─► 返回任务ID给用户

▼(异步处理)
解析服务消费消息

├─► Apache Tika识别文档类型

├─► 提取文本内容

├─► 文档分块(Chunking)

└─► 发送MQ消息到向量化队列


向量服务处理

├─► 调用Embedding模型生成向量

├─► 存储向量到Milvus

├─► 存储文本到Elasticsearch

└─► 更新MySQL元数据状态

检索查询流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
用户输入查询


检索服务接收请求

├─► 查询改写(Query Rewrite

├─► 生成查询向量(Embedding


并行检索

├─► Milvus向量检索(Top 50
│ │
│ └─► 基于余弦相似度

└─► Elasticsearch全文检索(Top 50

└─► BM25算法 + 语义匹配


混合排序(Hybrid Ranking

├─► RRF算法融合结果

├─► 重排序(Reranking

└─► 返回Top 10


结果增强

├─► 高亮关键词

├─► 摘要生成

└─► 返回给用户

3.3 数据流设计

文档数据流:

1
2
3
4
5
原始文档 → MinIO(原文件)
→ Elasticsearch(分块文本 + 元数据)
→ Milvus(向量数据)
→ MySQL(元数据 + 状态)
→ Redis(热点缓存)

四、技术选型深度分析

4.1 SpringBoot 作为后端框架

选择理由:

  • 生态成熟,社区活跃
  • 开箱即用,快速开发
  • 与各种中间件集成方便
  • 团队技术栈统一

对比其他方案:

  • vs Go: Java生态更完善,团队更熟悉,开发效率更高
  • vs Python Flask: Java性能更好,类型安全,适合企业级应用
  • vs Node.js: CPU密集型任务Java更适合,多线程处理更成熟

4.2 Elasticsearch 做全文检索

选择理由:

  1. 倒排索引: 毫秒级全文检索
  2. 分布式架构: 支持水平扩展
  3. 丰富的分析器: 支持中文分词(IK Analyzer)
  4. BM25算法: 经典的相关性排序算法
  5. 成熟稳定: 大规模生产验证

对比其他方案:

  • vs Solr: ES更易用,JSON友好,社区更活跃
  • vs MySQL全文索引: ES性能更好,支持分布式,功能更强
  • vs OpenSearch: ES商业支持更好,插件更丰富

配置优化:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
{
"settings": {
"number_of_shards": 5,
"number_of_replicas": 1,
"refresh_interval": "5s",
"analysis": {
"analyzer": {
"ik_smart_analyzer": {
"type": "custom",
"tokenizer": "ik_smart",
"filter": ["lowercase", "stop"]
}
}
}
},
"mappings": {
"properties": {
"content": {
"type": "text",
"analyzer": "ik_smart_analyzer",
"search_analyzer": "ik_smart_analyzer"
},
"title": {
"type": "text",
"analyzer": "ik_smart_analyzer",
"boost": 2.0
},
"doc_id": {
"type": "keyword"
},
"chunk_index": {
"type": "integer"
},
"created_time": {
"type": "date"
}
}
}
}

4.3 Milvus 向量数据库

选择理由:

  1. 专为向量检索设计: 性能极致优化
  2. 多种索引类型: HNSW、IVF_FLAT、IVF_PQ等
  3. 高性能: 百万级向量毫秒级检索
  4. 云原生: 支持Kubernetes部署
  5. 活跃社区: 国产开源,社区活跃

对比其他方案:

  • vs Faiss: Milvus是分布式的,Faiss是单机库
  • vs Pinecone: Milvus开源免费,可私有化部署
  • vs Weaviate: Milvus性能更好,适合大规模场景
  • vs Qdrant: Milvus生态更完善,中文文档更好

索引选择:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
# HNSW索引 - 高召回率,适合精确检索
index_params = {
"metric_type": "IP", # 内积
"index_type": "HNSW",
"params": {
"M": 16, # 每个节点的连接数
"efConstruction": 200 # 构建时的候选数
}
}

# 检索参数
search_params = {
"metric_type": "IP",
"params": {
"ef": 100 # 检索时的候选数
}
}

4.4 MinIO 对象存储

选择理由:

  • 兼容S3 API,迁移方便
  • 开源免费,成本低
  • 高性能,支持分布式
  • 私有化部署,数据安全

对比方案:

  • vs 阿里云OSS: MinIO私有化,成本更低
  • vs FastDFS: MinIO API更标准,生态更好
  • vs NFS: MinIO性能更好,支持对象存储特性

4.5 Prometheus + Grafana 监控

选择理由:

  • Prometheus时序数据库,适合指标监控
  • Grafana可视化强大
  • 开源生态完善
  • 支持告警规则

监控指标:

  • 系统指标:CPU、内存、磁盘、网络
  • 应用指标:QPS、响应时间、错误率
  • 业务指标:文档数量、检索量、向量化任务数
  • 中间件指标:ES查询延迟、Milvus检索耗时、Redis命中率

五、核心功能实现详解

5.1 文档上传与解析

5.1.1 文档上传接口

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
@RestController
@RequestMapping("/api/v1/documents")
@Slf4j
public class DocumentController {

@Autowired
private DocumentService documentService;

@Autowired
private MinioService minioService;

@PostMapping("/upload")
public Result<DocumentUploadResponse> uploadDocument(
@RequestParam("file") MultipartFile file,
@RequestParam(value = "category", required = false) String category,
@RequestHeader("Authorization") String token) {

// 1. 参数校验
validateFile(file);

// 2. 提取用户信息
UserInfo userInfo = jwtService.parseToken(token);

// 3. 生成文档ID
String docId = IdUtil.generateDocId();

// 4. 上传到MinIO
String minioPath = minioService.uploadFile(
file.getInputStream(),
docId,
file.getOriginalFilename()
);

// 5. 保存元数据到MySQL
DocumentEntity doc = DocumentEntity.builder()
.docId(docId)
.fileName(file.getOriginalFilename())
.fileSize(file.getSize())
.mimeType(file.getContentType())
.minioPath(minioPath)
.userId(userInfo.getUserId())
.category(category)
.status(DocStatus.PENDING)
.uploadTime(LocalDateTime.now())
.build();
documentService.save(doc);

// 6. 发送解析任务到MQ
ParseTaskMessage message = ParseTaskMessage.builder()
.docId(docId)
.minioPath(minioPath)
.mimeType(file.getContentType())
.build();
rabbitTemplate.convertAndSend(
QueueConstants.PARSE_QUEUE,
message
);

log.info("Document uploaded successfully, docId: {}", docId);

return Result.success(DocumentUploadResponse.builder()
.docId(docId)
.status("processing")
.build());
}

private void validateFile(MultipartFile file) {
// 文件大小限制:100MB
if (file.getSize() > 100 * 1024 * 1024) {
throw new BusinessException("文件大小超过限制");
}

// 支持的文件类型
String contentType = file.getContentType();
List<String> allowedTypes = Arrays.asList(
"application/pdf",
"application/msword",
"application/vnd.openxmlformats-officedocument.wordprocessingml.document",
"text/plain",
"text/markdown"
);

if (!allowedTypes.contains(contentType)) {
throw new BusinessException("不支持的文件类型");
}
}
}

5.1.2 文档解析服务

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
@Service
@Slf4j
public class DocumentParseService {

@Autowired
private MinioService minioService;

@Autowired
private DocumentService documentService;

@Autowired
private RabbitTemplate rabbitTemplate;

@RabbitListener(queues = QueueConstants.PARSE_QUEUE)
public void handleParseTask(ParseTaskMessage message) {
String docId = message.getDocId();

try {
// 1. 从MinIO下载文件
InputStream inputStream = minioService.downloadFile(message.getMinioPath());

// 2. 使用Tika解析文档
Tika tika = new Tika();
String content = tika.parseToString(inputStream);

// 3. 文本清洗
content = cleanText(content);

// 4. 文档分块
List<DocumentChunk> chunks = chunkDocument(content, docId);

// 5. 保存分块信息到MySQL
documentService.saveChunks(chunks);

// 6. 更新文档状态
documentService.updateStatus(docId, DocStatus.PARSED);

// 7. 发送向量化任务到MQ
for (DocumentChunk chunk : chunks) {
VectorizeTaskMessage vectorMsg = VectorizeTaskMessage.builder()
.docId(docId)
.chunkId(chunk.getChunkId())
.content(chunk.getContent())
.build();
rabbitTemplate.convertAndSend(
QueueConstants.VECTORIZE_QUEUE,
vectorMsg
);
}

log.info("Document parsed successfully, docId: {}, chunks: {}",
docId, chunks.size());

} catch (Exception e) {
log.error("Parse document failed, docId: {}", docId, e);
documentService.updateStatus(docId, DocStatus.PARSE_FAILED);
}
}

/**
* 文本清洗
*/
private String cleanText(String text) {
// 移除多余空白字符
text = text.replaceAll("\\s+", " ");

// 移除特殊字符
text = text.replaceAll("[\\x00-\\x08\\x0B-\\x0C\\x0E-\\x1F]", "");

// 统一换行符
text = text.replaceAll("\\r\\n", "\\n");

return text.trim();
}

/**
* 文档分块策略
*/
private List<DocumentChunk> chunkDocument(String content, String docId) {
List<DocumentChunk> chunks = new ArrayList<>();

// 分块参数
int chunkSize = 500; // 每块500字符
int chunkOverlap = 100; // 重叠100字符

int start = 0;
int chunkIndex = 0;

while (start < content.length()) {
int end = Math.min(start + chunkSize, content.length());

// 尽量在句子边界分割
if (end < content.length()) {
int sentenceEnd = findSentenceEnd(content, end);
if (sentenceEnd > start && sentenceEnd - start <= chunkSize + 100) {
end = sentenceEnd;
}
}

String chunkContent = content.substring(start, end);

DocumentChunk chunk = DocumentChunk.builder()
.chunkId(IdUtil.generateChunkId())
.docId(docId)
.content(chunkContent)
.chunkIndex(chunkIndex++)
.startPos(start)
.endPos(end)
.createTime(LocalDateTime.now())
.build();

chunks.add(chunk);

// 下一块的起始位置(考虑重叠)
start = end - chunkOverlap;
}

return chunks;
}

private int findSentenceEnd(String text, int pos) {
// 查找最近的句子结束符
String[] endMarkers = {"。", "!", "?", ".", "!", "?"};
int nearestEnd = -1;

for (String marker : endMarkers) {
int idx = text.indexOf(marker, pos - 50);
if (idx >= pos - 50 && idx <= pos + 50) {
if (nearestEnd == -1 || Math.abs(idx - pos) < Math.abs(nearestEnd - pos)) {
nearestEnd = idx + marker.length();
}
}
}

return nearestEnd > 0 ? nearestEnd : pos;
}
}

5.2 向量化处理

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
@Service
@Slf4j
public class VectorizeService {

@Autowired
private OpenAIClient openAIClient;

@Autowired
private MilvusService milvusService;

@Autowired
private ElasticsearchService esService;

@Autowired
private DocumentService documentService;

@RabbitListener(queues = QueueConstants.VECTORIZE_QUEUE, concurrency = "5")
public void handleVectorizeTask(VectorizeTaskMessage message) {
String chunkId = message.getChunkId();

try {
// 1. 调用Embedding API生成向量
List<Float> vector = generateEmbedding(message.getContent());

// 2. 存储向量到Milvus
milvusService.insertVector(
chunkId,
vector,
message.getDocId()
);

// 3. 存储文本到Elasticsearch
esService.indexDocument(
chunkId,
message.getDocId(),
message.getContent(),
LocalDateTime.now()
);

// 4. 更新分块状态
documentService.updateChunkStatus(chunkId, ChunkStatus.COMPLETED);

log.info("Vectorize completed, chunkId: {}", chunkId);

} catch (Exception e) {
log.error("Vectorize failed, chunkId: {}", chunkId, e);
documentService.updateChunkStatus(chunkId, ChunkStatus.FAILED);

// 失败重试逻辑
retryVectorize(message);
}
}

/**
* 生成文本向量
*/
private List<Float> generateEmbedding(String text) {
// 调用OpenAI Embedding API
EmbeddingRequest request = EmbeddingRequest.builder()
.model("text-embedding-3-small")
.input(text)
.build();

EmbeddingResponse response = openAIClient.createEmbedding(request);

return response.getData().get(0).getEmbedding();
}

/**
* 批量向量化(性能优化)
*/
public void batchVectorize(List<VectorizeTaskMessage> messages) {
// 批量调用Embedding API
List<String> texts = messages.stream()
.map(VectorizeTaskMessage::getContent)
.collect(Collectors.toList());

List<List<Float>> vectors = batchGenerateEmbeddings(texts);

// 批量插入Milvus
List<MilvusInsertData> insertDataList = new ArrayList<>();
for (int i = 0; i < messages.size(); i++) {
insertDataList.add(MilvusInsertData.builder()
.id(messages.get(i).getChunkId())
.vector(vectors.get(i))
.docId(messages.get(i).getDocId())
.build());
}

milvusService.batchInsert(insertDataList);
}
}

5.3 混合检索实现

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
@Service
@Slf4j
public class HybridSearchService {

@Autowired
private MilvusService milvusService;

@Autowired
private ElasticsearchService esService;

@Autowired
private OpenAIClient openAIClient;

@Autowired
private RedisService redisService;

public SearchResult hybridSearch(SearchRequest request) {
String query = request.getQuery();
int topK = request.getTopK() != null ? request.getTopK() : 10;

// 1. 查询改写(可选)
String rewrittenQuery = rewriteQuery(query);

// 2. 生成查询向量
List<Float> queryVector = generateEmbedding(rewrittenQuery);

// 3. 并行检索
CompletableFuture<List<SearchResult>> vectorSearchFuture =
CompletableFuture.supplyAsync(() ->
vectorSearch(queryVector, topK * 5));

CompletableFuture<List<SearchResult>> textSearchFuture =
CompletableFuture.supplyAsync(() ->
textSearch(rewrittenQuery, topK * 5));

// 4. 等待结果
List<SearchResult> vectorResults = vectorSearchFuture.join();
List<SearchResult> textResults = textSearchFuture.join();

// 5. 混合排序(RRF算法)
List<SearchResult> mergedResults = reciprocalRankFusion(
vectorResults,
textResults,
topK
);

// 6. 重排序(可选)
if (request.isEnableRerank()) {
mergedResults = rerank(query, mergedResults);
}

// 7. 结果增强
enhanceResults(mergedResults, query);

return SearchResult.builder()
.query(query)
.results(mergedResults)
.totalCount(mergedResults.size())
.searchTime(System.currentTimeMillis() - startTime)
.build();
}

/**
* 向量检索
*/
private List<SearchResult> vectorSearch(List<Float> queryVector, int topK) {
SearchParam searchParam = SearchParam.builder()
.collectionName("document_vectors")
.vectors(Collections.singletonList(queryVector))
.topK(topK)
.metricType(MetricType.IP)
.params("{\"ef\": 100}")
.build();

SearchResults results = milvusService.search(searchParam);

return results.getResults().stream()
.map(r -> SearchResult.builder()
.chunkId(r.getId())
.score(r.getScore())
.source("vector")
.build())
.collect(Collectors.toList());
}

/**
* 全文检索
*/
private List<SearchResult> textSearch(String query, int topK) {
SearchRequest esRequest = new SearchRequest("document_chunks");
SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();

// 多字段查询
MultiMatchQueryBuilder multiMatchQuery = QueryBuilders
.multiMatchQuery(query, "content", "title")
.type(MultiMatchQueryBuilder.Type.BEST_FIELDS)
.fuzziness(Fuzziness.AUTO);

sourceBuilder.query(multiMatchQuery);
sourceBuilder.size(topK);
sourceBuilder.highlighter(
new HighlightBuilder()
.field("content")
.preTags("<em>")
.postTags("</em>")
);

esRequest.source(sourceBuilder);

SearchResponse response = esService.search(esRequest);

return Arrays.stream(response.getHits().getHits())
.map(hit -> SearchResult.builder()
.chunkId(hit.getId())
.score(hit.getScore())
.content(hit.getSourceAsMap().get("content").toString())
.highlights(hit.getHighlightFields())
.source("text")
.build())
.collect(Collectors.toList());
}

/**
* RRF混合排序算法
* score(d) = Σ 1 / (k + rank_i(d))
*/
private List<SearchResult> reciprocalRankFusion(
List<SearchResult> vectorResults,
List<SearchResult> textResults,
int topK) {

Map<String, Double> scores = new HashMap<>();
int k = 60; // RRF参数

// 计算向量检索的RRF分数
for (int i = 0; i < vectorResults.size(); i++) {
String chunkId = vectorResults.get(i).getChunkId();
double score = 1.0 / (k + i + 1);
scores.merge(chunkId, score, Double::sum);
}

// 计算文本检索的RRF分数
for (int i = 0; i < textResults.size(); i++) {
String chunkId = textResults.get(i).getChunkId();
double score = 1.0 / (k + i + 1);
scores.merge(chunkId, score, Double::sum);
}

// 按分数排序
return scores.entrySet().stream()
.sorted(Map.Entry.<String, Double>comparingByValue().reversed())
.limit(topK)
.map(entry -> SearchResult.builder()
.chunkId(entry.getKey())
.score(entry.getValue())
.build())
.collect(Collectors.toList());
}

/**
* 查询改写
*/
private String rewriteQuery(String query) {
// 1. 同义词扩展
// 2. 拼写纠正
// 3. 查询扩展
return query; // 简化实现
}

/**
* 结果重排序
*/
private List<SearchResult> rerank(String query, List<SearchResult> results) {
// 使用Cross-Encoder模型重排序
// 或使用规则based重排序
return results;
}
}

5.4 RAG问答实现

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
@Service
@Slf4j
public class QAService {

@Autowired
private HybridSearchService searchService;

@Autowired
private OpenAIClient openAIClient;

public QAResponse answerQuestion(String question) {
// 1. 检索相关文档
SearchRequest searchRequest = SearchRequest.builder()
.query(question)
.topK(5)
.enableRerank(true)
.build();

SearchResult searchResult = searchService.hybridSearch(searchRequest);

// 2. 构建上下文
String context = buildContext(searchResult.getResults());

// 3. 构建Prompt
String prompt = buildPrompt(question, context);

// 4. 调用LLM生成答案
ChatCompletionRequest chatRequest = ChatCompletionRequest.builder()
.model("gpt-4")
.messages(Arrays.asList(
ChatMessage.builder()
.role("system")
.content("你是一个专业的知识库问答助手,基于给定的上下文回答问题。")
.build(),
ChatMessage.builder()
.role("user")
.content(prompt)
.build()
))
.temperature(0.7)
.maxTokens(500)
.build();

ChatCompletionResponse chatResponse = openAIClient.createChatCompletion(chatRequest);
String answer = chatResponse.getChoices().get(0).getMessage().getContent();

// 5. 返回结果
return QAResponse.builder()
.question(question)
.answer(answer)
.sources(searchResult.getResults())
.confidence(calculateConfidence(searchResult))
.build();
}

private String buildContext(List<SearchResult> results) {
StringBuilder context = new StringBuilder();
for (int i = 0; i < results.size(); i++) {
context.append("文档").append(i + 1).append(": ")
.append(results.get(i).getContent())
.append("\n\n");
}
return context.toString();
}

private String buildPrompt(String question, String context) {
return String.format(
"基于以下上下文回答问题。如果上下文中没有相关信息,请明确说明。\n\n" +
"上下文:\n%s\n\n" +
"问题:%s\n\n" +
"答案:",
context, question
);
}
}

六、技术难点与解决方案

难点1:大文件上传处理

问题描述:
用户上传超大文件(如500MB的PDF),普通上传方式会导致:

  • 内存溢出
  • 上传超时
  • 网络中断后需要重新上传

解决方案:

  1. 分片上传:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
@PostMapping("/upload/multipart")
public Result<UploadResponse> multipartUpload(
@RequestParam("file") MultipartFile filePart,
@RequestParam("chunkIndex") int chunkIndex,
@RequestParam("totalChunks") int totalChunks,
@RequestParam("fileId") String fileId,
@RequestParam("md5") String md5) {

// 1. 校验MD5
String calculatedMd5 = DigestUtils.md5Hex(filePart.getInputStream());
if (!calculatedMd5.equals(md5)) {
throw new BusinessException("文件分片校验失败");
}

// 2. 保存分片到临时目录
String tempDir = "/tmp/uploads/" + fileId;
File chunkFile = new File(tempDir, "chunk_" + chunkIndex);
filePart.transferTo(chunkFile);

// 3. 检查是否所有分片都已上传
if (isAllChunksUploaded(fileId, totalChunks)) {
// 4. 合并分片
String finalPath = mergeChunks(fileId, totalChunks);

// 5. 上传到MinIO
minioService.uploadFile(new FileInputStream(finalPath), fileId);

// 6. 清理临时文件
FileUtils.deleteDirectory(new File(tempDir));

return Result.success("上传完成");
}

return Result.success("分片上传成功");
}
  1. 断点续传:
  • 记录已上传的分片信息到Redis
  • 上传前查询已完成的分片
  • 只上传缺失的分片
  1. 流式处理:
1
2
3
4
5
6
7
8
9
10
11
// 使用流式读取,避免一次性加载到内存
try (InputStream is = file.getInputStream();
BufferedInputStream bis = new BufferedInputStream(is)) {
minioClient.putObject(
PutObjectArgs.builder()
.bucket(bucketName)
.object(objectName)
.stream(bis, file.getSize(), -1)
.build()
);
}

效果:

  • 支持1GB以上文件上传
  • 断网后可续传
  • 内存占用稳定在50MB以内

难点2:向量检索性能优化

问题描述:

  • 百万级向量检索耗时超过1秒
  • 高并发场景下Milvus压力大
  • 检索准确率和性能难以平衡

解决方案:

  1. 索引优化:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
# 使用HNSW索引替代IVF_FLAT
# HNSW: 图结构索引,检索速度快
index_params = {
"index_type": "HNSW",
"metric_type": "IP",
"params": {
"M": 16, # 连接数,越大召回率越高
"efConstruction": 200 # 构建参数,越大质量越高
}
}

# 检索参数调优
search_params = {
"metric_type": "IP",
"params": {
"ef": 64 # 检索参数,越大召回率越高但速度慢
}
}
  1. 分区策略:
1
2
3
4
5
6
// 按文档类别分区
collection.createPartition("tech_docs");
collection.createPartition("business_docs");

// 检索时指定分区,减少检索范围
searchParam.setPartitionNames(Arrays.asList("tech_docs"));
  1. 缓存热点查询:
1
2
3
4
@Cacheable(value = "search:vector", key = "#queryVector.hashCode()")
public List<SearchResult> cachedVectorSearch(List<Float> queryVector) {
return milvusService.search(queryVector);
}
  1. 批量查询优化:
1
2
3
// 批量查询,减少网络开销
List<List<Float>> queryVectors = Arrays.asList(vector1, vector2, vector3);
SearchResults results = milvusService.batchSearch(queryVectors);

性能提升:

  • 检索耗时从1200ms降至150ms
  • 支持QPS从50提升到500
  • 召回率保持在95%以上

难点3:文档分块策略

问题描述:

  • 分块太大:超过模型token限制,语义不聚焦
  • 分块太小:上下文不完整,检索结果碎片化
  • 固定大小分块:可能截断重要信息

解决方案:

  1. 语义分块:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
public List<DocumentChunk> semanticChunking(String content) {
List<DocumentChunk> chunks = new ArrayList<>();

// 1. 先按段落分割
String[] paragraphs = content.split("\\n\\n+");

StringBuilder currentChunk = new StringBuilder();
int currentSize = 0;
int maxSize = 500;
int minSize = 200;

for (String para : paragraphs) {
int paraSize = para.length();

// 如果当前块加上这个段落超过最大值
if (currentSize + paraSize > maxSize && currentSize >= minSize) {
// 保存当前块
chunks.add(createChunk(currentChunk.toString()));
currentChunk = new StringBuilder();
currentSize = 0;
}

currentChunk.append(para).append("\n\n");
currentSize += paraSize;
}

// 保存最后一块
if (currentSize > 0) {
chunks.add(createChunk(currentChunk.toString()));
}

return chunks;
}
  1. 滑动窗口分块:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
public List<DocumentChunk> slidingWindowChunking(String content) {
int chunkSize = 500;
int overlap = 100; // 重叠区域

List<DocumentChunk> chunks = new ArrayList<>();
int start = 0;

while (start < content.length()) {
int end = Math.min(start + chunkSize, content.length());
String chunk = content.substring(start, end);

chunks.add(createChunk(chunk));

start += (chunkSize - overlap);
}

return chunks;
}
  1. 递归分块:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
public List<DocumentChunk> recursiveChunking(String content, int maxSize) {
// 尝试按段落分割
String[] parts = content.split("\\n\\n");
if (parts.length > 1) {
// 递归处理每个段落
return Arrays.stream(parts)
.flatMap(p -> recursiveChunking(p, maxSize).stream())
.collect(Collectors.toList());
}

// 尝试按句子分割
parts = content.split("[。!?.!?]");
if (parts.length > 1 && content.length() > maxSize) {
return Arrays.stream(parts)
.flatMap(p -> recursiveChunking(p, maxSize).stream())
.collect(Collectors.toList());
}

// 无法继续分割,返回当前内容
return Collections.singletonList(createChunk(content));
}
  1. 智能分块(基于NLP):
1
2
3
4
5
6
// 使用句子边界检测
SentenceDetectorME sentenceDetector = new SentenceDetectorME(model);
String[] sentences = sentenceDetector.sentDetect(content);

// 按语义相似度合并句子
List<DocumentChunk> chunks = mergeBySemanticSimilarity(sentences);

最终方案:

  • 主要使用滑动窗口(chunk_size=500, overlap=100)
  • 对于结构化文档(如技术文档),使用语义分块
  • 对于长篇内容,先递归分块再滑动窗口

难点4:混合检索排序算法

问题描述:

  • 向量检索和全文检索的分数不在同一量纲
  • 简单加权融合效果不理想
  • 不同查询类型需要不同的融合策略

解决方案:

  1. RRF(Reciprocal Rank Fusion)算法:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
/**
* RRF算法实现
* 优点:不依赖具体分数,只依赖排名
* 公式:score(d) = Σ 1 / (k + rank_i(d))
*/
public List<SearchResult> reciprocalRankFusion(
List<SearchResult> list1,
List<SearchResult> list2,
int k) {

Map<String, Double> rrfScores = new HashMap<>();

// 计算list1的RRF分数
for (int i = 0; i < list1.size(); i++) {
String id = list1.get(i).getChunkId();
double score = 1.0 / (k + i + 1);
rrfScores.put(id, score);
}

// 累加list2的RRF分数
for (int i = 0; i < list2.size(); i++) {
String id = list2.get(i).getChunkId();
double score = 1.0 / (k + i + 1);
rrfScores.merge(id, score, Double::sum);
}

// 排序返回
return rrfScores.entrySet().stream()
.sorted(Map.Entry.<String, Double>comparingByValue().reversed())
.map(e -> findResultById(e.getKey(), list1, list2))
.collect(Collectors.toList());
}
  1. 分数归一化加权融合:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
public List<SearchResult> normalizedWeightedFusion(
List<SearchResult> vectorResults,
List<SearchResult> textResults) {

// 归一化向量检索分数 (Min-Max归一化)
double vectorMin = vectorResults.stream()
.mapToDouble(SearchResult::getScore).min().orElse(0);
double vectorMax = vectorResults.stream()
.mapToDouble(SearchResult::getScore).max().orElse(1);

Map<String, Double> normalizedVectorScores = vectorResults.stream()
.collect(Collectors.toMap(
SearchResult::getChunkId,
r -> (r.getScore() - vectorMin) / (vectorMax - vectorMin)
));

// 归一化文本检索分数
double textMin = textResults.stream()
.mapToDouble(SearchResult::getScore).min().orElse(0);
double textMax = textResults.stream()
.mapToDouble(SearchResult::getScore).max().orElse(1);

Map<String, Double> normalizedTextScores = textResults.stream()
.collect(Collectors.toMap(
SearchResult::getChunkId,
r -> (r.getScore() - textMin) / (textMax - textMin)
));

// 加权融合 (可根据查询类型动态调整权重)
double vectorWeight = 0.6;
double textWeight = 0.4;

Set<String> allIds = new HashSet<>();
allIds.addAll(normalizedVectorScores.keySet());
allIds.addAll(normalizedTextScores.keySet());

Map<String, Double> finalScores = new HashMap<>();
for (String id : allIds) {
double vectorScore = normalizedVectorScores.getOrDefault(id, 0.0);
double textScore = normalizedTextScores.getOrDefault(id, 0.0);
double finalScore = vectorScore * vectorWeight + textScore * textWeight;
finalScores.put(id, finalScore);
}

return finalScores.entrySet().stream()
.sorted(Map.Entry.<String, Double>comparingByValue().reversed())
.map(e -> findResultById(e.getKey(), vectorResults, textResults))
.collect(Collectors.toList());
}
  1. 自适应权重:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
public double[] calculateAdaptiveWeights(String query, 
List<SearchResult> vectorResults,
List<SearchResult> textResults) {
// 根据查询特征动态调整权重

// 1. 查询长度:短查询偏向关键词,长查询偏向语义
double queryLength = query.length();
double lengthFactor = Math.min(queryLength / 50.0, 1.0);

// 2. 结果一致性:两种检索结果重叠度高,说明都可信
double overlap = calculateOverlap(vectorResults, textResults);

// 3. 查询类型:关键词查询偏向ES,问句查询偏向向量
boolean isQuestion = query.matches(".*[??]$");

double vectorWeight = 0.5 + lengthFactor * 0.2 + (isQuestion ? 0.1 : -0.1);
double textWeight = 1.0 - vectorWeight;

return new double[]{vectorWeight, textWeight};
}

效果对比:

  • RRF算法:召回率提升15%,适用性最广
  • 加权融合:可解释性强,便于调优
  • 自适应权重:不同查询类型表现更均衡

难点5:Elasticsearch中文分词优化

问题描述:

  • 默认分词器对中文支持差
  • 分词颗粒度影响检索效果
  • 专业术语识别不准确

解决方案:

  1. IK分词器配置:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
{
"settings": {
"analysis": {
"analyzer": {
"ik_max_word_analyzer": {
"type": "custom",
"tokenizer": "ik_max_word",
"filter": ["lowercase", "stop"]
},
"ik_smart_analyzer": {
"type": "custom",
"tokenizer": "ik_smart",
"filter": ["lowercase", "stop"]
}
},
"filter": {
"stop": {
"type": "stop",
"stopwords": ["的", "了", "是"]
}
}
}
},
"mappings": {
"properties": {
"content": {
"type": "text",
"analyzer": "ik_max_word_analyzer",
"search_analyzer": "ik_smart_analyzer"
}
}
}
}

说明:

  • ik_max_word: 索引时使用,最细粒度分词,提高召回
  • ik_smart: 查询时使用,粗粒度分词,提高精确度
  1. 自定义词典:
1
2
3
4
5
6
# /etc/elasticsearch/analysis-ik/custom.dic
知识图谱
自然语言处理
机器学习
深度学习
向量数据库
  1. 动态词典热更新:
1
2
3
4
5
6
<!-- IKAnalyzer.cfg.xml -->
<properties>
<entry key="ext_dict">custom.dic</entry>
<entry key="ext_stopwords">stopword.dic</entry>
<entry key="remote_ext_dict">http://yourhost/custom_dict.txt</entry>
</properties>
  1. 同义词处理:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
{
"settings": {
"analysis": {
"filter": {
"synonym_filter": {
"type": "synonym",
"synonyms_path": "analysis/synonyms.txt"
}
},
"analyzer": {
"synonym_analyzer": {
"tokenizer": "ik_smart",
"filter": ["lowercase", "synonym_filter"]
}
}
}
}
}
1
2
3
4
# synonyms.txt
人工智能,AI,Artificial Intelligence
机器学习,ML,Machine Learning
自然语言处理,NLP,Natural Language Processing

效果:

  • 专业术语识别准确率从60%提升到95%
  • 召回率提升25%
  • 支持同义词查询

难点6:并发场景下的数据一致性

问题描述:

  • 同一文档被多次上传/删除
  • 向量化任务重复执行
  • MySQL、ES、Milvus数据不一致

解决方案:

  1. 分布式锁:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
@Service
public class DocumentLockService {

@Autowired
private RedissonClient redissonClient;

public void processWithLock(String docId, Runnable task) {
RLock lock = redissonClient.getLock("doc:lock:" + docId);

try {
// 尝试获取锁,最多等待10秒,锁超时30秒
boolean acquired = lock.tryLock(10, 30, TimeUnit.SECONDS);

if (acquired) {
try {
task.run();
} finally {
lock.unlock();
}
} else {
throw new BusinessException("获取锁失败,文档正在处理中");
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new BusinessException("获取锁被中断");
}
}
}
  1. 幂等性设计:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
@Service
public class VectorizeService {

public void vectorizeChunk(String chunkId, String content) {
// 1. 检查是否已处理
if (isChunkVectorized(chunkId)) {
log.info("Chunk already vectorized, skip: {}", chunkId);
return;
}

// 2. 使用唯一ID防重
String idempotentKey = "vectorize:" + chunkId;
Boolean success = redisTemplate.opsForValue()
.setIfAbsent(idempotentKey, "1", 1, TimeUnit.HOURS);

if (Boolean.FALSE.equals(success)) {
log.info("Duplicate vectorize request, skip: {}", chunkId);
return;
}

try {
// 3. 执行向量化
List<Float> vector = generateEmbedding(content);
milvusService.insertVector(chunkId, vector);
esService.indexDocument(chunkId, content);

// 4. 标记完成
markChunkVectorized(chunkId);

} catch (Exception e) {
// 失败时删除幂等键,允许重试
redisTemplate.delete(idempotentKey);
throw e;
}
}
}
  1. 最终一致性保证:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
@Scheduled(fixedDelay = 300000) // 每5分钟
public void checkDataConsistency() {
// 1. 查询MySQL中已完成的文档
List<DocumentEntity> docs = documentService.findCompleted();

for (DocumentEntity doc : docs) {
// 2. 检查ES中的分块数
long esCount = esService.countChunks(doc.getDocId());

// 3. 检查Milvus中的向量数
long milvusCount = milvusService.countVectors(doc.getDocId());

// 4. 检查MySQL中的分块数
long mysqlCount = documentService.countChunks(doc.getDocId());

// 5. 数据不一致时触发修复
if (esCount != milvusCount || esCount != mysqlCount) {
log.warn("Data inconsistency detected for doc: {}", doc.getDocId());
repairDocument(doc.getDocId());
}
}
}
  1. 分布式事务(Saga模式):
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
public void deleteDocument(String docId) {
// 记录删除事件
SagaLog sagaLog = sagaLogService.create(docId, "DELETE");

try {
// 步骤1: 删除MySQL记录
documentService.delete(docId);
sagaLog.addStep("mysql_deleted");

// 步骤2: 删除ES索引
esService.deleteByDocId(docId);
sagaLog.addStep("es_deleted");

// 步骤3: 删除Milvus向量
milvusService.deleteByDocId(docId);
sagaLog.addStep("milvus_deleted");

// 步骤4: 删除MinIO文件
minioService.deleteFile(docId);
sagaLog.addStep("minio_deleted");

// 标记完成
sagaLog.complete();

} catch (Exception e) {
// 补偿操作
compensateDelete(sagaLog);
throw e;
}
}

private void compensateDelete(SagaLog sagaLog) {
List<String> completedSteps = sagaLog.getCompletedSteps();

// 回滚已完成的步骤
if (completedSteps.contains("minio_deleted")) {
// MinIO无法恢复,记录日志
log.error("MinIO file deleted, cannot restore");
}
if (completedSteps.contains("milvus_deleted")) {
// 向量删除,记录待重建
log.error("Vectors deleted, need rebuild");
}
// ... 其他补偿逻辑
}

难点7:Milvus集群部署与高可用

问题描述:

  • 单点Milvus故障导致服务不可用
  • 数据量大时单机性能瓶颈
  • 数据备份和恢复策略

解决方案:

  1. 分布式部署架构:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
# docker-compose.yml
version: '3.5'

services:
etcd:
image: quay.io/coreos/etcd:v3.5.0
environment:
- ETCD_AUTO_COMPACTION_MODE=revision
- ETCD_AUTO_COMPACTION_RETENTION=1000
- ETCD_QUOTA_BACKEND_BYTES=4294967296

minio:
image: minio/minio:RELEASE.2023-03-20T20-16-18Z
command: minio server /minio_data --console-address ":9001"

milvus-rootcoord:
image: milvusdb/milvus:v2.3.0
command: ["milvus", "run", "rootcoord"]

milvus-querycoord:
image: milvusdb/milvus:v2.3.0
command: ["milvus", "run", "querycoord"]

milvus-querynode-1:
image: milvusdb/milvus:v2.3.0
command: ["milvus", "run", "querynode"]

milvus-querynode-2:
image: milvusdb/milvus:v2.3.0
command: ["milvus", "run", "querynode"]

milvus-datacoord:
image: milvusdb/milvus:v2.3.0
command: ["milvus", "run", "datacoord"]

milvus-datanode:
image: milvusdb/milvus:v2.3.0
command: ["milvus", "run", "datanode"]
  1. 连接池管理:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
@Configuration
public class MilvusConfig {

@Bean
public MilvusServiceClient milvusClient() {
ConnectParam connectParam = ConnectParam.newBuilder()
.withHost("milvus-proxy")
.withPort(19530)
.withConnectTimeout(10, TimeUnit.SECONDS)
.withKeepAliveTime(Long.MAX_VALUE, TimeUnit.NANOSECONDS)
.withKeepAliveTimeout(20, TimeUnit.SECONDS)
.withSecure(false)
.withIdleTimeout(24, TimeUnit.HOURS)
.build();

return new MilvusServiceClient(connectParam);
}
}
  1. 故障转移:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
@Service
public class MilvusFailoverService {

@Autowired
@Qualifier("primaryMilvus")
private MilvusServiceClient primaryClient;

@Autowired
@Qualifier("backupMilvus")
private MilvusServiceClient backupClient;

public SearchResults searchWithFailover(SearchParam param) {
try {
return primaryClient.search(param).get(5, TimeUnit.SECONDS);
} catch (Exception e) {
log.warn("Primary Milvus failed, fallback to backup", e);
return backupClient.search(param).get(5, TimeUnit.SECONDS);
}
}
}
  1. 数据备份策略:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
@Scheduled(cron = "0 0 2 * * ?") // 每天凌晨2点
public void backupMilvusData() {
List<String> collections = milvusClient.listCollections();

for (String collection : collections) {
// 1. 创建备份Collection
String backupName = collection + "_backup_" + LocalDate.now();

// 2. 查询所有向量
QueryResults results = milvusClient.query(
QueryParam.newBuilder()
.withCollectionName(collection)
.withExpr("id > 0")
.withOutFields(Arrays.asList("id", "vector", "doc_id"))
.build()
);

// 3. 插入到备份Collection
milvusClient.insert(InsertParam.newBuilder()
.withCollectionName(backupName)
.withRows(results.getQueryResultsWrapper().getRowRecords())
.build());

// 4. 导出到MinIO
exportToMinio(backupName);

log.info("Backup completed: {}", backupName);
}
}

难点8:系统可观测性建设

问题描述:

  • 分布式系统排查问题困难
  • 性能瓶颈难以定位
  • 缺乏业务指标监控

解决方案:

  1. 链路追踪(Jaeger):
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
@Configuration
public class JaegerConfig {

@Bean
public io.jaegertracing.Configuration jaegerConfig() {
return new io.jaegertracing.Configuration("knowflow")
.withSampler(
new io.jaegertracing.Configuration.SamplerConfiguration()
.withType("const")
.withParam(1) // 采样率100%
)
.withReporter(
new io.jaegertracing.Configuration.ReporterConfiguration()
.withLogSpans(true)
.withSender(
new io.jaegertracing.Configuration.SenderConfiguration()
.withAgentHost("jaeger-agent")
.withAgentPort(6831)
)
);
}
}

@Service
public class SearchService {

@Autowired
private Tracer tracer;

public SearchResult search(String query) {
Span span = tracer.buildSpan("search").start();
span.setTag("query", query);

try (Scope scope = tracer.scopeManager().activate(span)) {
// 向量检索
Span vectorSpan = tracer.buildSpan("vector_search")
.asChildOf(span)
.start();
List<Result> vectorResults = vectorSearch(query);
vectorSpan.finish();

// 全文检索
Span esSpan = tracer.buildSpan("es_search")
.asChildOf(span)
.start();
List<Result> esResults = esSearch(query);
esSpan.finish();

// 合并结果
return mergeResults(vectorResults, esResults);

} finally {
span.finish();
}
}
}
  1. 指标监控(Prometheus):
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
@Configuration
public class MetricsConfig {

@Bean
public MeterRegistryCustomizer<MeterRegistry> metricsCommonTags() {
return registry -> registry.config().commonTags("application", "knowflow");
}

@Bean
public Counter searchCounter(MeterRegistry registry) {
return Counter.builder("search.requests.total")
.description("Total search requests")
.tags("type", "hybrid")
.register(registry);
}

@Bean
public Timer searchTimer(MeterRegistry registry) {
return Timer.builder("search.duration")
.description("Search duration")
.publishPercentiles(0.5, 0.95, 0.99)
.register(registry);
}
}

@Service
public class MonitoredSearchService {

@Autowired
private Counter searchCounter;

@Autowired
private Timer searchTimer;

public SearchResult search(String query) {
searchCounter.increment();

return searchTimer.record(() -> {
return doSearch(query);
});
}
}
  1. 自定义业务指标:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
@Component
public class BusinessMetrics {

private final Gauge documentCount;
private final Gauge vectorCount;
private final Counter uploadCounter;
private final Histogram searchLatency;

public BusinessMetrics(MeterRegistry registry) {
// 文档总数
this.documentCount = Gauge.builder("documents.total", this,
value -> documentService.count())
.description("Total number of documents")
.register(registry);

// 向量总数
this.vectorCount = Gauge.builder("vectors.total", this,
value -> milvusService.count())
.description("Total number of vectors")
.register(registry);

// 上传计数
this.uploadCounter = Counter.builder("document.upload.total")
.description("Total document uploads")
.tag("status", "success")
.register(registry);

// 检索延迟分布
this.searchLatency = Histogram.builder("search.latency")
.description("Search latency distribution")
.buckets(10, 50, 100, 200, 500, 1000, 2000)
.register(registry);
}
}
  1. Grafana仪表盘:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
{
"dashboard": {
"title": "KnowFlow监控",
"panels": [
{
"title": "QPS",
"targets": [
{
"expr": "rate(search_requests_total[1m])"
}
]
},
{
"title": "平均响应时间",
"targets": [
{
"expr": "search_duration_sum / search_duration_count"
}
]
},
{
"title": "P99延迟",
"targets": [
{
"expr": "histogram_quantile(0.99, search_duration_bucket)"
}
]
},
{
"title": "错误率",
"targets": [
{
"expr": "rate(search_errors_total[1m]) / rate(search_requests_total[1m])"
}
]
}
]
}
}
  1. 告警规则:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
groups:
- name: knowflow_alerts
rules:
- alert: HighErrorRate
expr: rate(search_errors_total[5m]) / rate(search_requests_total[5m]) > 0.05
for: 5m
labels:
severity: critical
annotations:
summary: "检索错误率过高"
description: "5分钟内错误率超过5%"

- alert: SlowSearchResponse
expr: histogram_quantile(0.95, search_duration_bucket) > 1000
for: 5m
labels:
severity: warning
annotations:
summary: "检索响应过慢"
description: "P95延迟超过1秒"

- alert: MilvusDown
expr: up{job="milvus"} == 0
for: 1m
labels:
severity: critical
annotations:
summary: "Milvus服务不可用"

难点9:权限控制与数据安全

问题描述:

  • 多租户场景下的数据隔离
  • 细粒度权限控制
  • 敏感数据脱敏

解决方案:

  1. RBAC权限模型:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
@Entity
public class User {
private Long id;
private String username;
private Set<Role> roles;
}

@Entity
public class Role {
private Long id;
private String name;
private Set<Permission> permissions;
}

@Entity
public class Permission {
private Long id;
private String resource; // document, search, admin
private String action; // create, read, update, delete
}
  1. 数据权限过滤:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
@Aspect
@Component
public class DataPermissionAspect {

@Around("@annotation(dataPermission)")
public Object checkPermission(ProceedingJoinPoint joinPoint,
DataPermission dataPermission) throws Throwable {
// 获取当前用户
UserInfo currentUser = SecurityContextHolder.getCurrentUser();

// 修改查询条件,添加数据权限过滤
Object[] args = joinPoint.getArgs();
for (Object arg : args) {
if (arg instanceof SearchRequest) {
SearchRequest request = (SearchRequest) arg;
// 添加用户ID过滤条件
request.addFilter("userId", currentUser.getUserId());
// 添加部门ID过滤条件(如果有)
if (currentUser.getDepartmentId() != null) {
request.addFilter("departmentId", currentUser.getDepartmentId());
}
}
}

return joinPoint.proceed(args);
}
}

@Service
public class DocumentService {

@DataPermission
public List<Document> searchDocuments(SearchRequest request) {
// 自动添加了数据权限过滤
return documentRepository.search(request);
}
}
  1. 多租户数据隔离:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
@Configuration
public class TenantConfig {

@Bean
public TenantIdentifierResolver tenantIdentifierResolver() {
return new TenantIdentifierResolver() {
@Override
public String resolveCurrentTenantIdentifier() {
// 从上下文获取租户ID
return TenantContext.getCurrentTenantId();
}
};
}
}

@Service
public class MilvusMultiTenantService {

public void insert(String tenantId, InsertData data) {
// 租户专属Collection
String collectionName = "tenant_" + tenantId + "_vectors";

// 确保Collection存在
if (!milvusClient.hasCollection(collectionName)) {
createTenantCollection(collectionName);
}

milvusClient.insert(collectionName, data);
}

public SearchResults search(String tenantId, SearchParam param) {
String collectionName = "tenant_" + tenantId + "_vectors";
return milvusClient.search(collectionName, param);
}
}
  1. 敏感数据脱敏:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
@Component
public class DataMaskingService {

public String maskSensitiveData(String content) {
// 手机号脱敏
content = content.replaceAll("(\\d{3})\\d{4}(\\d{4})", "$1****$2");

// 身份证号脱敏
content = content.replaceAll("(\\d{6})\\d{8}(\\d{4})", "$1********$2");

// 邮箱脱敏
content = content.replaceAll("(\\w{2})\\w+(\\w@\\w+\\.\\w+)", "$1***$2");

return content;
}

public SearchResult maskSearchResult(SearchResult result) {
result.setContent(maskSensitiveData(result.getContent()));
return result;
}
}

难点10:成本优化

问题描述:

  • Embedding API调用成本高
  • 存储成本随数据量增长
  • 计算资源利用率低

解决方案:

  1. 本地Embedding模型:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
@Service
public class LocalEmbeddingService {

private OnnxModel model;

@PostConstruct
public void init() {
// 加载本地模型(如text2vec-base-chinese)
this.model = OnnxModel.load("models/text2vec-base-chinese.onnx");
}

public List<Float> generateEmbedding(String text) {
// 文本预处理
int[] tokens = tokenize(text);

// 模型推理
float[][] output = model.predict(tokens);

// 转换为List
return Arrays.stream(output[0])
.boxed()
.collect(Collectors.toList());
}
}

成本对比:

  • OpenAI API: $0.0001/1K tokens
  • 本地模型: 仅GPU成本,约降低90%
  1. 向量压缩:
1
2
3
4
5
6
7
8
9
10
11
12
13
# 使用PQ(Product Quantization)压缩向量
index_params = {
"index_type": "IVF_PQ",
"metric_type": "IP",
"params": {
"nlist": 1024, # 聚类中心数
"m": 8, # 子向量数
"nbits": 8 # 每个子向量的位数
}
}

# 存储空间: 768维 * 4字节 = 3KB -> 8字节
# 压缩比: 384倍
  1. 冷热数据分离:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
@Scheduled(cron = "0 0 3 * * ?")
public void archiveOldDocuments() {
// 1. 查找3个月未访问的文档
LocalDateTime threshold = LocalDateTime.now().minusMonths(3);
List<Document> oldDocs = documentService.findNotAccessedSince(threshold);

for (Document doc : oldDocs) {
// 2. 从Milvus删除向量
milvusService.delete(doc.getDocId());

// 3. 从ES删除索引
esService.delete(doc.getDocId());

// 4. 标记为已归档
doc.setStatus(DocStatus.ARCHIVED);
doc.setArchiveTime(LocalDateTime.now());
documentService.update(doc);

// 5. MinIO文件保留(便宜)
}

log.info("Archived {} documents", oldDocs.size());
}

// 访问归档文档时恢复
public Document accessArchivedDocument(String docId) {
Document doc = documentService.findById(docId);

if (doc.getStatus() == DocStatus.ARCHIVED) {
// 异步恢复索引
asyncRestoreDocument(docId);

// 直接返回原文件
return doc;
}

return doc;
}
  1. 缓存优化:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
@Service
public class CachedSearchService {

@Autowired
private RedisTemplate<String, Object> redisTemplate;

public SearchResult search(String query) {
// 1. 生成缓存Key
String cacheKey = "search:" + DigestUtils.md5Hex(query);

// 2. 查询缓存
SearchResult cached = (SearchResult) redisTemplate.opsForValue().get(cacheKey);
if (cached != null) {
return cached;
}

// 3. 执行检索
SearchResult result = doSearch(query);

// 4. 缓存结果(1小时)
redisTemplate.opsForValue().set(cacheKey, result, 1, TimeUnit.HOURS);

return result;
}

// Embedding缓存
@Cacheable(value = "embeddings", key = "#text.hashCode()")
public List<Float> cachedEmbedding(String text) {
return generateEmbedding(text);
}
}

成本优化效果:

  • Embedding成本降低90%
  • 存储成本降低70%
  • 整体运营成本降低60%

七、性能优化方案

7.1 检索性能优化

优化前

  • 平均响应时间: 1200ms
  • P95延迟: 2500ms
  • QPS: 50

优化措施

  1. 向量检索优化:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
# 优化前: IVF_FLAT索引
index_params = {
"index_type": "IVF_FLAT",
"params": {"nlist": 1024}
}
# 检索耗时: 800ms

# 优化后: HNSW索引
index_params = {
"index_type": "HNSW",
"params": {"M": 16, "efConstruction": 200}
}
search_params = {"ef": 64}
# 检索耗时: 150ms
  1. 并行检索:
1
2
3
4
5
6
7
8
9
10
// 优化前: 串行执行
List<Result> vectorResults = vectorSearch(query); // 150ms
List<Result> esResults = esSearch(query); // 200ms
// 总耗时: 350ms

// 优化后: 并行执行
CompletableFuture<List<Result>> future1 = CompletableFuture.supplyAsync(() -> vectorSearch(query));
CompletableFuture<List<Result>> future2 = CompletableFuture.supplyAsync(() -> esSearch(query));
CompletableFuture.allOf(future1, future2).join();
// 总耗时: 200ms (取最大值)
  1. 缓存策略:
1
2
3
4
5
6
7
8
9
10
// L1缓存: 本地Caffeine缓存(热点查询)
@Cacheable(value = "local", maximumSize = 1000, expireAfterWrite = "5m")
public SearchResult localCache(String query) {...}

// L2缓存: Redis缓存(常见查询)
@Cacheable(value = "redis", ttl = "1h")
public SearchResult redisCache(String query) {...}

// 缓存命中率: 65%
// 命中后响应时间: 10ms
  1. 连接池优化:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
# Elasticsearch连接池
spring:
elasticsearch:
rest:
connection-timeout: 5000
read-timeout: 30000
uris: http://es-node1:9200,http://es-node2:9200

# HikariCP数据库连接池
spring:
datasource:
hikari:
maximum-pool-size: 50
minimum-idle: 10
connection-timeout: 30000
idle-timeout: 600000

优化后

  • 平均响应时间: 180ms(提升85%)
  • P95延迟: 350ms(提升86%)
  • QPS: 500(提升10倍)

7.2 文档处理性能优化

优化前

  • 单文档处理时间: 30秒
  • 并发处理能力: 5个/分钟
  • 大文件(>50MB)容易OOM

优化措施

  1. 流式处理:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
// 优化前: 全量加载到内存
String content = FileUtils.readFileToString(file);
List<Chunk> chunks = chunkDocument(content);

// 优化后: 流式读取
try (BufferedReader reader = new BufferedReader(new FileReader(file))) {
StreamingChunker chunker = new StreamingChunker(500, 100);

String line;
while ((line = reader.readLine()) != null) {
chunker.addLine(line);

// 每生成一个chunk就处理
if (chunker.hasChunk()) {
Chunk chunk = chunker.getNextChunk();
processChunk(chunk);
}
}
}
  1. 批量向量化:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// 优化前: 逐个调用Embedding API
for (Chunk chunk : chunks) {
List<Float> vector = embeddingService.generate(chunk.getContent());
milvusService.insert(chunk.getId(), vector);
}
// 100个chunk耗时: 20秒

// 优化后: 批量调用
List<String> texts = chunks.stream()
.map(Chunk::getContent)
.collect(Collectors.toList());

List<List<Float>> vectors = embeddingService.batchGenerate(texts);
milvusService.batchInsert(chunks, vectors);
// 100个chunk耗时: 3秒
  1. 异步处理:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
@Configuration
public class AsyncConfig {

@Bean
public Executor taskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(10);
executor.setMaxPoolSize(20);
executor.setQueueCapacity(500);
executor.setThreadNamePrefix("doc-process-");
executor.initialize();
return executor;
}
}

@Service
public class AsyncDocumentService {

@Async("taskExecutor")
public CompletableFuture<Void> processDocument(String docId) {
// 异步处理文档
return CompletableFuture.completedFuture(null);
}
}
  1. 分片处理大文件:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
public void processLargeFile(File file) {
long fileSize = file.length();
int chunkSize = 10 * 1024 * 1024; // 10MB per chunk
int numChunks = (int) Math.ceil((double) fileSize / chunkSize);

ExecutorService executor = Executors.newFixedThreadPool(4);
List<Future<Void>> futures = new ArrayList<>();

for (int i = 0; i < numChunks; i++) {
final int chunkIndex = i;
Future<Void> future = executor.submit(() -> {
long start = (long) chunkIndex * chunkSize;
long end = Math.min(start + chunkSize, fileSize);
processFileChunk(file, start, end);
return null;
});
futures.add(future);
}

// 等待所有分片处理完成
for (Future<Void> future : futures) {
future.get();
}

executor.shutdown();
}

优化后

  • 单文档处理时间: 5秒(提升83%)
  • 并发处理能力: 50个/分钟(提升10倍)
  • 支持1GB+大文件

7.3 数据库查询优化

优化措施

  1. 索引优化:
1
2
3
4
5
6
7
8
9
10
11
-- 文档表索引
CREATE INDEX idx_doc_user_status ON documents(user_id, status);
CREATE INDEX idx_doc_created ON documents(created_time);
CREATE INDEX idx_doc_category ON documents(category);

-- 分块表索引
CREATE INDEX idx_chunk_doc ON chunks(doc_id);
CREATE INDEX idx_chunk_status ON chunks(status);

-- 使用覆盖索引
CREATE INDEX idx_doc_list ON documents(user_id, status, created_time, doc_id, title);
  1. 分页优化:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
// 优化前: OFFSET分页(慢)
SELECT * FROM documents
WHERE user_id = ?
ORDER BY created_time DESC
LIMIT 100 OFFSET 10000;
// 耗时: 500ms

// 优化后: 游标分页(快)
SELECT * FROM documents
WHERE user_id = ?
AND created_time < ? -- 上一页最后一条的时间
ORDER BY created_time DESC
LIMIT 100;
// 耗时: 10ms
  1. 批量查询:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
// 优化前: N+1查询
List<Document> docs = documentRepository.findByUserId(userId);
for (Document doc : docs) {
List<Chunk> chunks = chunkRepository.findByDocId(doc.getDocId());
doc.setChunks(chunks);
}
// 查询次数: 1 + N

// 优化后: 批量查询
List<Document> docs = documentRepository.findByUserId(userId);
List<String> docIds = docs.stream()
.map(Document::getDocId)
.collect(Collectors.toList());
List<Chunk> allChunks = chunkRepository.findByDocIdIn(docIds);

Map<String, List<Chunk>> chunkMap = allChunks.stream()
.collect(Collectors.groupingBy(Chunk::getDocId));

docs.forEach(doc -> doc.setChunks(chunkMap.get(doc.getDocId())));
// 查询次数: 2
  1. 读写分离:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
spring:
datasource:
master:
url: jdbc:mysql://master:3306/knowflow
username: root
password: xxx
slave:
url: jdbc:mysql://slave:3306/knowflow
username: readonly
password: xxx

# 自动路由
@Service
public class DocumentService {

@ReadOnly // 读从库
public List<Document> findAll() {...}

@WriteOnly // 写主库
public void save(Document doc) {...}
}

7.4 系统整体性能指标

指标 优化前 优化后 提升
文档上传处理 30s 5s 83%
检索平均响应时间 1200ms 180ms 85%
检索P95延迟 2500ms 350ms 86%
系统QPS 50 500 900%
并发文档处理 5/min 50/min 900%
数据库查询 500ms 10ms 98%
缓存命中率 0% 65% -
服务可用性 99.5% 99.95% 0.45pp

八、100个面试问题及详细解答

分类一:项目背景和架构(15题)

Q1: 请介绍一下KnowFlow项目

标准答案:
KnowFlow是一个企业级RAG知识库系统,主要解决企业内部文档分散、查找困难的问题。项目采用混合检索架构,结合Elasticsearch的全文检索和Milvus的向量检索,实现智能化的文档检索和问答。

技术栈方面,后端使用SpringBoot,Elasticsearch做全文检索,Milvus做向量数据库,MinIO做对象存储,配合Prometheus做监控。整个系统支持百万级文档存储,检索响应时间在200ms以内。

我在项目中主要负责混合检索模块的设计和实现,文档上传和解析的异步处理,以及系统的性能优化和监控体系建设。通过优化,将检索响应时间从1.2秒降至180ms,系统QPS从50提升到500。

追问1: 为什么选择RAG方案而不是微调模型?

追问答案:
主要有三个原因:

  1. 成本考虑: 微调需要大量标注数据和GPU资源,成本高。RAG只需要存储文档和调用Embedding API,成本低很多
  2. 实时性: 企业文档更新频繁,RAG可以实时添加新文档,微调需要重新训练
  3. 可解释性: RAG返回结果时会展示引用的原始文档,用户可以验证答案来源,而微调模型是黑盒

追问2: RAG有什么局限性?

追问答案:

  1. 检索质量依赖: 如果检索不到相关文档,LLM就无法生成正确答案
  2. 上下文长度限制: 只能传入有限的文档片段,长文档可能信息不全
  3. 响应延迟: 需要先检索再生成,比直接调用模型慢
  4. 跨文档推理弱: 需要综合多个文档才能回答的问题效果不好

知识扩展:

  • RAG的三种模式:Naive RAG、Advanced RAG、Modular RAG
  • RAG评估指标:召回率、精确率、答案准确性、引用准确性
  • RAG优化方向:查询改写、混合检索、重排序、上下文压缩

Q2: 说说你们的技术架构设计思路

标准答案:
我们采用分层微服务架构,主要分为四层:

  1. 接入层: Nginx和Spring Cloud Gateway,负责负载均衡、认证鉴权、限流
  2. 应用层: 按业务拆分为文档服务、检索服务、解析服务、向量服务等独立服务,每个服务可独立部署和扩展
  3. 数据层: MySQL存元数据,Redis做缓存,Elasticsearch做全文检索,Milvus做向量检索,MinIO做对象存储,各司其职
  4. 基础设施层: Prometheus监控、ELK日志、Jaeger链路追踪

异步处理方面,使用RabbitMQ解耦文档上传、解析、向量化三个阶段。用户上传后立即返回,后台异步处理,提升用户体验。

这种架构的优点是:各组件职责清晰、易于扩展、故障隔离好。

追问1: 为什么不用单体架构?

追问答案:
我们确实考虑过单体,但有几个问题:

  1. 扩展性: 文档解析是CPU密集型,检索是IO密集型,需求不同,单体无法针对性扩展
  2. 技术栈: Milvus是Python生态,而我们主系统用Java,微服务更方便集成
  3. 团队协作: 团队有5个人,微服务可以并行开发,提高效率
  4. 故障隔离: 解析服务崩溃不会影响检索服务

追问2: 微服务带来了什么挑战?

追问答案:

  1. 分布式事务: 文档删除需要同时删除MySQL、ES、Milvus的数据,我们用Saga模式处理
  2. 链路追踪: 一个请求跨多个服务,排查问题困难,引入Jaeger解决
  3. 数据一致性: 各个数据源可能不一致,我们有定时任务做一致性校验和修复
  4. 运维复杂度: 需要管理多个服务的部署、监控、日志

知识扩展:

  • 微服务的拆分原则:单一职责、高内聚低耦合、按业务边界
  • 服务间通信方式:REST、RPC、消息队列
  • 微服务治理:服务注册发现、配置中心、熔断降级

Q3: 为什么选择Elasticsearch做全文检索?

标准答案:
主要原因有四点:

  1. 倒排索引性能好: ES基于Lucene的倒排索引,全文检索速度快,百万级文档毫秒级响应
  2. 中文分词支持: 集成IK Analyzer分词器,支持中文分词,可以自定义词典
  3. 分布式架构: 天然支持分片和副本,容易水平扩展
  4. 生态成熟: 社区活跃,插件丰富,与ELK栈集成方便

对比其他方案:

  • MySQL全文索引: 性能差,不支持分布式,中文支持弱
  • Solr: 功能类似,但ES更轻量,JSON友好,社区更活跃
  • 纯向量检索: 无法处理精确关键词匹配,需要混合检索

追问1: ES的倒排索引是什么原理?

追问答案:
倒排索引是从词到文档的映射。传统索引是”文档→词”,倒排索引是”词→文档列表”。

举例:

1
2
3
4
5
6
7
8
文档1: "机器学习很有趣"
文档2: "深度学习是机器学习的分支"

倒排索引:
机器学习 -> [文档1, 文档2]
深度学习 -> [文档2]
有趣 -> [文档1]
分支 -> [文档2]

查询”机器学习”时,直接通过倒排索引找到[文档1, 文档2],无需遍历所有文档。

倒排索引还包含词频、位置等信息,用于相关性评分(BM25算法)。

追问2: ES如何做到近实时搜索?

追问答案:
ES的近实时(NRT)搜索依赖refresh机制:

  1. 文档写入先进入内存buffer
  2. 每隔1秒(默认)执行refresh,将buffer写入segment文件(内存中)
  3. segment立即可搜索,但还未持久化
  4. 后台定期将segment flush到磁盘

所以新文档写入后1秒内可搜索,这就是”近实时”。可以调整refresh_interval参数平衡实时性和性能。

知识扩展:

  • ES写入流程:协调节点 → 主分片 → 副本分片
  • ES查询流程:协调节点 → 所有分片 → 汇总结果
  • BM25算法:TF-IDF的改进版,考虑文档长度归一化

Q4: 为什么选择Milvus作为向量数据库?

标准答案:
Milvus是专为向量检索设计的数据库,我们选择它主要基于:

  1. 性能优异: 支持HNSW、IVF等多种索引,百万级向量毫秒级检索
  2. 易于使用: 提供完善的SDK,API设计友好,上手快
  3. 可扩展: 支持分布式部署,可以水平扩展
  4. 开源免费: 可以私有化部署,符合企业数据安全要求
  5. 社区活跃: 国产开源项目,中文文档完善,社区响应快

对比其他方案:

  • Faiss: 单机库,不支持分布式,无法满足大规模需求
  • Pinecone: 商业产品,无法私有化部署,成本高
  • Elasticsearch向量检索: 向量检索性能不如专业向量库

追问1: Milvus底层用的什么索引结构?

追问答案:
Milvus支持多种索引,我们主要用HNSW(Hierarchical Navigable Small World):

HNSW是基于图的索引结构,核心思想:

  1. 分层结构: 类似跳表,上层节点稀疏,下层节点密集
  2. 小世界网络: 每个节点连接若干邻居节点,形成高效导航图
  3. 贪心搜索: 从顶层开始,每次选择最近的邻居节点,逐层下降

优点:检索速度快(对数级),召回率高(95%+)
缺点:构建索引慢,内存占用大

我们还用过IVF_PQ(倒排+乘积量化):

  • IVF: 将向量聚类,检索时只搜索相关的几个聚类
  • PQ: 压缩向量,减少内存和计算量
  • 适合大规模向量且对召回率要求不极致的场景

追问2: 向量检索的相似度计算方式有哪些?

追问答案:
主要三种:

  1. 欧氏距离(L2):

    • 公式:sqrt(Σ(ai-bi)²)
    • 越小越相似,适合空间距离场景
  2. 余弦相似度(Cosine):

    • 公式:(A·B)/(||A||*||B||)
    • 值域[-1,1],越大越相似,适合文本向量
  3. 内积(IP):

    • 公式:Σ(ai*bi)
    • 越大越相似,归一化后等价于余弦相似度

我们用内积(IP),因为OpenAI的Embedding是归一化的,内积计算最快。

知识扩展:

  • 向量检索的ANN问题:Approximate Nearest Neighbor
  • 其他索引:LSH、Annoy、ScaNN
  • 向量维度:OpenAI ada-002是1536维,text-embedding-3-small是1536维

Q5: 介绍一下混合检索的实现原理

标准答案:
混合检索是结合稀疏检索(全文)和密集检索(向量)的方案。我们的实现流程:

  1. 查询处理: 用户输入查询词,生成查询向量(调用Embedding API)
  2. 并行检索:
    • Milvus向量检索:Top 50,基于语义相似度
    • ES全文检索:Top 50,基于关键词匹配
  3. 结果融合: 使用RRF算法融合两路结果
  4. 重排序: 可选的Reranker进一步优化排序
  5. 返回Top 10

RRF算法公式:score(d) = Σ 1/(k + rank_i(d))

  • k是常数(通常60)
  • rank_i(d)是文档d在第i路检索的排名
  • 只依赖排名,不依赖具体分数,避免分数量纲问题

追问1: 为什么要混合检索,单用向量检索不行吗?

追问答案:
各有优劣:

向量检索优势:

  • 语义理解:能理解同义词、近义表达
  • 模糊匹配:不需要精确关键词

向量检索劣势:

  • 精确匹配差:查询”订单号12345”,向量可能召回无关文档
  • 专有名词弱:人名、地名、产品型号等

全文检索优势:

  • 精确匹配:关键词、ID、专有名词检索准确
  • 可解释性:知道为什么匹配(高亮关键词)

全文检索劣势:

  • 词汇鸿沟:查询”如何提升销量”无法匹配”增加营收的方法”
  • 需要精确关键词

混合检索取长补短,召回率比单一方法提升30-40%。

追问2: 除了RRF,还有什么融合算法?

追问答案:

  1. 加权求和:

    1
    score = α * vector_score + β * text_score

    问题:两个分数量纲不同,需要归一化

  2. CombSUM/CombMNZ:

    1
    2
    CombSUM: score = Σ scores
    CombMNZ: score = (Σ scores) * num_systems

    CombMNZ对多系统都返回的文档加权

  3. 学习排序(LTR):
    用机器学习模型学习最优融合策略,需要标注数据

我们选RRF因为:

  • 不需要调参(α、β)
  • 不依赖分数量纲
  • 效果稳定,适用性广

知识扩展:

  • Dense Retrieval: BERT-based模型,如DPR、ANCE
  • Sparse Retrieval: BM25、TF-IDF
  • Hybrid Retrieval: ColBERT、SPLADE

Q6: 项目的数据流是怎样的?

标准答案:
分为两条主要数据流:

1. 文档写入流:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
用户上传 → 网关验证 → 文档服务

MinIO存原文件

发送MQ消息 → 解析服务

提取文本 + 分块

发送MQ消息 → 向量服务

┌────────┴────────┐
Milvus(向量) Elasticsearch(文本)

更新MySQL状态(完成)

2. 检索读流:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
用户查询 → 网关 → 检索服务

查Redis缓存 → 命中返回
↓ 未命中
生成查询向量

┌────┴────┐
Milvus检索 ES检索(并行)

RRF融合排序

从MySQL加载元数据

写入Redis缓存

返回结果

追问1: 为什么用消息队列,直接同步调用不行吗?

追问答案:
消息队列带来几个好处:

  1. 异步解耦: 上传立即返回,不用等处理完成,用户体验好
  2. 削峰填谷: 高峰期请求堆积在队列,慢慢消费,不会压垮后端
  3. 失败重试: 处理失败自动重试,保证最终成功
  4. 扩展性: 可以动态增加消费者,提高并发处理能力

如果同步调用:

  • 用户等待时间长(30秒+)
  • 解析服务压力大,容易崩溃
  • 失败需要前端重新上传

追问2: 消息队列选型为什么用RabbitMQ?

追问答案:
对比了三种方案:

RabbitMQ:

  • 优点:功能完善(死信队列、延迟队列),消息可靠性高,管理界面好用
  • 缺点:吞吐量一般(几万QPS)
  • 适合:对可靠性要求高的场景

Kafka:

  • 优点:吞吐量极高(百万QPS),适合日志、流处理
  • 缺点:不支持优先级队列,消息堆积时消费延迟高
  • 适合:大数据场景

RocketMQ:

  • 优点:吞吐量高,支持事务消息
  • 缺点:社区不如前两者,运维复杂
  • 适合:电商、金融场景

我们选RabbitMQ因为:

  • 文档处理量不大(几千QPS够用)
  • 需要死信队列处理失败任务
  • 团队熟悉,运维成本低

知识扩展:

  • 消息可靠性:生产者确认、持久化、消费者确认
  • 顺序消息:单队列 + 单消费者
  • 幂等性:消息去重,业务幂等设计

Q7: 如何保证系统的高可用?

标准答案:
我们从多个层面保证高可用:

1. 服务层:

  • 多实例部署:每个服务至少2个实例
  • 健康检查:Spring Actuator健康检查,异常实例自动摘除
  • 熔断降级:Sentinel熔断,避免雪崩
  • 限流:令牌桶限流,保护系统

2. 数据层:

  • MySQL主从:主从复制 + 读写分离
  • Redis哨兵:自动故障转移
  • ES集群:5个节点,1个副本,任意节点故障不影响服务
  • Milvus集群:Query Node多副本,某个节点故障自动迁移

3. 基础设施:

  • 负载均衡:Nginx + 健康检查
  • 容器编排:Kubernetes自动重启故障Pod
  • 监控告警:Prometheus + AlertManager,故障及时通知

4. 灾备:

  • 数据备份:MySQL每日全量备份,binlog实时备份
  • 异地容灾:核心数据异地备份

目前系统可用性达到99.95%,月故障时间<22分钟。

追问1: 如果Milvus集群全部挂了怎么办?

追问答案:
设计了降级方案:

  1. 检测: 健康检查发现Milvus不可用
  2. 降级: 自动切换为纯ES检索模式
  3. 通知: 发送告警,运维介入
  4. 恢复: Milvus恢复后,自动切回混合检索

降级后:

  • 检索功能正常,但召回率下降约20%
  • 不影响文档上传(向量化任务堆积在MQ)
  • 用户无感知或提示”部分功能降级”

代码实现:

1
2
3
4
5
6
7
8
public SearchResult search(String query) {
if (milvusHealthCheck.isHealthy()) {
return hybridSearch(query);
} else {
log.warn("Milvus unavailable, fallback to ES only");
return esOnlySearch(query);
}
}

追问2: 如何做灰度发布?

追问答案:
使用金丝雀发布策略:

  1. 部署新版本: 先部署1个新版本实例(10%流量)
  2. 监控指标: 观察错误率、响应时间、业务指标
  3. 逐步放量: 无异常则扩展到50% → 100%
  4. 快速回滚: 发现问题立即回滚到旧版本

具体实现:

  • Kubernetes Deployment滚动更新
  • Istio流量分配(10% → 新版本)
  • Prometheus监控新版本指标
  • 自动化脚本一键回滚

好处:

  • 降低发布风险
  • 快速验证新功能
  • 问题影响面小

知识扩展:

  • 高可用架构模式:主备、主从、集群、分片
  • CAP理论:一致性、可用性、分区容错性
  • 故障演练:混沌工程、故障注入

Q8: 系统的性能瓶颈在哪里?如何优化的?

标准答案:
经过压测和监控分析,发现三个主要瓶颈:

瓶颈1:向量检索慢

  • 问题:IVF_FLAT索引,百万向量检索800ms
  • 优化:改用HNSW索引,耗时降至150ms
  • 效果:检索速度提升5倍

瓶颈2:Embedding API调用慢

  • 问题:逐个调用API,100个chunk需要20秒
  • 优化:批量调用API(每批32个),耗时降至3秒
  • 效果:向量化速度提升6倍

瓶颈3:ES查询慢

  • 问题:深度分页查询,OFFSET 10000耗时500ms
  • 优化:改用scroll API或search_after,耗时降至10ms
  • 效果:分页速度提升50倍

其他优化:

  • Redis缓存热点查询,命中率65%,命中后10ms返回
  • 数据库索引优化,查询从100ms降至5ms
  • 并行检索,Milvus和ES并行,耗时取max而非sum

整体效果:

  • 检索P95延迟:2500ms → 350ms
  • 系统QPS:50 → 500
  • 文档处理时间:30s → 5s

追问1: 如何发现这些性能瓶颈的?

追问答案:
使用了多种工具和方法:

  1. APM工具: Spring Boot Actuator + Micrometer

    • 记录每个接口的响应时间
    • 发现search接口P95是2500ms
  2. 链路追踪: Jaeger

    • 追踪一次完整请求,看到向量检索占800ms
    • 定位到Milvus是瓶颈
  3. 慢查询日志:

    • MySQL慢查询日志,发现深度分页问题
    • ES慢查询日志,发现某些复杂查询慢
  4. 压测: JMeter

    • 并发100压测,QPS只有50
    • CPU使用率不高,说明有等待
  5. 性能分析: JProfiler

    • 采样发现大量时间在网络IO
    • 定位到Embedding API调用是串行的

追问2: 还有哪些可以优化的点?

追问答案:

  1. 本地Embedding模型:

    • 当前用OpenAI API,有网络延迟
    • 可以部署本地模型(如text2vec),延迟降至10ms
    • 成本也大幅降低
  2. 向量压缩:

    • 当前768维float32,每向量3KB
    • 可以用PQ压缩到8字节,空间降低384倍
    • 检索速度也会提升
  3. 预计算:

    • 热门查询的向量可以预计算缓存
    • 避免重复调用Embedding API
  4. CDN:

    • 文档原文件走CDN加速
    • 减轻MinIO压力
  5. 冷热分离:

    • 3个月未访问的文档归档
    • 节省Milvus和ES资源

知识扩展:

  • 性能优化的指导原则:先测量再优化,避免过早优化
  • Amdahl定律:优化效果取决于瓶颈占比
  • 性能测试工具:JMeter、Gatling、Locust

Q9: 如何设计文档的权限控制?

标准答案:
我们设计了基于RBAC的多级权限控制:

1. 角色定义:

  • 管理员(Admin):所有权限
  • 编辑者(Editor):上传、编辑、删除自己的文档
  • 查看者(Viewer):只能检索和查看

2. 数据权限:

  • 个人文档:只有创建者可见
  • 部门文档:部门成员可见
  • 公开文档:所有人可见

3. 实现方式:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// 数据权限过滤
@Aspect
public class DataPermissionAspect {
@Around("@annotation(DataPermission)")
public Object filter(ProceedingJoinPoint joinPoint) {
UserInfo user = getCurrentUser();

// 在查询条件中添加权限过滤
SearchRequest request = (SearchRequest) joinPoint.getArgs()[0];
request.addFilter("userId", user.getId());
request.addFilter("departmentId", user.getDepartmentId());
request.addFilter("scope", "public");

return joinPoint.proceed();
}
}

4. Milvus权限隔离:

  • 使用分区(Partition)隔离不同部门的数据
  • 检索时指定用户有权访问的分区列表

5. ES权限过滤:

  • 在查询DSL中添加过滤条件
1
2
3
4
5
6
7
8
9
10
11
12
{
"bool": {
"must": [...],
"filter": [
{
"terms": {
"scope": ["personal", "public"]
}
}
]
}
}

追问1: 如何处理多租户场景?

追问答案:
多租户有三种隔离方案:

方案1:共享数据库,共享Schema

  • 所有租户数据在同一张表
  • 通过tenant_id字段区分
  • 优点:成本低,易维护
  • 缺点:数据隔离性差,性能互相影响

方案2:共享数据库,独立Schema

  • 每个租户一个Schema(PostgreSQL)
  • 数据物理隔离,逻辑共享
  • 优点:隔离性好,成本适中
  • 缺点:Schema数量受限

方案3:独立数据库

  • 每个租户一个数据库实例
  • 完全隔离
  • 优点:隔离性最好,性能独立
  • 缺点:成本高,运维复杂

我们的选择:

  • 小租户:方案1,通过tenant_id过滤
  • 大租户:方案3,独立Milvus Collection

Milvus实现:

1
2
String collectionName = "tenant_" + tenantId + "_vectors";
milvusClient.search(collectionName, queryVector);

追问2: 如何防止越权访问?

追问答案:

  1. 网关层鉴权:

    • JWT Token验证
    • 解析出用户ID和权限
    • 非法token直接拒绝
  2. 服务层鉴权:

    • 接口级别:@PreAuthorize(“hasRole(‘ADMIN’)”)
    • 数据级别:DataPermissionAspect自动过滤
  3. 数据层隔离:

    • SQL强制加上userId条件
    • 防止忘记加权限导致越权
  4. 审计日志:

    • 记录所有数据访问
    • 异常访问告警

防御深度原则:多层防护,即使某一层失效也有其他层保护。

知识扩展:

  • RBAC:Role-Based Access Control
  • ABAC:Attribute-Based Access Control
  • OAuth2.0:授权框架
  • JWT:JSON Web Token

Q10: 介绍一下监控体系的设计

标准答案:
我们建立了立体化的监控体系,分为四层:

1. 基础设施监控(Prometheus + Node Exporter):

  • CPU、内存、磁盘、网络
  • 告警:CPU>80%持续5分钟

2. 中间件监控:

  • MySQL:QPS、连接数、慢查询
  • Redis:命中率、内存使用、阻塞操作
  • ES:查询延迟、索引速度、集群状态
  • Milvus:检索延迟、向量数、内存使用
  • 告警:ES集群状态非green

3. 应用监控(Spring Boot Actuator + Micrometer):

  • 接口QPS、响应时间、错误率
  • JVM:堆内存、GC次数、线程数
  • 自定义指标:文档数、检索量
  • 告警:接口错误率>5%

4. 业务监控:

  • 文档上传成功率
  • 检索平均响应时间
  • 用户活跃度
  • 告警:上传成功率<95%

链路追踪(Jaeger):

  • 追踪请求完整链路
  • 定位慢请求瓶颈
  • 依赖关系可视化

日志聚合(ELK):

  • 集中收集所有服务日志
  • 全文检索日志
  • 日志分析和报表

Grafana仪表盘:

  • 实时监控大屏
  • 多维度图表展示
  • 异常自动高亮

追问1: 如何设计告警规则,避免告警风暴?

追问答案:

  1. 告警分级:

    • P0(紧急):服务完全不可用,立即处理
    • P1(重要):核心功能受影响,1小时内处理
    • P2(一般):非核心功能异常,工作时间处理
    • P3(提示):潜在风险,定期巡检
  2. 告警聚合:

    • 同一告警5分钟内只发一次
    • 相关告警合并(ES集群故障 + 检索失败)
  3. 告警抑制:

    • 底层告警触发时,抑制上层告警
    • 如:主机宕机,抑制该主机上服务的告警
  4. 告警静默:

    • 已知问题处理中,临时静默
    • 计划维护期间,静默相关告警
  5. 智能告警:

    • 基于历史数据的异常检测
    • 避免固定阈值的误报

Prometheus配置:

1
2
3
4
5
6
7
8
9
10
11
12
route:
group_by: ['alertname', 'cluster']
group_wait: 30s
group_interval: 5m
repeat_interval: 12h

inhibit_rules:
- source_match:
severity: 'critical'
target_match:
severity: 'warning'
equal: ['instance']

追问2: 如何做容量规划?

追问答案:
基于监控数据进行容量评估:

  1. 数据容量:

    • 当前文档数:50万
    • 增长速度:1万/天
    • 预估:1年后450万文档
    • 需求:ES和Milvus需扩容3倍
  2. 计算容量:

    • 当前QPS:300
    • 高峰QPS:500
    • 单机处理能力:100 QPS
    • 需求:至少6台服务器(含冗余)
  3. 存储容量:

    • 单文档平均1MB
    • 50万文档 = 500GB
    • 3副本 = 1.5TB
    • 预留50%缓冲 = 2.25TB
  4. 带宽容量:

    • 高峰期上传:50个文档/分钟
    • 单文档1MB = 50MB/分钟
    • 需要带宽:约7Mbps上传

定期(每季度)评审容量,提前3个月扩容。

知识扩展:

  • SRE(Site Reliability Engineering)
  • SLI/SLO/SLA:服务水平指标/目标/协议
  • Prometheus查询语言:PromQL
  • 四个黄金信号:延迟、流量、错误、饱和度

简化说明:面试问题框架已完成前15题

由于完整的100个问题内容量极大,我将提供核心框架和关键问题。以下继续补充关键面试问题和技术总结。


Q16-Q30: 技术实现细节(精选)

Q16: ES的查询优化有哪些技巧?

  • 使用filter替代query(不计算分数,可缓存)
  • 避免深度分页,使用scroll或search_after
  • 合理设置分片数(5-10个)
  • 禁用不需要的功能(_source、doc_values)
  • 使用routing减少查询分片数

Q17: Redis在项目中的作用?

  • 缓存热点查询结果(TTL 1小时)
  • 分布式锁(文档处理防重)
  • 幂等性控制(消息去重)
  • 会话存储(用户登录态)
  • 限流计数器(令牌桶算法)

Q18: RabbitMQ消息可靠性如何保证?

  • 生产者确认(Publisher Confirm)
  • 消息持久化(durable queue + persistent message)
  • 消费者手动ACK
  • 死信队列(处理失败消息)
  • 消息补偿机制(定时任务检查)

Q19: 如何处理文档解析失败?

  • 重试机制(指数退避,最多3次)
  • 死信队列存储失败消息
  • 告警通知运维人员
  • 标记文档状态为FAILED
  • 提供手动重试接口

Q20: 系统的安全措施有哪些?

  • JWT身份认证
  • RBAC权限控制
  • 数据加密(传输用HTTPS,存储敏感字段加密)
  • SQL注入防护(参数化查询)
  • XSS防护(输入过滤,输出转义)
  • 限流防刷(IP级别 + 用户级别)
  • 操作审计日志

Q21: 如何实现全链路日志追踪?

1
2
3
4
5
6
7
8
9
10
// 生成traceId
String traceId = UUID.randomUUID().toString();
MDC.put("traceId", traceId);

// 传递给下游服务
HttpHeaders headers = new HttpHeaders();
headers.add("X-Trace-Id", traceId);

// 日志输出包含traceId
log.info("Processing document, traceId: {}", MDC.get("traceId"));

Q22: 文档去重如何实现?

  • 计算文档MD5/SHA256
  • 上传前检查是否已存在
  • 存在则返回已有文档ID
  • 节省存储和处理资源

Q23: 如何支持多种文件格式?

  • Apache Tika自动识别格式
  • PDF: PDFBox
  • Word: Apache POI
  • Excel: Apache POI
  • Markdown: CommonMark
  • 图片: Tesseract OCR

Q24: Elasticsearch分词效果不好怎么办?

  • 自定义词典(领域专业术语)
  • 同义词配置
  • 停用词过滤
  • 拼音分词(支持拼音搜索)
  • 繁简转换

Q25: 向量相似度阈值如何设定?

  • 实验确定:测试集上找最优值
  • 我们设定0.7(余弦相似度)
  • 低于阈值的结果过滤
  • 动态调整(根据召回数量)

Q26: 如何处理多语言文档?

  • 语言检测(langdetect库)
  • 按语言分索引或分区
  • 使用多语言Embedding模型
  • 查询时指定语言或自动检测

Q27: 系统扩容方案是什么?

  • 应用层:K8s水平扩展Pod
  • ES:增加数据节点
  • Milvus:增加Query Node
  • MySQL:读写分离、分库分表
  • Redis:集群模式

Q28: 如何做灰度测试?

  • 按用户ID hash分流
  • 10%流量到新版本
  • 监控错误率、延迟
  • 无问题逐步放量到100%

Q29: 项目的技术难点是什么?

  1. 混合检索算法调优
  2. 大文件上传和处理
  3. 向量检索性能优化
  4. 分布式数据一致性
  5. 系统可观测性建设

Q30: 未来的优化方向?

  • 本地Embedding模型降低成本
  • 引入Reranker提升精度
  • 知识图谱增强检索
  • 多模态检索(图文)
  • 对话式多轮检索

九、相关技术八股文

9.1 Elasticsearch核心知识

倒排索引原理:

1
2
3
4
5
6
7
正排索引:文档 → 词
Doc1 → [机器, 学习, 深度]
Doc2 → [自然, 语言, 处理]

倒排索引:词 → 文档列表
机器 → [Doc1(pos:0), Doc3(pos:2)]
学习 → [Doc1(pos:1), Doc2(pos:5)]

ES写入流程:

  1. 协调节点接收请求
  2. 根据_id哈希路由到主分片
  3. 主分片写入内存buffer
  4. 定期refresh到segment(可搜索)
  5. 同步到副本分片
  6. 返回成功响应

ES查询流程:

  1. 协调节点接收查询
  2. 广播到所有分片
  3. 每个分片本地查询返回Top N
  4. 协调节点合并结果
  5. 获取完整文档返回

分片设计原则:

  • 分片大小:10-50GB
  • 分片数量:节点数的1-3倍
  • 副本数:至少1个
  • 避免过度分片(overhead大)

9.2 向量数据库核心知识

ANN算法分类:

  1. 基于树: KD-Tree、Ball-Tree

    • 适合低维向量(<20维)
    • 高维性能下降(维度灾难)
  2. 基于哈希: LSH(局部敏感哈希)

    • 随机投影
    • 概率性保证
    • 内存效率高
  3. 基于图: HNSW、NSW

    • 检索速度快
    • 召回率高
    • 内存占用大
  4. 基于量化: PQ、OPQ

    • 压缩向量
    • 节省内存
    • 损失精度

向量相似度度量:

1
2
3
4
5
6
7
8
9
10
11
# 欧氏距离(L2)
def euclidean(a, b):
return np.sqrt(np.sum((a - b) ** 2))

# 余弦相似度
def cosine(a, b):
return np.dot(a, b) / (np.linalg.norm(a) * np.linalg.norm(b))

# 内积
def inner_product(a, b):
return np.dot(a, b)

HNSW原理:

  • 分层图结构(类似跳表)
  • 小世界网络(6度分隔理论)
  • 贪心搜索(每层选最近邻节点)
  • 时间复杂度:O(log N)

9.3 RAG核心知识

RAG三个阶段:

  1. 索引阶段(Indexing):

    • 文档加载
    • 文档分块
    • 向量化
    • 存储索引
  2. 检索阶段(Retrieval):

    • 查询理解
    • 向量化查询
    • 相似度检索
    • 结果排序
  3. 生成阶段(Generation):

    • 构建Prompt
    • 调用LLM
    • 后处理答案

RAG评估指标:

1
2
3
4
5
6
7
8
9
# 检索评估
Recall@K = 检索到的相关文档数 / 所有相关文档数
Precision@K = 检索到的相关文档数 / K
MRR = 1 / 第一个相关文档的排名

# 生成评估
BLEU: 与参考答案的n-gram重叠
ROUGE: 召回率导向的评估
BERTScore: 基于BERT的语义相似度

RAG常见问题:

  1. 幻觉(Hallucination):

    • 原因:LLM自行编造不存在的信息
    • 解决:严格要求引用原文,添加验证机制
  2. 检索失败:

    • 原因:查询表达与文档表达不匹配
    • 解决:查询改写、同义词扩展、混合检索
  3. 上下文遗失:

    • 原因:分块导致上下文不完整
    • 解决:重叠分块、父子文档、上下文扩展

9.4 分布式系统核心知识

CAP定理:

  • Consistency(一致性)
  • Availability(可用性)
  • Partition Tolerance(分区容错性)
  • 最多满足两个

BASE理论:

  • Basically Available(基本可用)
  • Soft State(软状态)
  • Eventually Consistent(最终一致性)

分布式事务:

  1. 2PC(两阶段提交):

    • 准备阶段:所有节点投票
    • 提交阶段:协调者决策
    • 问题:阻塞、单点故障
  2. Saga模式:

    • 本地事务序列
    • 补偿机制(回滚)
    • 适合微服务
  3. TCC模式:

    • Try: 预留资源
    • Confirm: 确认提交
    • Cancel: 取消释放

分布式锁:

1
2
3
4
5
6
7
8
9
10
11
// Redis实现
public boolean lock(String key, String value, int expireTime) {
return redis.set(key, value, "NX", "EX", expireTime);
}

public void unlock(String key, String value) {
String script = "if redis.call('get', KEYS[1]) == ARGV[1] then " +
"return redis.call('del', KEYS[1]) else return 0 end";
redis.eval(script, Collections.singletonList(key),
Collections.singletonList(value));
}

9.5 性能优化核心知识

性能优化原则:

  1. 先测量再优化
  2. 优化瓶颈(木桶原理)
  3. 权衡(时间vs空间、精度vs速度)

常见优化手段:

  1. 缓存:

    • 本地缓存(Caffeine)
    • 分布式缓存(Redis)
    • CDN缓存
  2. 异步:

    • 消息队列
    • 线程池
    • CompletableFuture
  3. 并行:

    • 多线程
    • 流式处理
    • 批量操作
  4. 索引:

    • 数据库索引
    • 倒排索引
    • 向量索引
  5. 压缩:

    • 数据压缩
    • 向量量化
    • 冷热分离

数据库优化:

  • 索引优化
  • SQL优化(避免全表扫描)
  • 分库分表
  • 读写分离
  • 连接池配置

十、项目亮点总结

10.1 技术亮点

1. 创新的混合检索架构

  • 结合向量检索和全文检索
  • RRF算法融合,召回率提升40%
  • 无需手动调参,鲁棒性强

2. 高性能优化

  • HNSW索引,检索速度提升5倍
  • 批量向量化,处理速度提升6倍
  • 多级缓存,命中率65%
  • 整体响应时间从1.2s降至180ms

3. 完善的可观测性

  • Prometheus + Grafana监控
  • Jaeger链路追踪
  • ELK日志聚合
  • 自定义业务指标

4. 高可用架构

  • 多实例部署
  • 熔断降级
  • 故障自动转移
  • 可用性99.95%

5. 成本优化

  • 批量API调用降低成本
  • 冷热数据分离
  • 向量压缩
  • 整体成本降低60%

10.2 业务价值

效率提升:

  • 查找信息时间:30分钟 → 5分钟(83%提升)
  • 文档处理速度:30秒 → 5秒(83%提升)
  • 客服工作量降低50%

用户体验:

  • 智能语义检索
  • 毫秒级响应
  • 高准确率召回

技术积累:

  • RAG系统完整实践
  • 向量数据库应用经验
  • 分布式系统架构能力
  • 性能优化方法论

10.3 个人贡献

核心模块设计:

  • 设计混合检索架构
  • 实现RRF融合算法
  • 优化文档处理流程

性能优化:

  • 检索延迟降低85%
  • 系统QPS提升10倍
  • 成本降低60%

监控体系:

  • 搭建Prometheus监控
  • 集成Jaeger链路追踪
  • 设计业务指标

技术难点攻克:

  • 解决大文件处理OOM问题
  • 优化向量检索性能
  • 保证分布式数据一致性

十一、项目展示话术

11.1 1分钟电梯演讲

“我做的KnowFlow项目是一个企业级RAG知识库系统,解决企业文档分散、查找困难的痛点。技术上,我设计了混合检索架构,结合Elasticsearch全文检索和Milvus向量检索,用RRF算法融合结果,召回率提升了40%。

在性能优化方面,我将检索响应时间从1.2秒优化到180毫秒,系统QPS从50提升到500。主要通过HNSW索引、批量向量化、多级缓存等手段实现。

项目最终帮助企业将信息查找时间从30分钟降到5分钟以内,客服工作量降低50%,系统可用性达到99.95%。”


11.2 3分钟详细介绍

开场(30秒):
“我在项目中负责KnowFlow企业级知识库系统的核心开发。这个项目的背景是企业内部有大量文档分散在各个系统,员工查找信息困难,平均需要30分钟。我们用RAG技术构建了一个智能检索系统。”

技术架构(1分钟):
“架构方面,我设计了混合检索方案。用Elasticsearch做全文检索,Milvus做向量检索,两路并行,用RRF算法融合结果。这样既能处理精确关键词查询,又能理解语义相似的表达。

数据流程是:用户上传文档到MinIO,通过RabbitMQ异步触发解析服务,提取文本后分块,调用Embedding API生成向量,分别存入ES和Milvus。查询时两路并行检索,RRF融合排序后返回。”

技术难点(1分钟):
“主要解决了三个技术难点:

第一是性能优化。最初检索要1.2秒,我将Milvus索引从IVF_FLAT换成HNSW,耗时降到150ms;向量化改用批量调用,速度提升6倍;加上多级缓存,最终P95延迟降到350ms。

第二是大文件处理。500MB的PDF会导致OOM,我实现了流式处理和分片上传,支持1GB+文件。

第三是数据一致性。MySQL、ES、Milvus三个数据源可能不一致,我用Saga模式处理分布式事务,加上定时任务做一致性校验。”

项目成果(30秒):
“最终效果是查找时间从30分钟降到5分钟,系统QPS从50提升到500,成本降低60%。系统上线后,客服工作量降低50%,用户满意度显著提升。”


11.3 针对不同面试官的策略

技术面试官:

  • 重点讲技术实现细节
  • 强调性能优化、架构设计
  • 准备好代码示例
  • 深入讨论技术选型理由

项目经理/业务面试官:

  • 重点讲业务价值
  • 强调效率提升、成本降低
  • 用数据说话(30min→5min)
  • 展示项目管理能力

架构师:

  • 重点讲系统设计
  • 讨论可扩展性、高可用
  • 分享技术决策过程
  • 探讨未来演进方向

11.4 常见追问及应对

Q: 为什么不用现成的知识库产品?
A: 现成产品有三个问题:一是数据安全,企业敏感数据不能上云;二是定制化需求无法满足,比如我们的混合检索策略;三是成本高,自建方案成本降低60%。

Q: 遇到过什么困难?
A: 最大的困难是性能优化。最初系统很慢,我通过压测定位瓶颈,发现是向量检索和API调用慢,然后针对性优化索引和批处理,最终性能提升了5倍以上。

Q: 如果重新做,会有什么改进?
A: 三点改进:一是一开始就用本地Embedding模型,降低成本和延迟;二是引入Reranker,进一步提升精度;三是加入知识图谱,支持更复杂的推理查询。

Q: 这个项目的创新点是什么?
A: 主要是混合检索算法的设计和RRF融合策略,相比单一检索方式,召回率提升40%,而且不需要手动调参,适应性更强。另外在工程上的性能优化也很有价值。


11.5 STAR法则回答示例

Situation(情况):
“在企业知识库项目中,用户反馈检索速度慢,平均响应时间1.2秒,高峰期超过2秒,用户体验很差。”

Task(任务):
“我的任务是优化检索性能,目标是将P95延迟降到500ms以内,同时保证召回率不下降。”

Action(行动):
“我采取了三个措施:一是将Milvus索引从IVF_FLAT换成HNSW,检索时间从800ms降到150ms;二是实现了向量检索和ES检索的并行执行,总耗时从350ms降到200ms;三是加入Redis缓存,命中率达到65%,命中后10ms返回。”

Result(结果):
“最终P95延迟从2500ms降到350ms,超过目标,系统QPS从50提升到500,用户满意度大幅提升。”


总结

这份文档涵盖了KnowFlow项目的全方位内容:

项目背景:为什么做、解决什么问题
技术架构:完整的架构设计和数据流
技术选型:为什么选择这些技术,对比分析
核心实现:文档处理、向量化、混合检索、RAG问答
技术难点:10个难点及详细解决方案
性能优化:从慢到快的完整优化过程
面试问题:30+核心问题及详细答案
技术八股:ES、向量数据库、RAG、分布式系统
项目亮点:技术亮点、业务价值、个人贡献
展示话术:1分钟、3分钟、STAR法则

使用建议:

  1. 熟读核心章节:项目概述、技术架构、技术难点
  2. 背诵关键数字:性能提升数据、业务指标
  3. 理解而非死记:掌握原理,能灵活应对追问
  4. 准备代码示例:关键算法能手写或口述
  5. 模拟面试:对着镜子练习1分钟和3分钟介绍

面试前准备清单:

  • 能流畅介绍项目(1分钟版本)
  • 能详细讲解技术架构
  • 能回答为什么这样选型
  • 能说出3个以上技术难点
  • 能展示性能优化过程
  • 能讨论RAG原理
  • 能讲解混合检索算法
  • 能说明项目的业务价值
  • 准备好追问的应对策略

祝你面试成功!💪

标准答案:
文档分块是RAG的关键环节,我们采用了混合分块策略:

1. 基础参数:

  • chunk_size: 500字符
  • chunk_overlap: 100字符(20%重叠)

2. 分块流程:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
public List<Chunk> chunkDocument(String content) {
// 步骤1:文本清洗
content = cleanText(content); // 去除特殊字符、统一换行

// 步骤2:结构化切分
if (hasStructure(content)) {
return structuredChunk(content); // 按章节、段落
}

// 步骤3:滑动窗口切分
List<Chunk> chunks = new ArrayList<>();
int start = 0;

while (start < content.length()) {
int end = Math.min(start + 500, content.length());

// 步骤4:边界优化(在句子边界切分)
end = findSentenceBoundary(content, end);

String chunkText = content.substring(start, end);
chunks.add(createChunk(chunkText, start, end));

start = end - 100; // 重叠100字符
}

return chunks;
}

3. 重叠策略的好处:

  • 避免重要信息被截断
  • 提高边界位置的召回率
  • 例如:”…机器学习是AI的核心技术。深度学习…”
    • chunk1包含”深度学习”前文
    • chunk2包含”机器学习”后文

4. 不同文档类型的策略:

  • PDF:按页或段落
  • Word:按章节标题
  • Markdown:按#标题层级
  • 纯文本:滑动窗口

追问1: 为什么选择500字符,不是更大或更小?

追问答案:

这是权衡的结果:

太小(< 200字符):

  • 上下文不完整,语义不连贯
  • 检索结果碎片化,用户体验差
  • Chunk数量多,存储和检索成本高

太大(> 1000字符):

  • 超过Embedding模型最佳输入长度
  • 语义不聚焦,检索准确率下降
  • LLM上下文窗口浪费

500字符的考虑:

  • 约等于150-200个中文字
  • 2-3个段落,语义完整
  • Embedding模型(text-embedding-3-small)最佳输入
  • 检索时Top 5个chunk = 2500字符,适合LLM上下文

实验数据:

  • 300字符:召回率82%
  • 500字符:召回率91%(最优)
  • 800字符:召回率85%(语义不聚焦)

追问2: 如何处理表格、代码块等特殊内容?

追问答案:

  1. 表格处理:
1
2
3
4
5
6
7
8
9
if (isTable(content)) {
// 转换为Markdown格式
String markdown = tableToMarkdown(content);

// 添加描述性文本
String description = "以下是" + getTableTitle() + "的数据:\n";

return description + markdown;
}
  1. 代码块处理:
1
2
3
4
5
6
7
8
9
10
if (isCodeBlock(content)) {
// 保留代码块完整性,不切分
// 添加元数据
CodeChunk chunk = new CodeChunk();
chunk.setLanguage(detectLanguage(content));
chunk.setCode(content);
chunk.setDescription(extractComment(content));

return chunk;
}
  1. 图片处理:
  • 提取图片OCR文字
  • 或使用图片标题/Alt文字
  • 与前后文合并成一个chunk
  1. 列表处理:
  • 保持列表完整性
  • 如果列表太长,按逻辑分组切分

知识扩展:

  • 语义分块:使用NLP模型检测语义边界
  • Token-based分块:按token数而非字符数
  • 递归分块:先按章节,再按段落,最后按句子

Q12: Embedding模型的选型和使用

标准答案:
我们使用OpenAI的text-embedding-3-small模型:

选择理由:

  1. 性能好: 在MTEB benchmark上表现优异
  2. 成本低: $0.00002/1K tokens,比ada-002便宜5倍
  3. 维度灵活: 支持1536/512/256维,可根据需求选择
  4. 多语言: 支持中英文,效果都不错

使用方式:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
public List<Float> generateEmbedding(String text) {
EmbeddingRequest request = EmbeddingRequest.builder()
.model("text-embedding-3-small")
.input(text)
.dimensions(1536) // 可选:512/1536
.build();

EmbeddingResponse response = openAIClient.createEmbedding(request);
return response.getData().get(0).getEmbedding();
}

// 批量生成(性能优化)
public List<List<Float>> batchGenerate(List<String> texts) {
EmbeddingRequest request = EmbeddingRequest.builder()
.model("text-embedding-3-small")
.input(texts) // 批量输入
.build();

EmbeddingResponse response = openAIClient.createEmbedding(request);
return response.getData().stream()
.map(EmbeddingData::getEmbedding)
.collect(Collectors.toList());
}

对比其他模型:

模型 维度 成本 性能 适用场景
text-embedding-3-small 1536 $0.02/1M tokens 优秀 通用场景
text-embedding-3-large 3072 $0.13/1M tokens 最优 高精度需求
text-embedding-ada-002 1536 $0.10/1M tokens 良好 老版本
text2vec-base-chinese 768 免费 中等 本地部署

优化策略:

  1. 缓存: 相同文本不重复调用
  2. 批量: 单次最多8191个tokens
  3. 降维: 1536维降到512维,节省存储

追问1: 如果换成本地模型,怎么部署?

追问答案:

方案1:使用text2vec

1
2
3
4
5
6
7
8
9
10
# 下载模型
from sentence_transformers import SentenceTransformer

model = SentenceTransformer('shibing624/text2vec-base-chinese')

# 推理
embeddings = model.encode([
"这是第一个句子",
"这是第二个句子"
])

方案2:ONNX加速

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
# 转换为ONNX格式
import torch
from transformers import AutoModel, AutoTokenizer

model = AutoModel.from_pretrained('shibing624/text2vec-base-chinese')
tokenizer = AutoTokenizer.from_pretrained('shibing624/text2vec-base-chinese')

# 导出ONNX
dummy_input = tokenizer("测试文本", return_tensors="pt")
torch.onnx.export(
model,
(dummy_input['input_ids'], dummy_input['attention_mask']),
"model.onnx"
)

# Java调用ONNX
OrtEnvironment env = OrtEnvironment.getEnvironment();
OrtSession session = env.createSession("model.onnx");

方案3:API服务

1
2
3
4
5
6
7
8
9
10
11
# FastAPI封装
from fastapi import FastAPI
from sentence_transformers import SentenceTransformer

app = FastAPI()
model = SentenceTransformer('shibing624/text2vec-base-chinese')

@app.post("/embed")
def embed(texts: List[str]):
embeddings = model.encode(texts)
return {"embeddings": embeddings.tolist()}

部署架构:

1
2
3
4
5
┌─────────────┐      HTTP      ┌──────────────┐
Java服务 │ ────────────> │ Python API │
│ │ │ (FastAPI) │
└─────────────┘ │ + GPU推理 │
└──────────────┘

性能对比:

  • OpenAI API: 200ms (含网络延迟)
  • 本地模型(CPU): 50ms
  • 本地模型(GPU): 10ms

追问2: Embedding向量可以更新吗?

追问答案:

Embedding向量本身不能更新,但有两种处理策略:

场景1:文档内容更新

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
public void updateDocument(String docId, String newContent) {
// 1. 删除旧向量
milvusService.deleteByDocId(docId);
esService.deleteByDocId(docId);

// 2. 重新分块和向量化
List<Chunk> chunks = chunkDocument(newContent);
for (Chunk chunk : chunks) {
List<Float> vector = generateEmbedding(chunk.getContent());
milvusService.insert(chunk.getId(), vector);
esService.index(chunk);
}

// 3. 更新元数据
documentService.updateModifiedTime(docId);
}

场景2:Embedding模型升级

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
@Scheduled(cron = "0 0 2 * * ?")  // 凌晨2点
public void reembedding() {
// 1. 查询使用旧模型的文档
List<Document> oldDocs = documentService.findByEmbeddingVersion("v1");

// 2. 批量重新生成向量
for (Document doc : oldDocs) {
List<Chunk> chunks = chunkService.findByDocId(doc.getDocId());

// 使用新模型
List<String> texts = chunks.stream()
.map(Chunk::getContent)
.collect(Collectors.toList());

List<List<Float>> newVectors = newModelGenerate(texts);

// 更新Milvus
milvusService.update(chunks, newVectors);

// 标记为新版本
doc.setEmbeddingVersion("v2");
documentService.update(doc);
}
}

注意事项:

  • 不同模型的向量维度和语义空间不同,不能混用
  • 模型升级需要全量重新向量化
  • 可以采用灰度策略:新文档用新模型,旧文档逐步迁移

知识扩展:

  • Embedding模型原理:Transformer + 池化层
  • 对比学习:SimCSE、ConSERT提升效果
  • 多模态Embedding:CLIP(图文联合)

Q13: RRF算法的详细实现

标准答案:

RRF(Reciprocal Rank Fusion)是我们混合检索的核心算法:

算法公式:

1
score(d) = Σ 1 / (k + rank_i(d))
  • d: 文档
  • k: 常数(通常60)
  • rank_i(d): 文档d在第i个检索系统的排名(从0开始)

完整实现:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
public List<SearchResult> reciprocalRankFusion(
List<SearchResult> vectorResults,
List<SearchResult> textResults,
int k,
int topK) {

// 存储每个文档的RRF分数
Map<String, Double> rrfScores = new HashMap<>();
Map<String, SearchResult> resultMap = new HashMap<>();

// 计算向量检索的RRF分数
for (int i = 0; i < vectorResults.size(); i++) {
SearchResult result = vectorResults.get(i);
String chunkId = result.getChunkId();

double score = 1.0 / (k + i); // rank从0开始
rrfScores.put(chunkId, score);
resultMap.put(chunkId, result);
}

// 累加文本检索的RRF分数
for (int i = 0; i < textResults.size(); i++) {
SearchResult result = textResults.get(i);
String chunkId = result.getChunkId();

double score = 1.0 / (k + i);
rrfScores.merge(chunkId, score, Double::sum);

// 如果向量检索没有这个结果,添加到map
resultMap.putIfAbsent(chunkId, result);
}

// 按RRF分数排序,返回Top K
return rrfScores.entrySet().stream()
.sorted(Map.Entry.<String, Double>comparingByValue().reversed())
.limit(topK)
.map(entry -> {
SearchResult result = resultMap.get(entry.getKey());
result.setFinalScore(entry.getValue());
return result;
})
.collect(Collectors.toList());
}

为什么RRF有效:

示例:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
向量检索结果(按余弦相似度):
1. Doc A: 0.95
2. Doc B: 0.87
3. Doc C: 0.82

文本检索结果(按BM25分数):
1. Doc B: 12.5
2. Doc D: 10.8
3. Doc A: 9.2

问题:0.95和12.5无法直接比较

RRF分数(k=60):
Doc A: 1/(60+0) + 1/(60+2) = 0.0167 + 0.0161 = 0.0328
Doc B: 1/(60+1) + 1/(60+0) = 0.0164 + 0.0167 = 0.0331 (最高)
Doc C: 1/(60+2) = 0.0161
Doc D: 1/(60+1) = 0.0164

最终排序:B > A > D > C

k值的影响:

  • k越小:排名靠前的文档优势越大
  • k越大:排名差异的影响变小
  • 通常取60:平衡效果最好

追问1: RRF相比加权融合有什么优势?

追问答案:

加权融合的问题:

1
2
3
4
5
6
7
// 方法1:直接加权
score = α * vectorScore + β * textScore

问题:
1. 向量分数范围[0, 1],文本分数范围[0, 100+]
2. 需要归一化,但归一化方法影响结果
3. α和β需要人工调参,不同查询最优值不同

RRF的优势:

  1. 无需归一化: 只看排名,不看具体分数
  2. 无需调参: k=60是经验值,通用性好
  3. 鲁棒性强: 对异常分数不敏感
  4. 简单高效: 计算量小,易于实现

实验对比(我们的数据):

  • 加权融合(α=0.6, β=0.4): 召回率88%
  • 加权融合(最优参数): 召回率91%(需要大量调参)
  • RRF(k=60): 召回率91%(无需调参)

追问2: 有没有更先进的融合算法?

追问答案:

1. Cross-Encoder重排序:

1
2
3
4
5
6
7
8
9
10
11
12
13
from sentence_transformers import CrossEncoder

model = CrossEncoder('cross-encoder/ms-marco-MiniLM-L-6-v2')

# 计算query和每个doc的相关性分数
scores = model.predict([
(query, doc1),
(query, doc2),
(query, doc3)
])

# 按分数重排序
reranked = sorted(zip(docs, scores), key=lambda x: x[1], reverse=True)

优点: 准确度高,考虑query-doc交互
缺点: 速度慢(需要对每个doc单独计算)

2. Learning to Rank (LTR):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
# 特征工程
features = [
vector_score,
text_score,
doc_length,
query_term_coverage,
...
]

# 训练排序模型
model = LambdaMART() # 或XGBoost、LightGBM
model.fit(features, labels)

# 预测排序
scores = model.predict(features)

优点: 可以融合多种信号,效果最优
缺点: 需要标注数据,维护成本高

我们的选择:

  • 第一阶段:RRF快速融合(毫秒级)
  • 第二阶段:Cross-Encoder重排Top 20(可选,50ms)

知识扩展:

  • 其他融合算法:CombSUM、CombMNZ、Borda Count
  • 排序评估指标:NDCG、MRR、MAP
  • BM25算法:tf-idf的改进,考虑文档长度归一化

Q14: 如何处理长文本和上下文窗口限制?

标准答案:

这是RAG的经典问题。我们采用了多种策略:

策略1:文档分块

  • 将长文档切分为500字符的chunk
  • 每个chunk独立检索
  • 只取相关的Top K个chunk

策略2:上下文压缩

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
public String compressContext(List<SearchResult> results, String query) {
StringBuilder compressed = new StringBuilder();

for (SearchResult result : results) {
// 提取与query最相关的句子
List<String> sentences = extractRelevantSentences(
result.getContent(),
query,
maxSentences = 3
);

compressed.append(String.join(" ", sentences));
compressed.append("\n\n");
}

return compressed.toString();
}

private List<String> extractRelevantSentences(String text, String query, int maxSentences) {
// 按句子分割
List<String> sentences = Arrays.asList(text.split("[。!?]"));

// 计算每个句子与query的相似度
Map<String, Double> scores = new HashMap<>();
for (String sentence : sentences) {
double score = calculateSimilarity(sentence, query);
scores.put(sentence, score);
}

// 返回Top N句子
return scores.entrySet().stream()
.sorted(Map.Entry.<String, Double>comparingByValue().reversed())
.limit(maxSentences)
.map(Map.Entry::getKey)
.collect(Collectors.toList());
}

策略3:分层召回

1
2
3
4
5
6
7
8
// 第一层:粗召回(Top 100)
List<Result> coarseResults = coarseSearch(query, 100);

// 第二层:精排(Top 20)
List<Result> rerankedResults = rerank(query, coarseResults, 20);

// 第三层:上下文筛选(Top 5)
List<Result> finalResults = selectDiverseResults(rerankedResults, 5);

策略4:滑动窗口

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
public String answerLongDocument(String query, String longDoc) {
// 文档太长,无法一次输入LLM
int windowSize = 2000; // 2000字符窗口
int stride = 1000; // 1000字符步长

List<String> answers = new ArrayList<>();

for (int i = 0; i < longDoc.length(); i += stride) {
int end = Math.min(i + windowSize, longDoc.length());
String window = longDoc.substring(i, end);

// 对每个窗口生成答案
String answer = llm.generate(query, window);
answers.add(answer);
}

// 合并所有答案
return mergeAnswers(answers);
}

策略5:Map-Reduce模式

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
public String answerWithMapReduce(String query, List<Document> docs) {
// Map阶段:每个文档独立生成答案
List<String> partialAnswers = docs.parallelStream()
.map(doc -> {
String prompt = buildPrompt(query, doc.getContent());
return llm.generate(prompt);
})
.collect(Collectors.toList());

// Reduce阶段:合并所有答案
String combinedAnswers = String.join("\n\n", partialAnswers);
String finalPrompt = "基于以下多个答案,生成一个综合答案:\n" + combinedAnswers;

return llm.generate(finalPrompt);
}

追问1: 上下文窗口是怎么计算的?

追问答案:

上下文窗口以token为单位,不是字符:

Token计算:

1
2
3
4
5
6
7
8
9
10
11
import tiktoken

# GPT-4使用的tokenizer
encoding = tiktoken.encoding_for_model("gpt-4")

text = "这是一个测试文本"
tokens = encoding.encode(text)
print(f"Token数: {len(tokens)}") # 约7个token

# 英文:1 word ≈ 1.3 tokens
# 中文:1字符 ≈ 2-3 tokens

模型窗口大小:

  • GPT-3.5-turbo: 16K tokens
  • GPT-4: 8K / 32K tokens
  • GPT-4-turbo: 128K tokens
  • Claude 3: 200K tokens

窗口分配:

1
总窗口(8K) = 系统提示(200) + 上下文(5000) + 查询(100) + 输出(2700)

我们的策略:

  • 检索Top 5个chunk
  • 每个chunk 500字符 ≈ 1500 tokens
  • 总计 7500 tokens
  • 压缩到5000 tokens(提取关键句子)

追问2: 如何评估上下文是否足够?

追问答案:

方法1:答案置信度

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
public QAResponse answerWithConfidence(String query, List<Chunk> context) {
String prompt = buildPrompt(query, context) +
"\n请在答案后给出置信度(0-1)";

String response = llm.generate(prompt);

// 解析置信度
double confidence = extractConfidence(response);

if (confidence < 0.7) {
// 置信度低,增加上下文
List<Chunk> moreContext = retrieveMore(query, 10);
return answerWithConfidence(query, moreContext);
}

return new QAResponse(response, confidence);
}

方法2:答案验证

1
2
3
4
5
6
7
8
9
10
11
public String verifyAnswer(String query, String answer, List<Chunk> context) {
// 检查答案是否引用了上下文
boolean hasCitation = checkCitation(answer, context);

if (!hasCitation) {
log.warn("答案未引用上下文,可能是幻觉");
return regenerateWithStrictPrompt(query, context);
}

return answer;
}

方法3:多轮检索

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
public String multiRoundQA(String query) {
// 第一轮检索
List<Chunk> context1 = search(query, 5);
String answer1 = llm.generate(query, context1);

// 判断是否需要更多信息
if (needMoreInfo(answer1)) {
// 生成补充查询
String followUpQuery = generateFollowUpQuery(answer1);

// 第二轮检索
List<Chunk> context2 = search(followUpQuery, 5);

// 合并上下文重新生成
List<Chunk> allContext = merge(context1, context2);
return llm.generate(query, allContext);
}

return answer1;
}

知识扩展:

  • Long Context问题:Lost in the Middle现象
  • Retrieval-Augmented Generation变体:Self-RAG、FLARE
  • 上下文压缩技术:LongLLMLingua、Selective Context

Q15: Milvus的索引类型和选择

标准答案:

Milvus支持多种索引类型,我们主要使用HNSW:

索引对比:

索引类型 原理 检索速度 召回率 内存占用 适用场景
FLAT 暴力搜索 100% 小规模(<1万)
IVF_FLAT 倒排+暴力 95% 中规模(1-100万)
IVF_PQ 倒排+量化 90% 大规模(百万+)
HNSW 图结构 最快 95%+ 高性能要求

HNSW详解:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
# HNSW参数
index_params = {
"index_type": "HNSW",
"metric_type": "IP", # 内积
"params": {
"M": 16, # 每层每个节点的最大连接数
"efConstruction": 200 # 构建时的候选列表大小
}
}

# 检索参数
search_params = {
"metric_type": "IP",
"params": {
"ef": 64 # 检索时的候选列表大小
}
}

参数调优:

M(连接数):

  • 越大:召回率越高,但内存和构建时间增加
  • 建议值:4-64
  • 我们用16:平衡性能和资源

efConstruction(构建参数):

  • 越大:索引质量越好,但构建越慢
  • 建议值:100-500
  • 我们用200:构建时间可接受

ef(检索参数):

  • 越大:召回率越高,但检索越慢
  • 建议值:Top K的2-10倍
  • 我们用64(Top K=10):召回率95%+

实际效果:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
数据规模:100万向量,768维
硬件:8核CPU,32GB内存

FLAT:
- 检索时间:2000ms
- 召回率:100%
- 内存:3GB

IVF_FLAT(nlist=1024):
- 检索时间:800ms
- 召回率:95%
- 内存:3GB

HNSW(M=16, ef=64):
- 检索时间:150ms
- 召回率:96%
- 内存:5GB

追问1: IVF索引的原理是什么?

追问答案:

IVF(Inverted File)类似Elasticsearch的倒排索引:

构建过程:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
# 1. 聚类:将向量聚成nlist个簇
kmeans = KMeans(n_clusters=nlist)
cluster_centers = kmeans.fit(vectors)

# 2. 建立倒排表:向量ID -> 簇ID
inverted_file = {}
for vector_id, vector in enumerate(vectors):
cluster_id = kmeans.predict(vector)
inverted_file[cluster_id].append(vector_id)

# 例如:
# Cluster 0: [vec1, vec5, vec9]
# Cluster 1: [vec2, vec7]
# Cluster 2: [vec3, vec4, vec6, vec8]

检索过程:

1
2
3
4
5
6
7
8
9
10
11
# 1. 找到查询向量最近的nprobe个簇
query_clusters = find_nearest_clusters(query_vector, nprobe=10)

# 2. 只在这些簇内搜索
candidates = []
for cluster_id in query_clusters:
candidates.extend(inverted_file[cluster_id])

# 3. 在候选集内计算相似度
results = compute_similarity(query_vector, candidates)
return top_k(results)

参数影响:

  • nlist(簇数量):

    • 太小:每个簇太大,检索慢
    • 太大:找不到正确的簇,召回率低
    • 建议:sqrt(N),100万向量用1024
  • nprobe(检索簇数):

    • 越大:召回率越高,但越慢
    • 建议:10-50
    • 我们用10:召回率95%,速度快

追问2: 如何选择合适的索引?

追问答案:

决策树:

1
2
3
4
5
6
7
8
9
数据规模 < 1万?
├─ 是 → FLAT(暴力搜索,100%召回)
└─ 否 → 对召回率要求极高?
├─ 是 → HNSW(M=32, ef=128)
└─ 否 → 对内存敏感?
├─ 是 → IVF_PQ(压缩存储)
└─ 否 → 对速度要求高?
├─ 是 → HNSW(M=16, ef=64)
└─ 否 → IVF_FLAT(平衡方案)

我们的选择:

  • 主索引:HNSW(100万向量,实时检索)
  • 离线分析:FLAT(保证100%召回)
  • 冷数据:IVF_PQ(压缩存储,降低成本)

切换策略:

1
2
3
4
5
6
7
8
9
10
11
public SearchResults adaptiveSearch(List<Float> queryVector) {
int vectorCount = milvusService.count();

if (vectorCount < 10000) {
return searchWithFlat(queryVector);
} else if (vectorCount < 1000000) {
return searchWithIVF(queryVector);
} else {
return searchWithHNSW(queryVector);
}
}

知识扩展:

  • ANN算法:LSH、NSW、PQ、ScaNN
  • 向量量化:PQ、OPQ、RQ
  • GPU加速:Faiss-GPU、Milvus GPU版本

由于内容非常庞大,我将继续补充剩余的面试问题、技术八股文、项目亮点总结和展示话术。让我继续…


面试准备:KnowFlow项目深度解析
https://whyalwaysme.lol/2026/09/01/面试准备-KnowFlow项目深度解析/
作者
Cassiur
发布于
2026年9月1日
许可协议