
在AI浪潮中,Python似乎占据了主导地位。但现实是,绝大多数企业级后端系统运行在Java生态上。如何在不重构技术栈的前提下,让Java应用具备AI能力?答案就是:将AI能力作为外部服务集成,并用Java构建坚实的业务编排层。
本文将以 Spring AI 和 LangChain4j 为核心,带你完整实现一个“智能客服+自动化任务”系统。你将掌握:
组件 | 选型 | 理由 |
|---|---|---|
Java框架 | Spring Boot 3.2 + Spring AI | 原生支持AI集成,函数调用标准化 |
LLM对接 | DeepSeek API(兼容OpenAI协议) | 性价比高,中文能力强 |
向量数据库 | Milvus (Docker部署) | 支持百万级向量检索,与Spring AI整合良好 |
文档解析 | Apache Tika + PDFBox | 支持PDF/Word/HTML等格式 |
文本分割 | LangChain4j文本分割器 | 支持递归分割、重叠窗口 |
Embedding | text-embedding-3-small (OpenAI) | 维度适中,效果均衡 |
缓存 | Caffeine + Redis | 多级缓存加速检索 |
RAG是企业落地AI最核心的场景。我们将构建一个产品手册智能问答系统,让Agent能够基于最新文档回答用户问题。
// 文档处理服务
@Service
@Slf4j
public class DocumentIngestionService {
@Autowired
private TikaDocumentParser tikaParser;
@Autowired
private TextSplitter textSplitter;
@Autowired
private EmbeddingService embeddingService;
@Autowired
private VectorStore vectorStore;
@Transactional
public String ingestDocument(MultipartFile file, String category) {
// 1. 解析文档
String rawText = tikaParser.parse(file.getInputStream());
// 2. 清洗与增强
String cleanedText = cleanText(rawText);
Map<String, Object> metadata = extractMetadata(cleanedText, category);
// 3. 智能分块(带重叠)
List<TextSegment> segments = textSplitter.split(
cleanedText,
SegmentSize.MEDIUM, // 500 tokens
SegmentOverlap.SMALL // 50 tokens
);
// 4. 生成向量并存储
List<Document> documents = segments.stream()
.map(seg -> {
float[] embedding = embeddingService.embed(seg.getContent());
return Document.builder()
.content(seg.getContent())
.embedding(embedding)
.metadata(metadata)
.build();
})
.collect(Collectors.toList());
vectorStore.add(documents);
return "文档处理完成,生成 " + documents.size() + " 个向量块";
}
}单一向量检索往往不够,我们采用 混合检索(Hybrid Search):
@Service
public class HybridRetrievalService {
@Autowired
private VectorStore vectorStore;
@Autowired
private ElasticsearchRepository esRepository;
@Autowired
private ReRankService reRankService;
public List<RelevantChunk> retrieve(String query, int topK) {
// 1. 向量检索(语义匹配)
float[] queryEmbedding = embeddingService.embed(query);
List<VectorResult> vectorResults = vectorStore.similaritySearch(
queryEmbedding, topK * 2
);
// 2. 关键词检索(BM25精确匹配)
List<ESResult> keywordResults = esRepository.bm25Search(query, topK * 2);
// 3. 融合与重排(RRF算法)
List<RelevantChunk> fused = fusionResults(vectorResults, keywordResults);
// 4. 重排序(Cross-Encoder精排)
return reRankService.rerank(query, fused, topK);
}
}RRF(倒数排名融合)算法实现:
private List<RelevantChunk> fusionResults(List<VectorResult> vResults,
List<ESResult> kResults) {
Map<String, Double> scoreMap = new HashMap<>();
int k = 60; // RRF常数
for (int i = 0; i < vResults.size(); i++) {
String id = vResults.get(i).getChunkId();
scoreMap.put(id, scoreMap.getOrDefault(id, 0.0) + 1.0 / (k + i + 1));
}
for (int i = 0; i < kResults.size(); i++) {
String id = kResults.get(i).getChunkId();
scoreMap.put(id, scoreMap.getOrDefault(id, 0.0) + 1.0 / (k + i + 1));
}
return scoreMap.entrySet().stream()
.sorted(Map.Entry.<String, Double>comparingByValue().reversed())
.limit(topK)
.map(entry -> buildChunk(entry.getKey()))
.collect(Collectors.toList());
}# application.yml
spring:
ai:
vectorstore:
milvus:
host: localhost
port: 19530
database: default
collection-name: product_docs
embedding-dimension: 1536
index-type: IVF_FLAT
metric-type: COSINE初始化Collection(带Schema):
@Configuration
public class MilvusConfig {
@Bean
public MilvusVectorStore milvusVectorStore(MilvusServiceClient client) {
return new MilvusVectorStore(client, MilvusVectorStoreOptions.builder()
.collectionName("product_docs")
.fields(Arrays.asList(
Field.builder().name("id").type(DataType.Int64).build(),
Field.builder().name("content").type(DataType.VarChar).maxLength(65535).build(),
Field.builder().name("embedding").type(DataType.FloatVector).dimension(1536).build(),
Field.builder().name("category").type(DataType.VarChar).maxLength(100).build(),
Field.builder().name("source").type(DataType.VarChar).maxLength(500).build()
))
.build()
);
}
}Agent是能自主决策的AI实体。我们将构建一个 "技术支持Agent" ,它能:
@Component
@Slf4j
public class SupportAgent {
@Autowired
private ChatClient chatClient;
@Autowired
private ToolRegistry toolRegistry;
@Autowired
private ConversationMemory memory;
@Autowired
private HybridRetrievalService retrievalService;
/**
* Agent主循环:思考-行动-观察
*/
public AgentResponse process(String userInput, String sessionId) {
// 1. 获取对话历史
List<Message> history = memory.getHistory(sessionId);
// 2. 构建系统指令(包含RAG上下文)
String systemPrompt = buildSystemPrompt(userInput);
// 3. 构建工具调用请求
ChatRequest request = ChatRequest.builder()
.systemPrompt(systemPrompt)
.userMessage(userInput)
.history(history)
.tools(toolRegistry.getToolDefinitions())
.build();
// 4. 执行Agent循环(最多3轮迭代)
AgentLoop loop = new AgentLoop(chatClient, toolRegistry);
AgentResponse response = loop.execute(request);
// 5. 保存记忆
memory.add(sessionId, userInput, response.getFinalAnswer());
return response;
}
private String buildSystemPrompt(String query) {
// 检索相关文档块
List<RelevantChunk> chunks = retrievalService.retrieve(query, 5);
String context = chunks.stream()
.map(RelevantChunk::getContent)
.collect(Collectors.joining("\n\n---\n\n"));
return """
你是一位专业的产品技术支持工程师。请基于以下知识库内容回答用户问题。
如果知识库中没有相关信息,请明确告知用户,并建议转接人工服务。
## 知识库内容
%s
## 可用工具
1. query_order(orderId) - 查询订单状态
2. create_ticket(title, description, priority) - 创建工单
3. search_known_issues(keyword) - 搜索已知问题库
请用中文友好回复,并主动提供帮助。
""".formatted(context);
}
}定义工具接口:
@FunctionalInterface
public interface AgentTool {
String execute(Map<String, Object> parameters);
}
@Component
public class ToolRegistry {
private final Map<String, AgentTool> tools = new ConcurrentHashMap<>();
@PostConstruct
public void init() {
register("query_order", this::queryOrder);
register("create_ticket", this::createTicket);
register("search_known_issues", this::searchKnownIssues);
}
public List<ToolDefinition> getToolDefinitions() {
return tools.keySet().stream()
.map(name -> ToolDefinition.builder()
.name(name)
.description(getDescription(name))
.parameters(getParameterSchema(name))
.build()
)
.collect(Collectors.toList());
}
// 工具实现示例
private String queryOrder(Map<String, Object> params) {
String orderId = (String) params.get("orderId");
// 调用订单服务
OrderDTO order = orderService.getOrder(orderId);
return "订单 %s 状态: %s,预计送达: %s"
.formatted(orderId, order.getStatus(), order.getEstDelivery());
}
}public class AgentLoop {
private final ChatClient client;
private final ToolRegistry toolRegistry;
private static final int MAX_ITERATIONS = 3;
public AgentResponse execute(ChatRequest request) {
List<Message> messages = new ArrayList<>(request.getHistory());
messages.add(new UserMessage(request.getUserMessage()));
for (int i = 0; i < MAX_ITERATIONS; i++) {
// 调用LLM
ChatResponse response = client.call(
request.getSystemPrompt(),
messages,
request.getTools()
);
// 检查是否有工具调用
if (response.hasToolCalls()) {
// 执行工具
List<ToolResult> results = executeTools(response.getToolCalls());
// 将工具结果加入上下文
messages.add(new ToolResponseMessage(results));
continue;
}
// 无工具调用,返回最终答案
return new AgentResponse(response.getContent(), messages);
}
// 达到最大迭代次数
return new AgentResponse("抱歉,我无法在限定步骤内完成您的请求,请简化问题或联系人工。", messages);
}
}@Service
public class ConversationMemory {
@Autowired
private RedisTemplate<String, Object> redisTemplate;
@Autowired
private Cache<String, List<Message>> localCache; // Caffeine
// 滑动窗口管理(保留最近20轮)
public void add(String sessionId, String userMsg, String assistantMsg) {
String key = "mem:" + sessionId;
List<Message> history = getHistory(sessionId);
history.add(new UserMessage(userMsg));
history.add(new AssistantMessage(assistantMsg));
// 截断至最大长度
if (history.size() > 40) {
history = history.subList(history.size() - 40, history.size());
}
// 先存本地缓存,异步写Redis
localCache.put(key, history);
CompletableFuture.runAsync(() ->
redisTemplate.opsForValue().set(key, history, 24, TimeUnit.HOURS)
);
}
public List<Message> getHistory(String sessionId) {
String key = "mem:" + sessionId;
List<Message> cached = localCache.getIfPresent(key);
if (cached != null) return cached;
List<Message> fromRedis = (List<Message>) redisTemplate.opsForValue().get(key);
if (fromRedis != null) {
localCache.put(key, fromRedis);
return fromRedis;
}
return new ArrayList<>();
}
// 摘要压缩(当对话过长时)
public String summarizeHistory(String sessionId) {
List<Message> history = getHistory(sessionId);
if (history.size() < 30) return null;
// 调用LLM生成摘要
String summary = summaryService.summarize(history);
// 用摘要替换早期对话
return summary;
}
}ai:
memory:
max-turns: 20 # 保留轮数
summary-trigger: 30 # 触发摘要的轮数
redis-ttl: 86400 # 24小时过期@RestController
@RequestMapping("/api/v1/ai")
@Slf4j
public class AIController {
@Autowired
private SupportAgent agent;
@Autowired
private DocumentIngestionService ingestionService;
// 对话接口
@PostMapping("/chat")
public ResponseEntity<ChatResponse> chat(@RequestBody ChatRequest request) {
String sessionId = request.getSessionId() != null ?
request.getSessionId() : UUID.randomUUID().toString();
AgentResponse response = agent.process(request.getMessage(), sessionId);
return ResponseEntity.ok(ChatResponse.builder()
.sessionId(sessionId)
.answer(response.getFinalAnswer())
.sources(response.getSources())
.tokensUsed(response.getTokenUsage())
.build()
);
}
// 文档上传接口
@PostMapping("/documents/upload")
public ResponseEntity<UploadResult> uploadDocument(
@RequestParam("file") MultipartFile file,
@RequestParam("category") String category) {
String docId = ingestionService.ingestDocument(file, category);
return ResponseEntity.ok(new UploadResult(docId, "上传成功"));
}
// 手动触发知识库更新
@PostMapping("/documents/refresh")
public ResponseEntity<Void> refreshKnowledgeBase() {
ingestionService.refreshAllDocuments();
return ResponseEntity.accepted().build();
}
}@GetMapping(value = "/chat/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<ServerSentEvent<String>> streamChat(@RequestParam String message,
@RequestParam String sessionId) {
return agent.processStream(message, sessionId)
.map(content -> ServerSentEvent.<String>builder()
.data(content)
.event("message")
.build()
)
.onErrorResume(e -> Flux.just(
ServerSentEvent.<String>builder()
.data("发生错误: " + e.getMessage())
.event("error")
.build()
));
}@Component
public class AICacheManager {
@Cacheable(value = "rag_results", key = "#query + '_' + #topK")
public List<RelevantChunk> getCachedResults(String query, int topK) {
return retrievalService.retrieve(query, topK);
}
// 使用Caffeine本地缓存 + Redis分布式缓存
}@Service
public class BatchEmbeddingService {
@Async("embeddingExecutor")
public CompletableFuture<List<float[]>> batchEmbed(List<String> texts) {
// 使用OpenAI批量接口(一次请求处理多段文本)
return CompletableFuture.completedFuture(
embeddingClient.embed(texts, EmbeddingOptions.DEFAULT)
);
}
}@Component
public class FallbackHandler {
@Autowired
private CircuitBreakerFactory circuitBreakerFactory;
public AgentResponse fallbackForLLM(String userInput, Throwable t) {
log.error("LLM服务不可用,使用降级方案", t);
// 方案1: 使用本地规则匹配
String ruleAnswer = ruleEngine.match(userInput);
if (ruleAnswer != null) {
return new AgentResponse(ruleAnswer, "基于规则");
}
// 方案2: 返回预置的兜底回复
return new AgentResponse(
"AI服务暂时繁忙,请稍后重试或联系人工客服。",
"降级响应"
);
}
}management:
endpoints:
web:
exposure:
include: health,prometheus,metrics
metrics:
tags:
application: ai-agent
# 自定义监控指标
@Bean
public MeterRegistryCustomizer<MeterRegistry> metricsCustomizer() {
return registry -> registry.config().commonTags("service", "ai-orchestrator");
}关键监控指标:
rag.retrieval.duration - 检索耗时agent.llm.tokens - Token消耗agent.tool.calls - 工具调用次数rag.hit.rate - 知识库命中率version: '3.8'
services:
milvus:
image: milvusdb/milvus:v2.4.0
ports:
- "19530:19530"
environment:
ETCD_ENDPOINTS: etcd:2379
MINIO_ADDRESS: minio:9000
depends_on:
- etcd
- minio
redis:
image: redis:7-alpine
ports:
- "6379:6379"
elasticsearch:
image: elasticsearch:8.10.2
environment:
- discovery.type=single-node
- xpack.security.enabled=false
ports:
- "9200:9200"
ai-service:
build: .
ports:
- "8080:8080"
environment:
- SPRING_PROFILES_ACTIVE=prod
- DEEPSEEK_API_KEY=${DEEPSEEK_API_KEY}
depends_on:
- milvus
- redis
- elasticsearch#!/bin/bash
# 构建Jar包
./gradlew clean bootJar
# 构建Docker镜像
docker build -t ai-agent:latest .
# 启动所有服务
docker-compose up -d
# 健康检查
curl -s http://localhost:8080/actuator/health | jq .用户输入:
"我的设备型号是X200,最近经常自动重启,而且日志里有'ERR-2024'错误码,能帮我查一下是什么问题吗?"
Agent执行流程:
search_known_issues("ERR-2024") 返回已知问题库信息"根据知识库和已知问题库的信息,X200设备的ERR-2024错误通常与固件版本低于v3.2.0有关。建议您执行以下操作:
如果问题仍未解决,我可以帮您创建一个技术支持工单。需要我协助创建吗?"
多轮对话(记忆发挥作用):
用户:"好的,我升级完了,但问题还在" Agent:(回忆上下文)"您已升级固件,但X200重启问题仍存在。这可能是硬件故障,我已经为您自动创建了工单 #T-2024-0815,高级工程师将在2小时内联系您。"
✅ 企业级RAG知识库:文档处理→向量检索→混合搜索→重排序 ✅ Agent智能体:工具调用→ReAct循环→多轮记忆 ✅ 生产级Java实现:Spring Boot + 缓存 + 监控 + 容错
完整的项目源码(包含本实战全部代码)请访问:
https://github.com/your-org/ai-agent-java
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。