AI PaaS 平台采用分层架构设计,主要包含以下层次:
系统支持两种记忆存储方式:
核心实现文件:RedisBasedChatMemory.java
序列化流程:
关键代码:
// 序列化消息列表
private byte[] serialize(List<Message> messages) {
}
// 反序列化消息列表
private List<Message> deserialize(byte[] data) {
}
实现方式:
chat:memory:{conversationId} 作为 Redis 键get 方法获取指定对话的历史消息核心实现文件:RagService.java
流程步骤:
实现方式:使用本地 Ollama 模型生成精准的检索查询词,考虑历史对话上下文
关键代码:
public String generateRetrievalQuery(String userCurrentMsg, List<Message> historyDialog) {
log.info("[RAG-Query] 开始生成检索查询词");
log.info("[RAG-Query] 用户原始消息:{}", userCurrentMsg);
// 构建提示词
StringBuilder promptBuilder = new StringBuilder();
promptBuilder.append("你是一个检索查询优化助手。请根据用户当前问题和历史对话,生成一个精准的检索查询词。\n");
promptBuilder.append("要求:\n");
promptBuilder.append("1. 提取核心关键词\n");
promptBuilder.append("2. 去除无关词汇\n");
promptBuilder.append("3. 保持简洁\n");
promptBuilder.append("4. 只返回查询词,不要有其他内容\n\n");
if (historyDialog != null && !historyDialog.isEmpty()) {
promptBuilder.append("历史对话:\n");
for (Message msg : historyDialog) {
if (msg instanceof UserMessage) {
promptBuilder.append("用户:").append(msg.getText()).append("\n");
}
}
promptBuilder.append("\n");
}
promptBuilder.append("当前问题:").append(userCurrentMsg).append("\n");
promptBuilder.append("生成的检索查询词:");
String prompt = promptBuilder.toString();
log.debug("[RAG-Query] 提示词:\n{}", prompt);
try {
// 使用本地模型生成查询词
ChatResponse response = ollamaChatModel.call(new Prompt(prompt));
String query = response.getResult().getOutput().getText().trim();
log.info("[RAG-Query] 生成的检索查询词:{}", query);
return query;
} catch (Exception e) {
log.error("[RAG-Query] 生成检索查询词失败,使用原始问题", e);
return userCurrentMsg;
}
}
实现方式:调用 RAG 客户端进行检索,对结果进行排序和限制
关键代码:
public Mono<RagRetrieveResponse> retrieve(String query, AgentConfiguration.RagConfig ragConfig) {
if (ragConfig == null || !Boolean.TRUE.equals(ragConfig.getEnabled())) {
log.debug("RAG 未启用,跳过检索");
return Mono.empty();
}
// 转换为 RAG 客户端请求
List<RagRetrieveRequest.VectorStoreConfig> vectorStores = ragConfig.getVectorStores().stream()
.map(config -> RagRetrieveRequest.VectorStoreConfig.builder()
.vectorStoreId(config.getVectorStoreId())
.topK(config.getTopK())
.scoreThreshold(config.getScoreThreshold())
.build())
.collect(Collectors.toList());
RagRetrieveRequest request = RagRetrieveRequest.builder()
.query(query)
.vectorStores(vectorStores)
.build();
return ragClient.retrieve(request);
}
public List<RagRetrieveResponse.Result> sortAndLimitResults(
RagRetrieveResponse response,
Integer maxResults) {
if (response == null || response.getResults() == null) {
log.warn("[RAG-Sort] 响应或结果为空");
return List.of();
}
log.info("[RAG-Sort] 开始排序,原始结果数:{},最大结果数:{}", response.getResults().size(), maxResults);
List<RagRetrieveResponse.Result> sorted = response.getResults().stream()
.sorted(Comparator.comparing(RagRetrieveResponse.Result::getScore).reversed())
.limit(maxResults)
.collect(Collectors.toList());
log.info("[RAG-Sort] 排序完成,保留结果数:{}", sorted.size());
return sorted;
}
核心实现文件:ReActAgent.java
设计理念:
关键代码:
/**
* 运行 Agent
*
* @param userPrompt 用户输入
* @return 执行结果
*/
public String run(String userPrompt) {
if (this.state != AgentState.IDLE) {
throw new RuntimeException("Cannot run agent from state: " + this.state);
}
if (userPrompt == null || userPrompt.trim().isEmpty()) {
throw new RuntimeException("Cannot run agent with empty user prompt");
}
state = AgentState.RUNNING;
messageList.add(new UserMessage(userPrompt));
List<String> results = new ArrayList<>();
try {
for (int i = 0; i < maxSteps && state != AgentState.FINISHED; i++) {
currentStep = i + 1;
log.info("[{}] 执行第 {} 步", name, currentStep);
String stepResult = step();
results.add("Step " + currentStep + ": " + stepResult);
}
if (currentStep >= maxSteps) {
state = AgentState.FINISHED;
results.add("Terminated: Reached max steps (" + maxSteps + ")");
}
return String.join("\n", results);
} catch (Exception e) {
state = AgentState.ERROR;
log.error("[{}] 执行错误: {}", name, e.getMessage(), e);
return "执行错误: " + e.getMessage();
} finally {
this.cleanup();
}
}
/**
* 单步执行
* 将 think() 和 act() 组合在一起
*
* @return 步骤执行结果
*/
public String step() {
try {
// 思考阶段:决定是否需要调用工具
boolean shouldAct = think();
if (!shouldAct) {
return "思考完成 - 无需行动";
}
// 行动阶段:执行工具调用
return act();
} catch (Exception e) {
log.error("[{}] 步骤执行失败: {}", name, e.getMessage(), e);
return "步骤执行失败: " + e.getMessage();
}
}
/**
* 思考阶段
* 分析当前情况,决定是否需要调用工具
*
* @return true - 需要调用工具,false - 无需调用工具
*/
public abstract boolean think();
/**
* 行动阶段
* 执行工具调用并处理结果
*
* @return 行动结果
*/
public abstract String act();
核心实现文件:
AgentConfigRegistry.java:Agent 配置注册表AgentConfiguration.java:Agent 配置类ConfigMapWatcher.java:配置监视器设计理念:
配置文件支持以下结构:
核心实现文件:
McpClientService.java:MCP 客户端服务,负责与 MCP Hub 通信McpTool.java:MCP 工具封装,将 MCP 工具转换为 Spring AI 工具McpServerConfig.java:MCP 服务器配置设计理念:
功能列表:
关键代码:
// MCP 连接管理
public synchronized String connect(String group) {
// 如果已有连接,先关闭
if (sessions.containsKey(group)) {
closeSession(sessions.get(group));
}
SessionInfo sessionInfo = new SessionInfo();
sessionInfo.setGroup(group);
sessionInfo.setSessionId(UUID.randomUUID().toString()); // 临时ID,会被覆盖
try {
// 获取 sessionId
String sessionId = getMCPSessionId();
sessionInfo.setSessionId(sessionId);
sessions.put(group, sessionInfo);
log.info("MCP 连接成功 [group={}, sessionId={}]", group, sessionId);
return sessionId;
} catch (Exception e) {
log.error("MCP 连接失败 [group={}]: {}", group, e.getMessage());
throw new RuntimeException("无法建立 MCP 连接: " + e.getMessage(), e);
}
}
// 工具调用
public ToolCallResult callTool(String group, String toolName, Map<String, Object> arguments) {
SessionInfo session = getOrCreateSession(group);
ToolCallParams params = ToolCallParams.builder()
.name(toolName)
.arguments(arguments)
.build();
JsonRpcRequest request = JsonRpcRequest.builder()
.jsonrpc("2.0")
.id(requestIdGenerator.incrementAndGet())
.method("tools/call")
.params(objectMapper.valueToTree(params))
.build();
JsonNode response = sendRequest(session, request);
return parseToolCallResult(response);
}
// 智能路由搜索工具
public ToolCallResult searchTools(String group, String query, int limit) {
SessionInfo session = getOrCreateSession(group);
if (!"intelligence".equals(session.getMode())) {
// 如果不是智能路由模式,先切换到智能路由
connectIntelligence(group);
session = sessions.get(group);
}
Map<String, Object> arguments = Map.of(
"query", query,
"limit", limit
);
ToolCallParams params = ToolCallParams.builder()
.name("search_tools")
.arguments(arguments)
.build();
JsonRpcRequest request = JsonRpcRequest.builder()
.jsonrpc("2.0")
.id(requestIdGenerator.incrementAndGet())
.method("tools/call")
.params(objectMapper.valueToTree(params))
.build();
JsonNode response = sendRequest(session, request);
return parseToolCallResult(response);
}
McpTool 实现:
public class McpTool implements ToolCallback {
private final String name;
private final String description;
private final String group;
private final McpClientService mcpClientService;
@Override
public ToolDefinition getToolDefinition() {
return ToolDefinition.builder()
.name(name)
.description(description)
// MCP 工具的参数 schema 由 MCP Hub 提供
// 这里使用一个通用的 object schema
.inputSchema("""
{
"type": "object",
"properties": {
"args": {
"type": "object",
"description": "工具参数,由 MCP Hub 定义"
}
}
}
""")
.build();
}
@Override
public String call(String toolInput) {
log.info("执行 MCP 工具 [name={}, group={}, input={}]", name, group, toolInput);
try {
// 解析输入参数
Map<String, Object> args = parseInput(toolInput);
// 调用 MCP Hub
ToolCallResult result;
// 这里简化处理,实际应该根据会话的模式来选择调用方法
// 暂时使用普通的 callTool 方法
result = mcpClientService.callTool(group, name, args);
if (result.isError()) {
log.error("MCP 工具执行失败 [name={}]: {}", name, result.getContent());
return "工具执行失败: " + result.getContent();
}
// 返回结果内容
String content = result.getContent();
log.info("MCP 工具执行成功 [name={}], 结果长度: {}", name,
content != null ? content.length() : 0);
return content != null ? content : "工具执行成功,但无返回内容";
} catch (Exception e) {
log.error("MCP 工具执行异常 [name={}]: {}", name, e.getMessage(), e);
return "工具执行异常: " + e.getMessage();
}
}
}
配置项:
@Data
@Configuration
@ConfigurationProperties(prefix = "mcp.hub")
public class McpServerConfig {
/**
* MCP Hub 基础 URL
*/
private String baseUrl = "https://ai-paas-mcp-endpoint.njuu.top";
/**
* 鉴权 token
*/
private String authorization = "sqGYuMvKgdxmzmTM5lNBgLdVpl6XNnPX";
/**
* 连接超时(秒)
*/
private int connectTimeout = 10;
/**
* 请求超时(秒)
*/
private int requestTimeout = 60;
/**
* 是否自动重连
*/
private boolean autoReconnect = true;
}
系统使用 Spring 的异步处理能力,将保全信息等耗时操作异步解耦,提高系统响应速度。
| 技术 | 版本 | 用途 |
|---|---|---|
| Spring Boot | 3.x | 后端框架 |
| Spring AI | 1.x | AI 集成框架 |
| Redis | 7.x | 持久化存储 |
| Kryo | 5.x | 高效序列化 |
| Ollama | 0.1.x | 本地模型 |
| Maven | 3.x | 构建工具 |
| MCP Hub | - | 工具服务平台 |
AI PaaS 平台采用分层架构设计,实现了配置管理、Agent 核心、记忆系统、RAG 功能和 MCP 集成。系统架构清晰,代码组织合理,功能完整。
系统已经具备了 AI PaaS 平台的核心功能,为未来的扩展和优化奠定了基础。MCP 功能的集成使得系统能够利用外部工具服务,大大扩展了 Agent 的能力范围。