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-001300 tokens1536通用多语言中等
ops-text-embedding-zh-0011024 tokens768纯中文场景优秀
ops-text-embedding-0028192 tokens1024长文本良好

对于中文知识库,我推荐用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费用超了或者回答质量下降的时候,已经晚了。

最后给个实用建议:别想着一口吃成胖子。先从一个小场景开始,比如先把产品手册做成可问答的,跑通了再慢慢加功能。过程中多收集用户反馈,特别是那些“答非所问”的情况,这些都是优化系统的最好材料。

Logo

DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。

更多推荐