Spring Boot 3.3 + 通义千问API:5步搞定企业知识库智能客服(附避坑指南)
Spring Boot 3.3 + 通义千问:构建企业级智能知识库的实战心法
最近在帮几个中小企业的技术团队落地AI客服项目,发现一个挺有意思的现象:大家一提到“智能问答”,第一反应往往是去研究各种复杂的Python框架和前沿算法,却忽略了Java生态里那些已经相当成熟的解决方案。其实,对于大多数Java背景的团队来说,用Spring Boot加上阿里云的通义千问,完全可以在几天内搭建出一个能跑起来的智能知识库系统,而且性能、稳定性都相当不错。
今天我就结合自己最近几个项目的实战经验,聊聊怎么用Spring Boot 3.3和通义千问API,快速构建一个企业级的智能知识库问答系统。我会重点分享那些文档里不会写的“坑”,以及怎么在实际项目中做出更好的选择。
1. 环境准备与框架选型:别在起跑线上浪费时间
开始之前,我们先得把环境搭好。很多人觉得这步简单,结果往往在这里卡住半天。
1.1 JDK与Spring Boot版本选择
我强烈建议直接上JDK 17和Spring Boot 3.3.x。别再用JDK 8了,虽然还能跑,但Spring AI Alibaba的一些新特性在旧版本上支持得不太好,而且性能差距明显。
<!-- pom.xml中的关键配置 -->
<properties>
<java.version>17</java.version>
<spring-boot.version>3.3.4</spring-boot.version>
<spring-ai.version>1.0.0-M6</spring-ai.version>
</properties>
注意:Spring AI目前还处于Milestone阶段,API可能会有变动。如果你要上生产环境,建议锁定具体的版本号,别用
latest这种模糊的版本声明。
1.2 依赖仓库配置
因为Spring AI Alibaba还在快速迭代中,你得在pom.xml里加上Spring的Snapshot和Milestone仓库:
<repositories>
<repository>
<id>spring-milestones</id>
<name>Spring Milestones</name>
<url>https://repo.spring.io/milestone</url>
<snapshots>
<enabled>false</enabled>
</snapshots>
</repository>
<repository>
<id>spring-snapshots</id>
<name>Spring Snapshots</name>
<url>https://repo.spring.io/snapshot</url>
<releases>
<enabled>false</enabled>
</releases>
</repository>
</repositories>
1.3 核心依赖引入
接下来是核心依赖。这里有个小技巧:除了Spring AI Alibaba的starter,我建议把spring-boot-starter-webflux也加上,因为后面做流式响应的时候会用到。
<dependencies>
<!-- Spring AI Alibaba -->
<dependency>
<groupId>com.alibaba.cloud.ai</groupId>
<artifactId>spring-ai-alibaba-starter</artifactId>
<version>1.0.0-M2</version>
</dependency>
<!-- WebFlux用于流式响应 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-webflux</artifactId>
</dependency>
<!-- 向量数据库客户端(这里以Redis为例) -->
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-redis-store</artifactId>
<version>${spring-ai.version}</version>
</dependency>
</dependencies>
1.4 通义千问API密钥配置
去阿里云的控制台开通通义千问服务,拿到API Key。配置的时候,我建议用环境变量而不是硬编码在配置文件里:
# application.yml
spring:
ai:
dashscope:
api-key: ${DASHSCOPE_API_KEY:your-default-key-here}
chat:
options:
model: qwen-turbo
temperature: 0.7
然后在启动应用前设置环境变量:
export DASHSCOPE_API_KEY=sk-your-actual-key-here
2. 向量数据库选型与配置:别被概念吓到
很多人一听到“向量数据库”就觉得特别高大上,其实对于中小规模的知识库(比如几万到几十万条文档),用Redis完全够用,而且部署简单、成本低。
2.1 Redis向量存储配置
Spring AI提供了Redis的向量存储实现,配置起来特别简单:
@Configuration
public class VectorStoreConfig {
@Bean
public VectorStore vectorStore(RedisConnectionFactory connectionFactory) {
RedisVectorStoreConfig config = RedisVectorStoreConfig.builder()
.withIndexName("company-knowledge-base")
.withPrefix("vec:")
.build();
return new RedisVectorStore(connectionFactory, config);
}
@Bean
public EmbeddingModel embeddingModel() {
// Spring AI Alibaba会自动配置DashScope的Embedding模型
// 这里只需要确保spring.ai.dashscope.api-key配置正确
return null; // 实际由Spring自动注入
}
}
2.2 文档切分策略
这是RAG系统里特别关键的一步。切得太碎,上下文信息不够;切得太大,检索精度下降。我一般用这种混合策略:
@Component
public class DocumentSplitter {
private final Tokenizer tokenizer;
public DocumentSplitter() {
// 使用通义千问的tokenizer
this.tokenizer = new DashScopeTokenizer();
}
public List<TextSegment> splitDocument(Document document) {
// 第一层:按段落切分
List<TextSegment> paragraphs = splitByParagraph(document);
// 第二层:对长段落再按句子切分
List<TextSegment> finalSegments = new ArrayList<>();
for (TextSegment paragraph : paragraphs) {
if (tokenizer.countTokens(paragraph.getText()) > 300) {
// 段落太长,按句子切分
finalSegments.addAll(splitBySentence(paragraph));
} else {
finalSegments.add(paragraph);
}
}
// 第三层:确保每个片段有重叠,避免信息断层
return addOverlap(finalSegments, 50); // 重叠50个token
}
private List<TextSegment> splitByParagraph(Document document) {
// 实现按段落切分的逻辑
return Arrays.stream(document.getText().split("\n\n"))
.filter(para -> !para.trim().isEmpty())
.map(para -> new TextSegment(para.trim()))
.collect(Collectors.toList());
}
private List<TextSegment> splitBySentence(TextSegment segment) {
// 简单的按句号、问号、感叹号切分
String[] sentences = segment.getText().split("[。?!]");
return Arrays.stream(sentences)
.filter(s -> !s.trim().isEmpty())
.map(s -> new TextSegment(s.trim()))
.collect(Collectors.toList());
}
private List<TextSegment> addOverlap(List<TextSegment> segments, int overlapTokens) {
// 为相邻片段添加重叠部分
List<TextSegment> overlapped = new ArrayList<>();
for (int i = 0; i < segments.size(); i++) {
TextSegment current = segments.get(i);
if (i > 0) {
TextSegment previous = segments.get(i - 1);
// 从前一个片段末尾取一部分作为重叠
String overlap = extractOverlap(previous.getText(), overlapTokens);
current = new TextSegment(overlap + current.getText());
}
overlapped.add(current);
}
return overlapped;
}
}
2.3 向量化参数调优
不同的Embedding模型有不同的最佳实践。通义千问的文本向量化服务有几个版本,我对比过它们的表现:
| 模型版本 | 输入长度限制 | 输出维度 | 适用场景 | 中文优化 |
|---|---|---|---|---|
| ops-text-embedding-001 | 300 tokens | 1536 | 通用多语言 | 中等 |
| ops-text-embedding-zh-001 | 1024 tokens | 768 | 纯中文场景 | 优秀 |
| ops-text-embedding-002 | 8192 tokens | 1024 | 长文本 | 良好 |
对于中文知识库,我推荐用ops-text-embedding-zh-001,虽然维度低一些,但在中文语义理解上表现更好。
3. RAG核心实现:从理论到实践的三个关键点
3.1 检索策略:不只是相似度搜索
很多人做RAG就是简单的向量相似度搜索,但实际项目中,我发现了几个可以显著提升效果的方法:
@Service
public class EnhancedRetrievalService {
private final VectorStore vectorStore;
private final EmbeddingModel embeddingModel;
private final ChatClient chatClient;
public EnhancedRetrievalService(VectorStore vectorStore,
EmbeddingModel embeddingModel,
ChatClient chatClient) {
this.vectorStore = vectorStore;
this.embeddingModel = embeddingModel;
this.chatClient = chatClient;
}
public List<TextSegment> retrieveRelevantDocuments(String query) {
// 1. 基础向量检索
List<TextSegment> vectorResults = vectorStore.similaritySearch(
SearchRequest.query(query)
.withTopK(10)
.withSimilarityThreshold(0.7)
);
// 2. 查询扩展:让大模型帮我们改写查询
String expandedQuery = expandQuery(query);
List<TextSegment> expandedResults = vectorStore.similaritySearch(
SearchRequest.query(expandedQuery)
.withTopK(5)
);
// 3. 混合检索:结合关键词匹配
List<TextSegment> keywordResults = keywordSearch(query);
// 4. 结果去重和重排序
return rerankAndDeduplicate(query,
Stream.of(vectorResults, expandedResults, keywordResults)
.flatMap(List::stream)
.collect(Collectors.toList())
);
}
private String expandQuery(String originalQuery) {
// 用大模型生成相关的查询变体
String prompt = """
用户的问题是:%s
请生成3个语义相同但表达不同的查询,用于在知识库中检索相关信息。
只输出查询语句,用换行分隔。
""".formatted(originalQuery);
String response = chatClient.prompt()
.user(prompt)
.call()
.content();
// 取第一个扩展查询
return response.split("\n")[0];
}
private List<TextSegment> keywordSearch(String query) {
// 简单的关键词匹配(实际项目中可以用Elasticsearch等)
// 这里简化实现
return Collections.emptyList();
}
private List<TextSegment> rerankAndDeduplicate(String query,
List<TextSegment> candidates) {
// 基于多种策略重排序
return candidates.stream()
.distinct()
.sorted((a, b) -> {
// 综合评分:相似度 + 长度惩罚 + 位置权重
double scoreA = calculateScore(query, a);
double scoreB = calculateScore(query, b);
return Double.compare(scoreB, scoreA);
})
.limit(5) // 最终返回top 5
.collect(Collectors.toList());
}
}
3.2 提示工程:让大模型“听话”的关键
通义千问的能力很强,但如果你不给它明确的指令,它可能会自由发挥。下面是我在实际项目中总结出来的几个有效的提示模板:
@Component
public class PromptTemplateManager {
private static final String KNOWLEDGE_QA_TEMPLATE = """
你是一个专业的客服助手,基于以下提供的公司知识库信息回答问题。
知识库信息:
%s
用户问题:%s
请严格按照以下要求回答:
1. 答案必须基于上述知识库信息,不要添加知识库中没有的内容
2. 如果知识库信息不足以回答问题,请明确告知“根据现有资料无法回答此问题”
3. 答案要简洁明了,分点说明(如果适用)
4. 不要使用“根据知识库”、“根据提供的信息”等冗余表述
5. 如果问题涉及具体操作步骤,请按顺序列出
现在请回答问题:
""";
private static final String SUMMARIZATION_TEMPLATE = """
请将以下文档内容总结为3-5个关键要点:
文档内容:
%s
要求:
1. 每个要点不超过2句话
2. 使用中文
3. 避免专业术语,用通俗语言表达
4. 按重要性排序
关键要点:
""";
public String buildKnowledgeQAPrompt(String context, String question) {
return String.format(KNOWLEDGE_QA_TEMPLATE, context, question);
}
public String buildSummarizationPrompt(String content) {
return String.format(SUMMARIZATION_TEMPLATE, content);
}
// 更多专业领域的模板...
public String buildTechnicalSupportPrompt(String context, String question) {
return """
你是一名技术专家,正在帮助用户解决技术问题。
相关技术文档:
%s
用户的问题:%s
请按照以下结构回答:
1. 问题诊断(分析可能的原因)
2. 解决方案(分步骤说明)
3. 预防措施(如何避免再次出现)
4. 相关参考(如果有)
如果文档中没有相关信息,请说“这个问题需要进一步的技术支持”。
""".formatted(context, question);
}
}
3.3 流式响应实现:提升用户体验
对于长回答,流式响应能让用户感觉响应更快。Spring AI配合WebFlux实现起来很优雅:
@RestController
@RequestMapping("/api/chat")
public class ChatController {
private final ChatService chatService;
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> streamChat(@RequestParam String question,
@RequestParam(required = false) String conversationId) {
return chatService.streamAnswer(question, conversationId)
.map(chunk -> {
// 处理SSE格式
return "data: " + chunk + "\n\n";
})
.onErrorResume(e -> {
return Flux.just("data: [错误] " + e.getMessage() + "\n\n");
});
}
}
@Service
public class ChatService {
public Flux<String> streamAnswer(String question, String conversationId) {
// 1. 检索相关文档
List<TextSegment> relevantDocs = retrievalService.retrieveRelevantDocuments(question);
// 2. 构建上下文
String context = buildContextFromSegments(relevantDocs);
// 3. 构建提示
String prompt = promptTemplateManager.buildKnowledgeQAPrompt(context, question);
// 4. 流式调用大模型
return chatClient.prompt()
.user(prompt)
.stream()
.map(AiResponse::getContent)
.filter(content -> content != null && !content.trim().isEmpty());
}
private String buildContextFromSegments(List<TextSegment> segments) {
// 合并片段,添加来源标记
StringBuilder context = new StringBuilder();
for (int i = 0; i < segments.size(); i++) {
context.append("[文档片段 ").append(i + 1).append("]\n");
context.append(segments.get(i).getText());
context.append("\n\n");
}
return context.toString();
}
}
前端用EventSource接收流式响应:
// 前端示例代码
function setupStreamingChat() {
const eventSource = new EventSource('/api/chat/stream?question=' +
encodeURIComponent(userQuestion));
const messageContainer = document.getElementById('message-container');
eventSource.onmessage = function(event) {
const chunk = event.data;
if (chunk.startsWith('data: ')) {
const content = chunk.substring(6);
messageContainer.innerHTML += content;
// 自动滚动到底部
messageContainer.scrollTop = messageContainer.scrollHeight;
}
};
eventSource.onerror = function(error) {
console.error('EventSource failed:', error);
eventSource.close();
};
// 用户停止输入或点击停止按钮时
return () => eventSource.close();
}
4. 性能优化与监控:让系统真正可用
4.1 缓存策略设计
大模型API调用不便宜,而且有速率限制。合理的缓存能显著降低成本和提高响应速度。
@Service
@Slf4j
public class CacheAwareChatService {
private final CacheManager cacheManager;
private final ChatService delegate;
// 使用Caffeine作为本地缓存
private final Cache<String, String> responseCache = Caffeine.newBuilder()
.maximumSize(1000)
.expireAfterWrite(1, TimeUnit.HOURS) // 1小时过期
.recordStats()
.build();
private final Cache<String, List<TextSegment>> embeddingCache = Caffeine.newBuilder()
.maximumSize(5000)
.expireAfterWrite(24, TimeUnit.HOURS)
.build();
public String getCachedAnswer(String question) {
String cacheKey = generateCacheKey(question);
// 先查缓存
return responseCache.get(cacheKey, key -> {
log.info("缓存未命中,调用大模型API: {}", question);
String answer = delegate.getAnswer(question);
// 异步更新相关问题的缓存
CompletableFuture.runAsync(() -> {
updateRelatedCache(question, answer);
});
return answer;
});
}
public Embedding getCachedEmbedding(String text) {
return embeddingCache.get(text, key -> {
return embeddingModel.embed(text);
});
}
private String generateCacheKey(String question) {
// 标准化问题:转小写、去除标点、排序词语
String normalized = question.toLowerCase()
.replaceAll("[^\\p{L}\\p{N}\\s]", "")
.trim();
// 按词语排序,确保语义相同的问题命中同一个缓存
String[] words = normalized.split("\\s+");
Arrays.sort(words);
return String.join(" ", words);
}
private void updateRelatedCache(String originalQuestion, String answer) {
// 生成相似问题并缓存
List<String> similarQuestions = generateSimilarQuestions(originalQuestion);
for (String similar : similarQuestions) {
String key = generateCacheKey(similar);
responseCache.put(key, answer);
}
}
}
4.2 异步处理与批量化
对于文档入库这种耗时操作,一定要用异步处理:
@Service
public class AsyncDocumentProcessor {
private final ExecutorService executor = Executors.newFixedThreadPool(
Runtime.getRuntime().availableProcessors() * 2
);
@Async
public CompletableFuture<Void> processDocumentsBatch(List<Document> documents) {
return CompletableFuture.runAsync(() -> {
// 分批处理,避免内存溢出
List<List<Document>> batches = Lists.partition(documents, 50);
batches.forEach(batch -> {
try {
// 1. 文本清洗
List<Document> cleaned = cleanDocuments(batch);
// 2. 切分
List<TextSegment> segments = splitDocuments(cleaned);
// 3. 批量向量化(减少API调用次数)
List<Embedding> embeddings = batchEmbed(segments);
// 4. 批量存储
storeVectorsBatch(segments, embeddings);
log.info("处理完成一批文档,数量:{}", batch.size());
} catch (Exception e) {
log.error("文档处理失败", e);
// 记录失败,但不中断整个流程
}
});
}, executor);
}
private List<Embedding> batchEmbed(List<TextSegment> segments) {
// 通义千问的Embedding API支持批量调用
// 这里可以一次性处理多个文本
List<String> texts = segments.stream()
.map(TextSegment::getText)
.collect(Collectors.toList());
return embeddingModel.embed(texts);
}
}
4.3 监控与告警
生产环境必须要有监控。我通常用Spring Boot Actuator加上自定义的指标:
@Component
public class ChatMetrics {
private final MeterRegistry meterRegistry;
private final DistributionSummary responseTimeSummary;
private final Counter errorCounter;
public ChatMetrics(MeterRegistry meterRegistry) {
this.meterRegistry = meterRegistry;
// 响应时间分布
this.responseTimeSummary = DistributionSummary
.builder("chat.response.time")
.description("聊天响应时间分布")
.baseUnit("milliseconds")
.publishPercentiles(0.5, 0.95, 0.99)
.register(meterRegistry);
// 错误计数
this.errorCounter = Counter
.builder("chat.errors")
.description("聊天服务错误次数")
.tag("type", "api")
.register(meterRegistry);
}
public void recordResponseTime(long milliseconds) {
responseTimeSummary.record(milliseconds);
// 同时记录到日志,方便排查慢查询
if (milliseconds > 5000) {
log.warn("慢响应检测: {}ms", milliseconds);
}
}
public void incrementError(String errorType) {
errorCounter.increment();
meterRegistry.counter("chat.errors.detail", "type", errorType).increment();
}
// API调用统计
public void recordApiCall(String model, boolean success, int tokenCount) {
Timer.Sample sample = Timer.start(meterRegistry);
// ... API调用逻辑
sample.stop(meterRegistry.timer("dashscope.api.calls",
"model", model,
"success", String.valueOf(success)));
// Token使用量
meterRegistry.counter("dashscope.tokens.used",
"model", model,
"type", "total").increment(tokenCount);
}
}
然后在application.yml中配置Actuator端点:
management:
endpoints:
web:
exposure:
include: health,metrics,prometheus
metrics:
export:
prometheus:
enabled: true
distribution:
percentiles-histogram:
http.server.requests: true
4.4 限流与降级
大模型API有调用频率限制,客户端也需要限流保护:
@Configuration
public class RateLimitConfig {
@Bean
public RateLimiter dashscopeRateLimiter() {
// 通义千问API限制:QPS根据套餐不同
return RateLimiter.create(10); // 10 QPS
}
@Bean
public Filter rateLimitFilter(RateLimiter rateLimiter) {
return new OncePerRequestFilter() {
@Override
protected void doFilterInternal(HttpServletRequest request,
HttpServletResponse response,
FilterChain filterChain)
throws ServletException, IOException {
if (request.getRequestURI().startsWith("/api/chat")) {
if (!rateLimiter.tryAcquire()) {
response.setStatus(429); // Too Many Requests
response.getWriter().write("请求过于频繁,请稍后再试");
return;
}
}
filterChain.doFilter(request, response);
}
};
}
}
@Service
public class CircuitBreakerChatService {
private final CircuitBreaker circuitBreaker;
private final ChatService delegate;
public CircuitBreakerChatService() {
this.circuitBreaker = CircuitBreaker.ofDefaults("dashscope-api");
}
@TimeLimiter(name = "chatTimeout", fallbackMethod = "timeoutFallback")
@CircuitBreaker(name = "dashscopeApi", fallbackMethod = "apiFallback")
public String getAnswerWithResilience(String question) {
return delegate.getAnswer(question);
}
public String timeoutFallback(String question, TimeoutException e) {
log.warn("大模型API响应超时,返回默认回答", e);
return "系统正在处理中,请稍后重试或联系客服人员。";
}
public String apiFallback(String question, Exception e) {
log.error("大模型API调用失败,使用备用方案", e);
// 可以返回缓存的常见问题答案
// 或者调用备用的大模型服务
return getFallbackAnswer(question);
}
}
5. 实际部署中的那些“坑”
5.1 中文编码问题
Spring Boot默认的字符编码可能不完整,特别是处理中文文档时:
@Configuration
public class EncodingConfig {
@Bean
public CharacterEncodingFilter characterEncodingFilter() {
CharacterEncodingFilter filter = new CharacterEncodingFilter();
filter.setEncoding("UTF-8");
filter.setForceEncoding(true);
return filter;
}
@Bean
public HttpMessageConverter<String> responseBodyConverter() {
StringHttpMessageConverter converter = new StringHttpMessageConverter(
StandardCharsets.UTF_8
);
converter.setWriteAcceptCharset(false);
return converter;
}
}
5.2 文档格式处理
不同格式的文档需要不同的处理方式:
@Component
public class DocumentProcessor {
public Document processDocument(MultipartFile file) throws IOException {
String filename = file.getOriginalFilename();
String content;
if (filename.endsWith(".pdf")) {
content = parsePdf(file.getInputStream());
} else if (filename.endsWith(".docx")) {
content = parseDocx(file.getInputStream());
} else if (filename.endsWith(".txt") || filename.endsWith(".md")) {
content = new String(file.getBytes(), StandardCharsets.UTF_8);
} else if (filename.endsWith(".html") || filename.endsWith(".htm")) {
content = parseHtml(file.getInputStream());
} else {
throw new IllegalArgumentException("不支持的文件格式: " + filename);
}
// 清理文本
content = cleanText(content);
return new Document(content, Map.of(
"filename", filename,
"size", String.valueOf(file.getSize()),
"uploadTime", LocalDateTime.now().toString()
));
}
private String cleanText(String text) {
// 移除多余的空格和换行
text = text.replaceAll("\\s+", " ")
.replaceAll("\\n{3,}", "\n\n")
.trim();
// 处理全角/半角字符
text = fullWidthToHalfWidth(text);
// 移除不可见字符
text = text.replaceAll("[\\x00-\\x08\\x0B\\x0C\\x0E-\\x1F\\x7F]", "");
return text;
}
private String fullWidthToHalfWidth(String text) {
// 全角转半角
char[] chars = text.toCharArray();
for (int i = 0; i < chars.length; i++) {
if (chars[i] == '\u3000') {
chars[i] = ' ';
} else if (chars[i] > '\uFF00' && chars[i] < '\uFF5F') {
chars[i] = (char) (chars[i] - 65248);
}
}
return new String(chars);
}
}
5.3 内存管理
处理大文档时容易内存溢出,需要特别注意:
@Service
public class MemorySafeDocumentProcessor {
private static final int MAX_DOCUMENT_SIZE = 10 * 1024 * 1024; // 10MB
public void processLargeDocument(Path filePath) throws IOException {
try (BufferedReader reader = Files.newBufferedReader(filePath)) {
String line;
StringBuilder currentChunk = new StringBuilder();
int chunkSize = 0;
while ((line = reader.readLine()) != null) {
// 按段落处理,避免一次性加载整个文件
if (line.trim().isEmpty()) {
// 遇到空行,处理当前块
if (currentChunk.length() > 0) {
processChunk(currentChunk.toString());
currentChunk.setLength(0);
chunkSize = 0;
}
} else {
if (chunkSize + line.length() > 5000) { // 每个块最多5000字符
processChunk(currentChunk.toString());
currentChunk.setLength(0);
chunkSize = 0;
}
currentChunk.append(line).append("\n");
chunkSize += line.length();
}
}
// 处理最后一块
if (currentChunk.length() > 0) {
processChunk(currentChunk.toString());
}
}
}
private void processChunk(String chunk) {
// 异步处理每个块
CompletableFuture.runAsync(() -> {
Document doc = new Document(chunk);
embeddingAndStore(doc);
});
}
}
5.4 错误处理与重试
网络调用总会有失败,必须有完善的错误处理:
@Service
@Slf4j
public class RetryableChatService {
private final ChatClient chatClient;
private final RetryTemplate retryTemplate;
public RetryableChatService(ChatClient chatClient) {
this.chatClient = chatClient;
this.retryTemplate = RetryTemplate.builder()
.maxAttempts(3)
.exponentialBackoff(1000, 2, 5000) // 初始1秒,指数退避
.retryOn(DashScopeApiException.class)
.retryOn(SocketTimeoutException.class)
.retryOn(IOException.class)
.withListener(new RetryListener() {
@Override
public <T, E extends Throwable> void onError(
RetryContext context, RetryCallback<T, E> callback, Throwable throwable) {
log.warn("第{}次重试失败: {}", context.getRetryCount(),
throwable.getMessage());
}
})
.build();
}
public String getAnswerWithRetry(String question) {
return retryTemplate.execute(context -> {
try {
return chatClient.prompt()
.user(question)
.call()
.content();
} catch (DashScopeApiException e) {
// 检查是否是配额不足
if (e.getErrorCode() == 429) {
log.error("API配额不足,需要升级套餐或等待重置");
throw new QuotaExceededException("API调用次数超限");
}
throw e;
}
});
}
@Retryable(value = {SocketTimeoutException.class, IOException.class},
maxAttempts = 3,
backoff = @Backoff(delay = 1000, multiplier = 2))
public Embedding getEmbeddingWithRetry(String text) {
return embeddingModel.embed(text);
}
@Recover
public Embedding recoverEmbedding(Exception e, String text) {
log.error("获取Embedding失败,使用默认值", e);
// 返回一个零向量或缓存的值
return getCachedEmbeddingOrDefault(text);
}
}
6. 进阶功能:让系统更智能
6.1 多轮对话支持
简单的问答不够,还需要支持多轮对话:
@Service
public class ConversationService {
private final ConversationMemory memory;
public String handleConversation(String sessionId, String userMessage) {
// 1. 获取对话历史
List<ChatMessage> history = memory.getConversationHistory(sessionId, 10);
// 2. 如果有历史,可以优化当前问题
String optimizedQuery = optimizeQueryWithHistory(userMessage, history);
// 3. 检索相关文档(考虑历史上下文)
List<TextSegment> relevantDocs = retrieveWithContext(optimizedQuery, history);
// 4. 构建包含历史的提示
String prompt = buildPromptWithHistory(userMessage, relevantDocs, history);
// 5. 调用大模型
String response = chatClient.prompt()
.messages(history)
.user(prompt)
.call()
.content();
// 6. 保存到对话历史
memory.saveMessage(sessionId, "user", userMessage);
memory.saveMessage(sessionId, "assistant", response);
return response;
}
private String optimizeQueryWithHistory(String currentQuery,
List<ChatMessage> history) {
if (history.isEmpty()) {
return currentQuery;
}
// 用大模型优化查询,考虑上下文
String historySummary = summarizeHistory(history);
String optimizationPrompt = """
用户之前的对话历史:
%s
用户当前的问题:%s
请根据对话历史,重新组织当前问题,使其更完整、明确。
只输出优化后的问题,不要其他内容。
""".formatted(historySummary, currentQuery);
return chatClient.prompt()
.user(optimizationPrompt)
.call()
.content();
}
}
6.2 答案验证与评分
不是所有大模型的回答都可靠,需要验证:
@Component
public class AnswerValidator {
public ValidationResult validateAnswer(String question,
String answer,
List<TextSegment> sources) {
ValidationResult result = new ValidationResult();
// 1. 检查答案是否基于提供的来源
boolean isGrounding = checkGrounding(answer, sources);
result.setGroundingScore(isGrounding ? 1.0 : 0.0);
// 2. 检查答案的完整性
double completeness = evaluateCompleteness(question, answer);
result.setCompletenessScore(completeness);
// 3. 检查是否有幻觉(编造信息)
boolean hasHallucination = detectHallucination(answer, sources);
result.setHasHallucination(hasHallucination);
// 4. 置信度评分
double confidence = calculateConfidence(answer);
result.setConfidenceScore(confidence);
// 综合评分
result.setOverallScore(
(result.getGroundingScore() * 0.4 +
result.getCompletenessScore() * 0.3 +
result.getConfidenceScore() * 0.3)
);
return result;
}
private boolean checkGrounding(String answer, List<TextSegment> sources) {
// 简单的实现:检查答案中的关键实体是否在来源中出现
Set<String> sourceEntities = extractEntities(sources);
Set<String> answerEntities = extractEntities(answer);
// 如果答案中的实体大部分都在来源中,认为是基于来源的
long matched = answerEntities.stream()
.filter(sourceEntities::contains)
.count();
return matched >= answerEntities.size() * 0.7; // 70%匹配
}
private double evaluateCompleteness(String question, String answer) {
// 使用大模型评估答案是否完整回答了问题
String evaluationPrompt = """
问题:%s
回答:%s
请评估这个回答是否完整地解决了问题(0-1分):
1. 是否直接回答了问题
2. 是否提供了足够的细节
3. 是否解决了问题的所有方面
只输出一个0到1之间的数字,不要其他内容。
""".formatted(question, answer);
String scoreStr = chatClient.prompt()
.user(evaluationPrompt)
.call()
.content();
try {
return Double.parseDouble(scoreStr.trim());
} catch (NumberFormatException e) {
return 0.5; // 默认值
}
}
}
6.3 知识库更新与版本管理
知识库不是一次性的,需要持续更新:
@Service
public class KnowledgeBaseManager {
private final VectorStore vectorStore;
private final DocumentVersionRepository versionRepo;
public void updateDocument(String docId, Document newVersion) {
// 1. 标记旧版本为失效
versionRepo.markAsObsolete(docId);
// 2. 处理新文档
List<TextSegment> segments = splitter.splitDocument(newVersion);
List<Embedding> embeddings = embeddingModel.embed(
segments.stream()
.map(TextSegment::getText)
.collect(Collectors.toList())
);
// 3. 存储新版本
for (int i = 0; i < segments.size(); i++) {
TextSegment segment = segments.get(i);
Embedding embedding = embeddings.get(i);
// 添加版本信息到元数据
Map<String, Object> metadata = new HashMap<>(segment.getMetadata());
metadata.put("docId", docId);
metadata.put("version", getNextVersion(docId));
metadata.put("updateTime", Instant.now().toString());
TextSegment versionedSegment = new TextSegment(
segment.getText(),
metadata
);
vectorStore.add(embedding, versionedSegment);
}
// 4. 记录版本历史
versionRepo.save(new DocumentVersion(docId, getCurrentVersion(docId)));
}
public void cleanupOldVersions(String docId, int keepVersions) {
// 只保留最近N个版本
List<DocumentVersion> versions = versionRepo.findByDocId(docId);
if (versions.size() > keepVersions) {
versions.stream()
.sorted(Comparator.comparing(DocumentVersion::getCreatedAt).reversed())
.skip(keepVersions)
.forEach(oldVersion -> {
// 从向量库中删除旧版本
deleteVectorsByVersion(docId, oldVersion.getVersion());
versionRepo.delete(oldVersion);
});
}
}
}
7. 测试与验证:确保系统可靠
7.1 单元测试
@SpringBootTest
class ChatServiceTest {
@Autowired
private ChatService chatService;
@MockBean
private VectorStore vectorStore;
@MockBean
private EmbeddingModel embeddingModel;
@Test
void testBasicQa() {
// 准备测试数据
TextSegment testSegment = new TextSegment(
"公司的年假政策是:员工入职满一年后享受15天年假。",
Map.of("source", "员工手册")
);
when(vectorStore.similaritySearch(any()))
.thenReturn(List.of(testSegment));
when(embeddingModel.embed(anyString()))
.thenReturn(new Embedding(new float[768]));
// 执行测试
String answer = chatService.getAnswer("年假有多少天?");
// 验证结果
assertThat(answer).contains("15天");
assertThat(answer).doesNotContain("根据知识库"); // 验证提示工程生效
}
@Test
void testUnknownQuestion() {
when(vectorStore.similaritySearch(any()))
.thenReturn(Collections.emptyList());
String answer = chatService.getAnswer("公司什么时候上市?");
assertThat(answer).contains("无法回答");
assertThat(answer).contains("现有资料");
}
}
7.2 集成测试
@SpringBootTest(webEnvironment = WebEnvironment.RANDOM_PORT)
class ChatIntegrationTest {
@LocalServerPort
private int port;
@Test
void testChatEndpoint() {
// 测试流式响应
WebTestClient client = WebTestClient.bindToServer()
.baseUrl("http://localhost:" + port)
.build();
client.get()
.uri("/api/chat/stream?question=你好")
.accept(MediaType.TEXT_EVENT_STREAM)
.exchange()
.expectStatus().isOk()
.expectHeader().contentTypeCompatibleWith(MediaType.TEXT_EVENT_STREAM)
.expectBody(String.class)
.consumeWith(response -> {
String body = response.getResponseBody();
assertThat(body).contains("data: ");
});
}
@Test
void testDocumentUpload() {
MultipartBodyBuilder builder = new MultipartBodyBuilder();
builder.part("file", new ClassPathResource("test-document.pdf"))
.contentType(MediaType.APPLICATION_PDF);
WebTestClient client = WebTestClient.bindToServer()
.baseUrl("http://localhost:" + port)
.build();
client.post()
.uri("/api/documents/upload")
.contentType(MediaType.MULTIPART_FORM_DATA)
.body(BodyInserters.fromMultipartData(builder.build()))
.exchange()
.expectStatus().isOk()
.expectBody()
.jsonPath("$.success").isEqualTo(true)
.jsonPath("$.documentId").exists();
}
}
7.3 性能测试
@SpringBootTest
class PerformanceTest {
@Autowired
private ChatService chatService;
@Test
void testResponseTimeUnderLoad() {
// 模拟并发请求
int concurrentUsers = 50;
ExecutorService executor = Executors.newFixedThreadPool(concurrentUsers);
List<CompletableFuture<Long>> futures = new ArrayList<>();
for (int i = 0; i < concurrentUsers; i++) {
futures.add(CompletableFuture.supplyAsync(() -> {
long start = System.currentTimeMillis();
chatService.getAnswer("测试问题 " + ThreadLocalRandom.current().nextInt());
return System.currentTimeMillis() - start;
}, executor));
}
// 收集结果
List<Long> responseTimes = futures.stream()
.map(CompletableFuture::join)
.collect(Collectors.toList());
// 计算统计信息
double avg = responseTimes.stream()
.mapToLong(Long::longValue)
.average()
.orElse(0);
double p95 = calculatePercentile(responseTimes, 95);
assertThat(avg).isLessThan(3000); // 平均响应时间小于3秒
assertThat(p95).isLessThan(5000); // 95%的请求小于5秒
}
}
8. 部署与运维建议
8.1 容器化部署
# Dockerfile
FROM eclipse-temurin:17-jre-alpine
# 安装中文字体(处理中文PDF等文档需要)
RUN apk add --no-cache fontconfig ttf-dejavu ttf-droid ttf-freefont ttf-liberation \
&& mkdir -p /usr/share/fonts/chinese \
&& apk add --no-cache wget \
&& wget -O /usr/share/fonts/chinese/simsun.ttc "https://example.com/fonts/simsun.ttc" \
&& fc-cache -fv
WORKDIR /app
# 复制JAR文件
COPY target/ai-knowledge-base.jar app.jar
# 设置时区
ENV TZ=Asia/Shanghai
# JVM参数优化
ENV JAVA_OPTS="-Xmx2g -Xms1g -XX:+UseG1GC -XX:MaxGCPauseMillis=200 -XX:+HeapDumpOnOutOfMemoryError"
# 健康检查
HEALTHCHECK --interval=30s --timeout=3s --start-period=60s --retries=3 \
CMD curl -f http://localhost:8080/actuator/health || exit 1
EXPOSE 8080
ENTRYPOINT ["sh", "-c", "java $JAVA_OPTS -jar app.jar"]
8.2 Kubernetes配置
# deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: ai-knowledge-base
spec:
replicas: 3
selector:
matchLabels:
app: ai-knowledge-base
template:
metadata:
labels:
app: ai-knowledge-base
spec:
containers:
- name: app
image: your-registry/ai-knowledge-base:latest
ports:
- containerPort: 8080
env:
- name: DASHSCOPE_API_KEY
valueFrom:
secretKeyRef:
name: dashscope-secret
key: api-key
- name: REDIS_URL
value: "redis://redis-master:6379"
resources:
requests:
memory: "2Gi"
cpu: "1000m"
limits:
memory: "4Gi"
cpu: "2000m"
livenessProbe:
httpGet:
path: /actuator/health/liveness
port: 8080
initialDelaySeconds: 60
periodSeconds: 10
readinessProbe:
httpGet:
path: /actuator/health/readiness
port: 8080
initialDelaySeconds: 30
periodSeconds: 5
---
# service.yaml
apiVersion: v1
kind: Service
metadata:
name: ai-knowledge-base
spec:
selector:
app: ai-knowledge-base
ports:
- port: 80
targetPort: 8080
type: ClusterIP
8.3 监控告警配置
# prometheus-rules.yaml
apiVersion: monitoring.coreos.com/v1
kind: PrometheusRule
metadata:
name: ai-knowledge-base-rules
spec:
groups:
- name: chat-service
rules:
- alert: HighErrorRate
expr: rate(chat_errors_total[5m]) > 0.1
for: 5m
labels:
severity: warning
annotations:
summary: "聊天服务错误率过高"
description: "错误率超过10%,当前值 {{ $value }}"
- alert: SlowResponse
expr: histogram_quantile(0.95, rate(chat_response_time_seconds_bucket[5m])) > 5
for: 5m
labels:
severity: warning
annotations:
summary: "聊天服务响应时间过长"
description: "95分位响应时间超过5秒,当前值 {{ $value }}s"
- alert: HighTokenUsage
expr: rate(dashscope_tokens_used_total[1h]) > 100000
for: 10m
labels:
severity: warning
annotations:
summary: "API Token使用量过高"
description: "每小时Token使用量超过10万,当前值 {{ $value }}"
8.4 备份与恢复
@Service
@Slf4j
public class BackupService {
private final VectorStore vectorStore;
private final ObjectMapper objectMapper;
@Scheduled(cron = "0 0 2 * * ?") // 每天凌晨2点备份
public void scheduledBackup() {
log.info("开始执行向量数据库备份");
try {
// 1. 导出所有向量数据
List<VectorDocument> allVectors = exportAllVectors();
// 2. 序列化到文件
String backupData = objectMapper.writeValueAsString(allVectors);
// 3. 上传到云存储
uploadToCloudStorage(backupData,
"backup/vectors-" + LocalDateTime.now().format(
DateTimeFormatter.ofPattern("yyyyMMdd-HHmmss")) + ".json");
// 4. 清理旧备份(保留最近7天)
cleanupOldBackups(7);
log.info("向量数据库备份完成");
} catch (Exception e) {
log.error("备份失败", e);
// 发送告警
sendAlert("向量数据库备份失败: " + e.getMessage());
}
}
public void restoreFromBackup(String backupFile) {
log.info("开始从备份恢复: {}", backupFile);
try {
// 1. 从云存储下载备份文件
String backupData = downloadFromCloudStorage(backupFile);
// 2. 反序列化
List<VectorDocument> vectors = objectMapper.readValue(
backupData,
new TypeReference<List<VectorDocument>>() {}
);
// 3. 清空现有数据
vectorStore.deleteAll();
// 4. 恢复数据
for (VectorDocument doc : vectors) {
vectorStore.add(doc.getEmbedding(), doc.getSegment());
}
log.info("恢复完成,共恢复 {} 个向量", vectors.size());
} catch (Exception e) {
log.error("恢复失败", e);
throw new RuntimeException("恢复失败", e);
}
}
}
这几个项目做下来,最大的感受是:技术选型很重要,但更重要的是对业务场景的深入理解。Spring Boot 3.3 + 通义千问这个组合,对于Java团队来说确实是个快速上手的好选择,但真想做出好用的系统,还得在细节上下功夫——比如怎么切分文档、怎么设计提示词、怎么处理多轮对话。那些看起来简单的配置参数,往往对最终效果影响最大。
另外就是监控和运维,AI应用和传统应用不太一样,除了要看CPU、内存这些常规指标,还得关注Token使用量、API调用延迟、回答质量这些业务指标。我见过有的团队一开始没做监控,等发现API费用超了或者回答质量下降的时候,已经晚了。
最后给个实用建议:别想着一口吃成胖子。先从一个小场景开始,比如先把产品手册做成可问答的,跑通了再慢慢加功能。过程中多收集用户反馈,特别是那些“答非所问”的情况,这些都是优化系统的最好材料。
更多推荐



所有评论(0)