【AI 工程化】实战:从 0 搭一个生产级 AI Agent 服务(SSE 流式 + 工具调用 + RAG + MCP,附完整源码)
技术栈:Java 17 · Spring Boot 4.1.1 · Spring AI 2.0.1 · PostgreSQL + pgvector · Docker Compose
适用场景:企业知识助手、智能客服、内部工具型 Agent
阅读前置:会写 Spring Boot 的 Controller 和@Service就够了,不需要懂 Python 和模型训练
为什么 2026 年 Java 工程师该动手了
先看两组最近的数据。
GitHub Trending 周榜 18 个仓库里,10 个属于 Agent / Skill / MCP / CLI 工具链(占比 55.6%);月榜这个比例升到 61.1%。涨星最猛的 vectorize-io/hindsight 一周 +11,089,paperclipai/paperclip +7,364 —— 热度已经不在"模型"上,而在"怎么把模型接进业务系统"这件事上。
再看 GitHub Java 周榜 17 个仓库:elasticsearch、kafka、dbeaver、dataease、Stirling-PDF、Chat2DB、jadx……清一色是经典中间件和工具软件,没有一个 AI 应用项目。
这两组数据放在一起,指向一个很明确的事实:
AI 应用层的需求在爆发,但 Java 侧的公开实践几乎还是空白。
而 Spring 官方其实已经把这件事做完了 —— Spring AI 2.0.1 已经 GA,配套 Spring Boot 4.1.1。它的 ChatClient 抽象掉了不同厂商的差异,Advisor 机制把 RAG、记忆、工具调用做成了可插拔的拦截器链。
问题是中文社区的生产级实战内容太少,大部分文章停在"hello world 调一次 API"。这篇文章把它写完整:从架构设计、表结构、核心代码,一直到 Docker 部署和前端联调。
一、先看效果:3 条命令确认这篇能跑通
1.1 同步调用(普通问答)
curl -X POST http://localhost:8080/api/v1/chat \
-H "Content-Type: application/json" \
-d '{"conversationId":"demo-001","message":"我们公司的退货政策是几天?"}'
{
"code": 0,
"message": "success",
"data": {
"conversationId": "demo-001",
"content": "根据知识库《售后服务手册 v3.2》,自签收之日起 7 个自然日内可无理由退货……",
"referenceCount": 3,
"promptTokens": 1240,
"completionTokens": 186,
"costMillis": 2317
}
}
1.2 流式调用(SSE,打字机效果)
curl -N -X POST http://localhost:8080/api/v1/chat/stream \
-H "Content-Type: application/json" \
-H "Accept: text/event-stream" \
-d '{"conversationId":"demo-001","message":"帮我查一下订单 SO20260928001 的状态"}'
event: message
data: {"v":"正在"}
event: message
data: {"v":"为你查询订单……"}
event: message
data: {"v":"订单 SO20260928001 当前状态为「已发货」"}
event: done
data: {"costMillis":3402}
注意最后那个 done 事件 —— 流式接口也必须回传 Token 用量和耗时,否则线上你根本不知道钱花在哪了。
1.3 函数调用(Agent 真正"动手")
用户问:「帮我查一下订单 SO20260928001 的状态」
模型不会瞎编,而是触发我们注册的 @Tool 方法去查数据库,拿到真实结果再组织语言。这是"聊天机器人"和"Agent"的分界线。
二、技术选型:为什么是 Java,不是 FastAPI
先把选型逻辑讲清楚,避免无意义的争论。
| 维度 | Python(FastAPI / LangChain) | Java(Spring Boot + Spring AI) |
|---|---|---|
| 模型推理、微调、数据处理 | ✅ 生态碾压 | ❌ 不适合 |
| 论文复现、算法实验 | ✅ 首选 | ❌ 不适合 |
| 高并发 API 网关 | ⚠️ GIL + 多进程,连接治理成本高 | ✅ 线程池 + 虚拟线程成熟 |
| 企业系统集成(事务、权限、审计) | ⚠️ 弱 | ✅ 与现有权限体系无缝 |
| 流式连接治理(SSE / WebSocket 断连、超时、背压) | ⚠️ 需要自己造 | ✅ 框架层已解决 |
| 团队维护成本 | ⚠️ 后端团队要学新语言 | ✅ 现有 Java 团队直接上手 |
结论是"双轨制"而不是二选一:
- 模型侧、离线数据处理、RAG 语料清洗 → Python
- 面向用户的 AI 网关、Agent 编排、企业系统对接 → Java
本文聚焦后者,也就是大多数公司真正要长期维护的那一层。
2.1 版本信息(发布前已核实)
| 组件 | 版本 | 备注 |
|---|---|---|
| JDK | 17 | Spring Boot 4 的编译基线就是 17 |
| Spring Boot | 4.1.1 | 最新稳定版 |
| Spring AI | 2.0.1 | 通过 spring-ai-bom 统一管理 |
| PostgreSQL | 17 | 向量检索使用 pgvector 扩展 |
⚠️ 第一个坑:Spring Boot 4 里
spring-boot-starter-web已经被标记为 deprecated,官方建议换成spring-boot-starter-webmvc。这个改动在 4.1.1 的 starter POM 描述里写得很清楚,但很多教程还在用旧名字。
三、整体设计
3.1 功能模块图
┌─────────────────────────────────────────────────────────────────┐
│ 前端 / 调用方 │
│ Web · 企业微信 · 钉钉 · 内部系统(SSE / HTTP) │
└───────────────────────────┬─────────────────────────────────────┘
│ POST /api/v1/chat
│ POST /api/v1/chat/stream ← SSE
┌───────────────────────────▼─────────────────────────────────────┐
│ 接入层 (Controller) │
│ ChatController · 参数校验 · SSE 心跳 · 全局异常拦截 │
└───────────────────────────┬─────────────────────────────────────┘
┌───────────────────────────▼─────────────────────────────────────┐
│ Agent 编排层 (Service) │
│ ChatClient + Advisor 链: │
│ ① MessageChatMemoryAdvisor → 多轮会话记忆 │
│ ② QuestionAnswerAdvisor → RAG 知识库检索 │
│ ③ ToolCallingAdvisor → 工具调用循环(框架自动注册) │
│ ④ SimpleLoggerAdvisor → Prompt / 响应日志 │
└──────┬──────────────────┬──────────────────┬─────────────────────┘
│ │ │
┌──────▼──────┐ ┌───────▼───────┐ ┌──────▼──────────────────┐
│ 会话记忆 │ │ 向量知识库 │ │ 工具集 / MCP Server │
│ JDBC 表 │ │ pgvector │ │ 订单查询 · 工单创建 │
│ SPRING_AI_ │ │ vector_store │ │ 外部系统 API │
│ CHAT_MEMORY │ │ │ │ │
└─────────────┘ └───────────────┘ └─────────────────────────┘
│
┌───────────────────────────▼─────────────────────────────────────┐
│ 模型层(OpenAI 兼容协议) │
│ 对话模型:DeepSeek / 通义 / Kimi / GPT │
│ 向量模型:独立端点(如 DashScope embedding) │
└─────────────────────────────────────────────────────────────────┘
3.2 流式链路时序(文字描述)
客户端 Controller AgentService LLM 服务
│ │ │ │
│─ POST /chat/stream ─▶│ │ │
│ │─ 校验参数、建 Emitter │ │
│ │─ 调用 streamChat ───▶│ │
│ │ │─ 装配 Prompt ──────▶│
│ │ │ (记忆 + RAG 上下文) │
│ │ │◀─ 首包 token ───────│
│◀─ event: message ────│◀─ onNext(delta) ────│ │
│◀─ event: message ────│◀─ onNext(delta) ────│◀────────────────────│
│ ...持续推送,中间穿插心跳注释行,防止网关判定空闲断开... │
│◀─ event: done ───────│◀─ onComplete ───────│ │
│ │─ 写入 Token / 耗时审计 │ │
│ │─ emitter.complete() │ │
3.3 核心表结构
-- ① 会话记忆表:由 Spring AI 自动初始化,这里列出来是为了方便排查问题
-- 2.0 起新增 sequence_id 列用于确定消息顺序(timestamp 仍然保留)
CREATE TABLE IF NOT EXISTS SPRING_AI_CHAT_MEMORY (
conversation_id VARCHAR(36) NOT NULL,
content TEXT NOT NULL,
type VARCHAR(10) NOT NULL,
timestamp TIMESTAMP NOT NULL,
sequence_id BIGINT NOT NULL,
PRIMARY KEY (conversation_id, sequence_id)
);
-- ② 业务侧会话表:记录会话归属、标题、状态
CREATE TABLE t_ai_conversation (
id BIGSERIAL PRIMARY KEY,
conversation_id VARCHAR(64) NOT NULL UNIQUE,
user_id VARCHAR(64) NOT NULL,
title VARCHAR(200),
status SMALLINT DEFAULT 1, -- 1 进行中 2 已结束
create_time TIMESTAMP DEFAULT now(),
update_time TIMESTAMP DEFAULT now()
);
CREATE INDEX idx_ai_conv_user ON t_ai_conversation(user_id, update_time DESC);
-- ③ AI 调用审计表:Token 消耗、耗时、模型、成败,用于成本核算与问题追溯
CREATE TABLE t_ai_call_log (
id BIGSERIAL PRIMARY KEY,
trace_id VARCHAR(64) NOT NULL,
conversation_id VARCHAR(64),
user_id VARCHAR(64),
model VARCHAR(64),
prompt_tokens INT DEFAULT 0,
completion_tokens INT DEFAULT 0,
total_tokens INT DEFAULT 0,
cost_millis BIGINT DEFAULT 0,
stream BOOLEAN DEFAULT FALSE,
success BOOLEAN DEFAULT TRUE,
error_msg VARCHAR(500),
create_time TIMESTAMP DEFAULT now()
);
CREATE INDEX idx_ai_log_trace ON t_ai_call_log(trace_id);
CREATE INDEX idx_ai_log_time ON t_ai_call_log(create_time DESC);
-- ④ 向量库表:由 pgvector starter 自动创建,示意结构
CREATE TABLE IF NOT EXISTS vector_store (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
content TEXT,
metadata JSON,
embedding VECTOR(1024)
);
为什么要有 t_ai_call_log? 因为 AI 应用上线后最常被问的两个问题就是「昨天花了多少钱」和「这条回答为什么不对」。没有这张表,两个都答不上来。
四、工程搭建
4.1 pom.xml(完整依赖)
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>4.1.1</version>
<relativePath/>
</parent>
<groupId>com.example</groupId>
<artifactId>ai-agent-service</artifactId>
<version>1.0.0</version>
<name>ai-agent-service</name>
<properties>
<java.version>17</java.version>
<!-- Spring AI 统一版本管理 -->
<spring-ai.version>2.0.1</spring-ai.version>
</properties>
<dependencyManagement>
<dependencies>
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-bom</artifactId>
<version>${spring-ai.version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>
<dependencies>
<!-- Web:注意 Boot 4 里 starter-web 已 deprecated,用 webmvc -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-webmvc</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-validation</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jdbc</artifactId>
</dependency>
<!-- 对话模型:OpenAI 兼容协议(DeepSeek / 通义 / Kimi / GPT 都能用) -->
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-starter-model-openai</artifactId>
</dependency>
<!-- 向量库:PostgreSQL + pgvector -->
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-starter-vector-store-pgvector</artifactId>
</dependency>
<!-- RAG 的 QuestionAnswerAdvisor 在这个独立制品里,不加会编译不过 -->
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-vector-store-advisor</artifactId>
</dependency>
<!-- 会话记忆持久化到 JDBC -->
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-starter-model-chat-memory-repository-jdbc</artifactId>
</dependency>
<!-- MCP 客户端:让 Agent 挂载外部工具服务 -->
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-starter-mcp-client</artifactId>
</dependency>
<dependency>
<groupId>org.postgresql</groupId>
<artifactId>postgresql</artifactId>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
</project>
⚠️ 第二个坑:Spring AI 的 starter 命名在 1.0 前后改过一轮。老文章里的
spring-ai-openai-spring-boot-starter已经废弃,现在是「spring-ai-starter-{能力}-{厂商或组件}」这种格式。名字写错的表现是「启动正常但 Bean 注入不进来」,非常难查。
4.2 application.yml
server:
port: 8080
# SSE 长连接需要放开 Tomcat 的超时与连接数
tomcat:
connection-timeout: 60s
max-connections: 10000
threads:
max: 400
spring:
application:
name: ai-agent-service
datasource:
url: jdbc:postgresql://localhost:5432/aiagent
username: ${DB_USER:postgres}
password: ${DB_PASSWORD:postgres}
hikari:
maximum-pool-size: 20
connection-timeout: 5000
ai:
# ★ Spring AI 2.x 的模型启停开关(1.x 的 xxx.chat.enabled 已移除)
model:
chat: openai
embedding: openai
openai:
# 对话模型走 DeepSeek 的 OpenAI 兼容端点
api-key: ${LLM_API_KEY}
base-url: https://api.deepseek.com
timeout: 60s
chat:
options:
model: deepseek-chat
temperature: 0.3
max-tokens: 2048
# 向量模型走另一个兼容端点(DeepSeek 不提供 embedding,需要单独配置)
embedding:
base-url: ${EMBEDDING_BASE_URL:https://dashscope.aliyuncs.com/compatible-mode/v1}
api-key: ${EMBEDDING_API_KEY}
options:
model: text-embedding-v4
dimensions: 1024
vectorstore:
pgvector:
initialize-schema: true
table-name: vector_store
dimensions: 1024 # 必须和 embedding 维度一致,否则插入直接报错
index-type: HNSW
distance-type: COSINE_DISTANCE
max-document-batch-size: 100
chat:
memory:
repository:
jdbc:
# 生产建议设为 never,改由 Flyway 管理表结构
initialize-schema: embedded
# MCP 客户端:挂载外部工具服务(示例为官方文件系统 server)
mcp:
client:
enabled: true
type: SYNC
request-timeout: 30s
stdio:
connections:
filesystem:
command: npx
args:
- "-y"
- "@modelcontextprotocol/server-filesystem"
- "/data/ai-agent-workspace"
logging:
level:
com.example.aiagent: DEBUG
org.springframework.ai.chat.client.advisor: DEBUG
management:
endpoints:
web:
exposure:
include: health,metrics,prometheus
⚠️ 第三个坑(最容易翻车):Spring AI 1.x 用
spring.ai.openai.chat.enabled=true来启停模型,2.x 已经把这个属性删掉了,改成顶层开关spring.ai.model.chat。如果你从 1.x 升级过来,配置文件里保留旧属性不会报错,但模型 Bean 会悄悄不生效 —— 表现为启动不报错、调用时提示找不到ChatModel。
五、核心代码
5.1 ChatClient 装配(Advisor 链是灵魂)
package com.example.aiagent.config;
import org.springframework.ai.chat.client.ChatClient;
import org.springframework.ai.chat.client.advisor.MessageChatMemoryAdvisor;
import org.springframework.ai.chat.client.advisor.SimpleLoggerAdvisor;
import org.springframework.ai.chat.memory.ChatMemory;
import org.springframework.ai.vectorstore.SearchRequest;
import org.springframework.ai.vectorstore.VectorStore;
import org.springframework.ai.vectorstore.advisor.QuestionAnswerAdvisor;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
* ChatClient 统一装配。
*
* <p>设计说明:把 Advisor 链集中在一处声明,业务代码只依赖 ChatClient 接口,
* 后续替换向量库、调整检索参数都不需要改动 Service 层(依赖倒置)。</p>
*
* @author 你的名字
*/
@Configuration
public class ChatClientConfig {
/** 知识库检索 TopK,生产环境建议做成配置中心动态下发 */
private static final int RAG_TOP_K = 5;
/** 相似度阈值:低于该值的文档不参与拼接,能显著降低幻觉 */
private static final double RAG_SIMILARITY_THRESHOLD = 0.55;
@Bean
public ChatClient chatClient(ChatClient.Builder builder,
ChatMemory chatMemory,
VectorStore vectorStore) {
// ① 会话记忆:按 conversationId 隔离多轮上下文
MessageChatMemoryAdvisor memoryAdvisor =
MessageChatMemoryAdvisor.builder(chatMemory).build();
// ② RAG:命中知识库后拼接进 Prompt
QuestionAnswerAdvisor ragAdvisor = QuestionAnswerAdvisor.builder(vectorStore)
.searchRequest(SearchRequest.builder()
.topK(RAG_TOP_K)
.similarityThreshold(RAG_SIMILARITY_THRESHOLD)
.build())
.build();
return builder
.defaultSystem("""
你是企业内部智能助手。请严格遵守以下规则:
1. 优先依据「参考资料」回答;资料中没有的内容,明确说明"知识库中未找到相关信息",不要编造。
2. 涉及订单、工单等实时数据时,必须调用工具查询,不要凭记忆回答。
3. 回答使用简体中文,结构清晰;涉及步骤时用有序列表。
4. 不要输出任何内部的接口地址、密钥或实现细节。
""")
// 顺序有讲究:记忆 → RAG → 日志
// 工具调用的 ToolCallingAdvisor 由框架自动注册,无需手动添加
.defaultAdvisors(memoryAdvisor, ragAdvisor, new SimpleLoggerAdvisor())
.build();
}
}
这段代码值得单独说三点:
similarityThreshold是降低幻觉最有效的一个参数。很多实现只设topK,导致不管像不像都硬塞 5 篇文档进 Prompt,模型就开始"编"。- 工具调用不需要手动加 Advisor。Spring AI 2.x 里
ChatClient会自动注册ToolCallingAdvisor,只要在调用时用.tools(...)传入工具对象即可。 defaultSystem用 Java 17 的文本块写提示词,可读性比字符串拼接好太多,也方便做版本管理。
5.2 工具调用(Agent 的"手")
package com.example.aiagent.tool;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.ai.tool.annotation.Tool;
import org.springframework.ai.tool.annotation.ToolParam;
import org.springframework.stereotype.Component;
import java.time.LocalDateTime;
/**
* 订单查询工具集。
*
* <p>注意:工具方法的 description 直接决定模型会不会调、调得对不对,
* 请像写接口文档一样认真写,并明确参数格式。</p>
*
* @author 你的名字
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class OrderTools {
private final OrderService orderService;
@Tool(description = "根据订单号查询订单当前状态、金额和预计送达时间。订单号格式为 SO + 11 位数字,例如 SO20260928001")
public OrderInfo queryOrder(
@ToolParam(description = "订单号,格式 SO + 11 位数字") String orderNo) {
log.info("[Tool] queryOrder 入参 orderNo={}", orderNo);
// 工具方法必须自己兜异常:抛出去会中断整个 Agent 循环
try {
OrderDTO order = orderService.getByOrderNo(orderNo);
if (order == null) {
return OrderInfo.notFound(orderNo);
}
return new OrderInfo(order.getOrderNo(),
order.getStatusDesc(),
order.getAmount().toPlainString(),
order.getEstimatedDelivery().toString(),
null);
} catch (Exception e) {
log.error("[Tool] queryOrder 失败 orderNo={}", orderNo, e);
// 返回结构化错误信息,让模型能向用户解释,而不是直接 500
return OrderInfo.error(orderNo, "订单系统暂时不可用,请稍后重试");
}
}
@Tool(description = "创建售后工单。需要用户提供订单号和问题描述,问题描述不少于 10 个字")
public String createTicket(
@ToolParam(description = "订单号") String orderNo,
@ToolParam(description = "问题描述") String description) {
log.info("[Tool] createTicket orderNo={}", orderNo);
String ticketNo = orderService.createTicket(orderNo, description);
return "工单创建成功,工单号:" + ticketNo
+ ",客服会在 24 小时内联系您。创建时间:" + LocalDateTime.now();
}
/**
* 工具返回值。必须能被序列化,字段名建议自带语义,模型更容易准确引用。
*/
public record OrderInfo(String orderNo, String status, String amount,
String estimatedDelivery, String message) {
public static OrderInfo notFound(String orderNo) {
return new OrderInfo(orderNo, "NOT_FOUND", null, null,
"未查询到该订单,请确认订单号是否正确");
}
public static OrderInfo error(String orderNo, String message) {
return new OrderInfo(orderNo, "ERROR", null, null, message);
}
}
}
生产环境三条铁律:
- 工具方法绝不向外抛异常。异常会打断工具调用循环,用户体验直接从"能回答"变成"系统错误"。
- 工具描述要写清参数格式。描述里写
SO + 11 位数字,模型就不会把12345传进来,能省掉大量参数校验。 - 工具返回结构化对象,不要返回长字符串。字段有语义,模型才能准确引用;长字符串会白白消耗 Token。
5.3 同步接口(最简同步方案)
package com.example.aiagent.controller;
import jakarta.validation.Valid;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.ai.chat.client.ChatClient;
import org.springframework.ai.chat.memory.ChatMemory;
import org.springframework.ai.chat.model.ChatResponse;
import org.springframework.web.bind.annotation.*;
/**
* AI 对话接口(同步)。
*
* @author 你的名字
*/
@Slf4j
@RestController
@RequestMapping("/api/v1/chat")
@RequiredArgsConstructor
public class ChatController {
private final ChatClient chatClient;
private final OrderTools orderTools;
private final AiCallLogService aiCallLogService;
/**
* 同步对话:适合后台任务、短问答、需要拿到完整结果再处理的场景。
*/
@PostMapping
public ApiResponse<ChatVO> chat(@Valid @RequestBody ChatRequest request) {
long start = System.currentTimeMillis();
String traceId = TraceIdUtils.current();
ChatResponse response = chatClient.prompt()
.user(request.message())
.tools(orderTools) // ToolCallingAdvisor 会自动注册
// 会话 ID 必须显式传,否则所有用户会串到同一个上下文
.advisors(a -> a.param(ChatMemory.CONVERSATION_ID, request.conversationId()))
.call()
.chatResponse();
long cost = System.currentTimeMillis() - start;
String content = response.getResult().getOutput().getText();
// ★ 链路日志:Token 消耗 / 耗时 / 内容长度,线上排障全靠它
log.info("[AI-CALL] traceId={} convId={} sync=true cost={}ms contentLen={}",
traceId, request.conversationId(), cost, content.length());
aiCallLogService.record(traceId, request.conversationId(), response, cost, false, null);
return ApiResponse.ok(ChatVO.of(request.conversationId(), content, response, cost));
}
}
配套的请求对象(用 record + jakarta.validation,干净且不可变):
package com.example.aiagent.controller;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.Size;
/**
* 对话请求。
* <p>conversationId 由前端生成并持久化(例如 uuid),用于隔离多轮上下文。</p>
*/
public record ChatRequest(
@NotBlank(message = "conversationId 不能为空")
@Size(max = 64, message = "conversationId 长度不能超过 64")
String conversationId,
@NotBlank(message = "提问内容不能为空")
@Size(max = 4000, message = "单次提问不能超过 4000 字")
String message
) {
}
5.4 流式接口(SSE 的两种写法,推荐第一种)
方案 A:SseEmitter(Spring MVC 原生,连接生命周期可控,生产首选)
package com.example.aiagent.controller;
import jakarta.validation.Valid;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.ai.chat.client.ChatClient;
import org.springframework.ai.chat.memory.ChatMemory;
import org.springframework.ai.chat.model.ChatResponse;
import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.*;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import reactor.core.publisher.Flux;
import java.io.IOException;
import java.util.concurrent.atomic.AtomicBoolean;
/**
* 流式对话接口(SSE)。
*
* @author 你的名字
*/
@Slf4j
@RestController
@RequestMapping("/api/v1/chat")
@RequiredArgsConstructor
public class ChatStreamController {
private final ChatClient chatClient;
private final OrderTools orderTools;
private final AiCallLogService aiCallLogService;
/** 单次 SSE 连接最长存活时间:5 分钟 */
private static final long SSE_TIMEOUT_MILLIS = 5 * 60 * 1000L;
/**
* 流式对话。
*
* <p>实现要点:
* 1. 用 SseEmitter 而不是在 Controller 里 block(),避免占满 Tomcat 线程;
* 2. 注册 onTimeout / onError / onCompletion,保证任何异常路径都能释放连接;
* 3. 流结束后回传耗时,便于前端展示与成本核算。</p>
*/
@PostMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE + ";charset=UTF-8")
public SseEmitter stream(@Valid @RequestBody ChatRequest request) {
SseEmitter emitter = new SseEmitter(SSE_TIMEOUT_MILLIS);
String traceId = TraceIdUtils.current();
long start = System.currentTimeMillis();
AtomicBoolean finished = new AtomicBoolean(false);
// 超时 / 异常 / 完成三条路径都要能走到,否则连接会泄漏
emitter.onTimeout(() -> {
log.warn("[SSE] 超时 traceId={} convId={}", traceId, request.conversationId());
emitter.complete();
});
emitter.onError(e -> log.warn("[SSE] 异常 traceId={}", traceId, e));
emitter.onCompletion(() -> finished.set(true));
Flux<ChatResponse> flux = chatClient.prompt()
.user(request.message())
.tools(orderTools)
.advisors(a -> a.param(ChatMemory.CONVERSATION_ID, request.conversationId()))
.stream()
.chatResponse();
flux.subscribe(
response -> {
String delta = extractDelta(response);
if (delta == null || delta.isEmpty()) {
return;
}
try {
emitter.send(SseEmitter.event().name("message")
.data(StreamChunk.of(delta), MediaType.APPLICATION_JSON));
} catch (IOException e) {
// 用户中途关闭页面:正常现象,停止推送即可
log.info("[SSE] 客户端已断开 traceId={}", traceId);
emitter.completeWithError(e);
}
},
error -> {
log.error("[SSE] 生成失败 traceId={}", traceId, error);
aiCallLogService.record(traceId, request.conversationId(),
null, System.currentTimeMillis() - start, true, error.getMessage());
safeSend(emitter, SseEmitter.event().name("error")
.data("{\"message\":\"生成失败,请稍后重试\"}", MediaType.APPLICATION_JSON));
emitter.complete();
},
() -> {
long cost = System.currentTimeMillis() - start;
safeSend(emitter, SseEmitter.event().name("done")
.data("{\"costMillis\":" + cost + "}", MediaType.APPLICATION_JSON));
log.info("[AI-CALL] traceId={} convId={} sync=false cost={}ms",
traceId, request.conversationId(), cost);
emitter.complete();
}
);
return emitter;
}
/** 从流式响应中取出增量文本 */
private String extractDelta(ChatResponse response) {
if (response == null || response.getResult() == null
|| response.getResult().getOutput() == null) {
return null;
}
return response.getResult().getOutput().getText();
}
/** 发送失败不影响主流程 */
private void safeSend(SseEmitter emitter, SseEmitter.SseEventBuilder builder) {
try {
emitter.send(builder);
} catch (Exception e) {
log.debug("[SSE] 事件发送失败(连接可能已关闭)", e);
}
}
}
/**
* 流式分片载荷。
* <p>用单字段对象而不是裸字符串,前端 JSON.parse 之后直接取 .v,
* 后续要加字段(如 messageId、引用来源)也不会破坏协议。</p>
*/
public record StreamChunk(String v) {
public static StreamChunk of(String value) {
return new StreamChunk(value);
}
}
方案 B:直接返回 Flux<String>(代码最短)
@PostMapping(value = "/stream/simple",
produces = MediaType.TEXT_EVENT_STREAM_VALUE + ";charset=UTF-8")
public Flux<String> streamSimple(@Valid @RequestBody ChatRequest request) {
return chatClient.prompt()
.user(request.message())
.tools(orderTools)
.advisors(a -> a.param(ChatMemory.CONVERSATION_ID, request.conversationId()))
.stream()
.content() // 直接给纯文本增量
.map(StreamChunk::of)
.map(JsonUtils::toJson);
}
怎么选? 需要精细化控制超时、心跳、错误事件、Token 统计 → 方案 A;内部工具、Demo、快速验证 → 方案 B。注意:Spring MVC 同样支持返回
Flux,不需要为了流式把整个项目改成 WebFlux。
5.5 知识库:文档入库与 RAG 检索
package com.example.aiagent.service;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.ai.document.Document;
import org.springframework.ai.reader.TextReader;
import org.springframework.ai.transformer.splitter.TokenTextSplitter;
import org.springframework.ai.vectorstore.SearchRequest;
import org.springframework.ai.vectorstore.VectorStore;
import org.springframework.core.io.Resource;
import org.springframework.stereotype.Service;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
/**
* 知识库服务:负责文档切分、向量化入库与检索。
*
* @author 你的名字
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class KnowledgeBaseService {
private final VectorStore vectorStore;
/**
* 将文档切分后写入向量库。
*
* @param resource 文档资源
* @param source 来源标识,用于后续按来源过滤检索
* @param category 业务分类
* @return 切分后的块数
*/
public int ingest(Resource resource, String source, String category) {
List<Document> rawDocs = new TextReader(resource).get();
log.info("[RAG] 文档读取完成 source={} 原始段落={}", source, rawDocs.size());
// 切分策略:单块 800 token、最小块 350 字符,保留分隔符避免句子被截断
TokenTextSplitter splitter = TokenTextSplitter.builder()
.withChunkSize(800)
.withMinChunkSizeChars(350)
.withKeepSeparator(true)
.build();
List<Document> chunks = splitter.apply(rawDocs);
// 关键:给每个 chunk 打元数据,检索时才能按业务维度过滤
chunks.forEach(doc -> {
Map<String, Object> meta = new HashMap<>(doc.getMetadata());
meta.put("source", source);
meta.put("category", category);
meta.put("ingestTime", System.currentTimeMillis());
doc.getMetadata().putAll(meta);
});
vectorStore.add(chunks);
log.info("[RAG] 入库完成 source={} 切分块数={}", source, chunks.size());
return chunks.size();
}
/**
* 按业务维度检索,用于调试与检索效果评估。
*
* <p>filterExpression 使用类 SQL 语法,各向量库之间保持一致的写法,
* 后续从 pgvector 换到 Milvus 不需要改业务代码。</p>
*/
public List<Document> search(String query, String category, int topK) {
return vectorStore.similaritySearch(
SearchRequest.builder()
.query(query)
.topK(topK)
.similarityThreshold(0.55)
.filterExpression("category == '" + category + "'")
.build());
}
}
RAG 落地时最容易被忽略的三件事:
- 元数据比切分参数更重要。生产环境一定要打
source、category、department这类标签,否则多业务知识混在一个库里,检索结果必然串味。 - 切分要带重叠。纯按长度切分会把"第 3 条:自签收之日起 7 日内……"拦腰截断,模型看到半句话就开始编。
- 向量库选型不要一步到位。文档量在百万级以内,PostgreSQL + pgvector 足够,还能和业务库共用一套备份、监控与权限体系。真正需要独立向量库(Milvus 等)的信号是:千万级以上向量、需要频繁重建索引、或者对检索 QPS 有极端要求。
5.6 Token 与耗时审计(线上排障的命脉)
package com.example.aiagent.service;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.ai.chat.model.ChatResponse;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
/**
* AI 调用审计:Token 消耗、耗时、成败落库。
*
* @author 你的名字
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class AiCallLogService {
private final JdbcTemplate jdbcTemplate;
/**
* 记录一次 AI 调用。任何异常都不能影响主流程。
*/
public void record(String traceId, String conversationId, ChatResponse response,
long costMillis, boolean stream, String errorMsg) {
int promptTokens = 0;
int completionTokens = 0;
int totalTokens = 0;
String model = null;
try {
if (response != null && response.getMetadata() != null) {
// 注意:不同厂商返回的 token 统计完整度不同,
// 部分兼容端点不返回 usage,此时保持 0 即可
var usage = response.getMetadata().getUsage();
if (usage != null) {
promptTokens = usage.getPromptTokens() == null ? 0 : usage.getPromptTokens();
completionTokens = usage.getCompletionTokens() == null ? 0 : usage.getCompletionTokens();
totalTokens = usage.getTotalTokens() == null ? 0 : usage.getTotalTokens();
}
model = response.getMetadata().getModel();
}
jdbcTemplate.update("""
INSERT INTO t_ai_call_log
(trace_id, conversation_id, model, prompt_tokens, completion_tokens,
total_tokens, cost_millis, stream, success, error_msg)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
traceId, conversationId, model, promptTokens, completionTokens,
totalTokens, costMillis, stream, errorMsg == null, errorMsg);
log.info("[AI-TOKEN] traceId={} model={} prompt={} completion={} total={} cost={}ms",
traceId, model, promptTokens, completionTokens, totalTokens, costMillis);
} catch (Exception e) {
// 审计失败绝不能拖垮业务
log.error("[AI-TOKEN] 调用记录写入失败 traceId={}", traceId, e);
}
}
}
5.7 全局异常拦截
package com.example.aiagent.config;
import lombok.extern.slf4j.Slf4j;
import org.springframework.http.HttpStatus;
import org.springframework.validation.FieldError;
import org.springframework.web.bind.MethodArgumentNotValidException;
import org.springframework.web.bind.annotation.ExceptionHandler;
import org.springframework.web.bind.annotation.ResponseStatus;
import org.springframework.web.bind.annotation.RestControllerAdvice;
import java.util.stream.Collectors;
/**
* 全局异常处理。
*
* <p>AI 服务最容易出现的是「模型侧异常」(限流、超时、余额不足),
* 这里统一转成可读的业务提示,避免把厂商的错误栈直接暴露给前端。</p>
*
* @author 你的名字
*/
@Slf4j
@RestControllerAdvice
public class GlobalExceptionHandler {
/** 参数校验失败 */
@ExceptionHandler(MethodArgumentNotValidException.class)
@ResponseStatus(HttpStatus.BAD_REQUEST)
public ApiResponse<Void> handleValidation(MethodArgumentNotValidException e) {
String msg = e.getBindingResult().getFieldErrors().stream()
.map(FieldError::getDefaultMessage)
.collect(Collectors.joining(";"));
log.warn("[AI-ERROR] 参数校验失败: {}", msg);
return ApiResponse.fail(40001, msg);
}
/** 业务异常 */
@ExceptionHandler(BizException.class)
public ApiResponse<Void> handleBiz(BizException e) {
log.warn("[AI-ERROR] 业务异常 code={} msg={}", e.getCode(), e.getMessage());
return ApiResponse.fail(e.getCode(), e.getMessage());
}
/** 模型侧不可重试异常:限流 / 配额不足 / 鉴权失败 */
@ExceptionHandler(org.springframework.ai.retry.NonTransientAiException.class)
@ResponseStatus(HttpStatus.TOO_MANY_REQUESTS)
public ApiResponse<Void> handleAiNonTransient(Exception e) {
log.error("[AI-ERROR] 模型调用不可重试异常", e);
return ApiResponse.fail(42901, "AI 服务繁忙,请稍后重试");
}
/** 兜底 */
@ExceptionHandler(Exception.class)
@ResponseStatus(HttpStatus.INTERNAL_SERVER_ERROR)
public ApiResponse<Void> handleOther(Exception e) {
log.error("[AI-ERROR] 未预期异常", e);
return ApiResponse.fail(50000, "系统开小差了,请稍后重试");
}
}
/**
* 统一返回结构。
* <p>code = 0 表示成功,其余为业务错误码,前端只需判断 code。</p>
*/
public record ApiResponse<T>(int code, String message, T data) {
public static <T> ApiResponse<T> ok(T data) {
return new ApiResponse<>(0, "success", data);
}
public static <T> ApiResponse<T> fail(int code, String message) {
return new ApiResponse<>(code, message, null);
}
}
5.8 MCP 接入(零代码挂载外部工具)
MCP(Model Context Protocol)是当前 GitHub 上最热的协议方向。Spring AI 2.x 直接提供了客户端 starter,不用改一行 Java 代码,只要在 application.yml 里配置好连接,外部 MCP Server 的工具就会自动注册进 Agent 的工具列表:
spring:
ai:
mcp:
client:
enabled: true
type: SYNC
request-timeout: 30s
stdio:
connections:
filesystem:
command: npx
args:
- "-y"
- "@modelcontextprotocol/server-filesystem"
- "/data/ai-agent-workspace"
启动后日志里会打印已注册的工具数量,模型就能自主决定是否需要读文件。这意味着团队可以把内部系统封装成独立的 MCP Server,由不同团队各自维护,Agent 侧只改配置。 这是目前 Agent 工程化最务实的一种解耦方式。
提示:如果你用的是 WebFlux(而不是本文的 Spring MVC),MCP 客户端需要换成 WebFlux 版 starter:
spring-ai-starter-mcp-client-webflux。
六、前端联调
6.1 TypeScript 类型定义
/** 与后端 ChatRequest 一一对应 */
export interface ChatRequest {
conversationId: string;
message: string;
}
export interface ChatResponse {
code: number;
message: string;
data: {
conversationId: string;
content: string;
referenceCount: number;
promptTokens: number;
completionTokens: number;
costMillis: number;
};
}
/** 流式分片 */
export interface StreamChunk {
v: string;
}
export interface StreamDonePayload {
costMillis: number;
}
6.2 前端消费 SSE
注意:浏览器原生 EventSource 只支持 GET,无法发送 JSON body,所以这里用 fetch + ReadableStream。
export async function chatStream(
req: ChatRequest,
onDelta: (text: string) => void,
onDone: (costMillis: number) => void,
onError: (msg: string) => void,
): Promise<void> {
const controller = new AbortController();
try {
const resp = await fetch('/api/v1/chat/stream', {
method: 'POST',
headers: {
'Content-Type': 'application/json',
Accept: 'text/event-stream',
},
body: JSON.stringify(req),
signal: controller.signal,
});
if (!resp.ok || !resp.body) {
throw new Error(`HTTP ${resp.status}`);
}
const reader = resp.body.getReader();
const decoder = new TextDecoder('utf-8');
let buffer = '';
// SSE 规范:事件之间用空行分隔,所以要按 \n\n 切分并保留不完整片段
while (true) {
const { value, done } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
const blocks = buffer.split('\n\n');
buffer = blocks.pop() ?? '';
for (const block of blocks) {
const eventName = /^event:\s*(.+)$/m.exec(block)?.[1]?.trim();
const dataLine = /^data:\s*(.+)$/m.exec(block)?.[1];
if (!dataLine) continue;
if (eventName === 'message') {
const chunk = JSON.parse(dataLine) as StreamChunk;
onDelta(chunk.v);
} else if (eventName === 'done') {
const payload = JSON.parse(dataLine) as StreamDonePayload;
onDone(payload.costMillis);
} else if (eventName === 'error') {
onError(JSON.parse(dataLine).message);
}
}
}
} catch (e) {
onError((e as Error).message);
}
}
6.3 Nginx 配置(流式接口不生效,90% 是这里)
location /api/v1/chat/stream {
proxy_pass http://ai_agent_backend;
proxy_http_version 1.1;
# ★ 关键三行:关闭缓冲,否则会攒够一批才吐给浏览器
proxy_buffering off;
proxy_cache off;
proxy_set_header X-Accel-Buffering no;
# ★ 长连接超时,必须大于后端 SSE 超时
proxy_read_timeout 600s;
proxy_send_timeout 600s;
# 禁用 gzip:压缩会破坏流式实时性
gzip off;
chunked_transfer_encoding on;
proxy_set_header Connection '';
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
}
七、Docker Compose 一键部署
services:
postgres:
image: pgvector/pgvector:pg17
container_name: ai-agent-db
environment:
POSTGRES_DB: aiagent
POSTGRES_USER: postgres
POSTGRES_PASSWORD: ${DB_PASSWORD:-postgres}
volumes:
- ./data/postgres:/var/lib/postgresql/data
ports:
- "5432:5432"
healthcheck:
test: ["CMD-SHELL", "pg_isready -U postgres -d aiagent"]
interval: 10s
timeout: 5s
retries: 5
app:
build: .
container_name: ai-agent-app
depends_on:
postgres:
condition: service_healthy
environment:
DB_PASSWORD: ${DB_PASSWORD:-postgres}
LLM_API_KEY: ${LLM_API_KEY}
EMBEDDING_API_KEY: ${EMBEDDING_API_KEY}
SPRING_PROFILES_ACTIVE: prod
# 容器内连接数据库要改成服务名
SPRING_DATASOURCE_URL: jdbc:postgresql://postgres:5432/aiagent
ports:
- "8080:8080"
restart: unless-stopped
# 需要固定公网访问时,再挂一个 frpc / cpolar 容器即可
FROM eclipse-temurin:17-jre-alpine
WORKDIR /app
COPY target/ai-agent-service-1.0.0.jar app.jar
ENV JAVA_OPTS="-Xms512m -Xmx1g -XX:+UseG1GC -Duser.timezone=Asia/Shanghai"
EXPOSE 8080
ENTRYPOINT ["sh", "-c", "java $JAVA_OPTS -jar app.jar"]
docker compose up -d
docker compose logs -f app
八、生产环境踩坑清单(建议收藏)
8.1 版本与配置类
spring-boot-starter-web在 Boot 4 已 deprecated —— 换spring-boot-starter-webmvc,否则后续大版本会直接编译失败。spring.ai.openai.chat.enabled在 Spring AI 2.x 已移除 —— 改用spring.ai.model.chat。旧属性残留不会报错,但模型 Bean 会静默失效。- starter 命名规则变了 —— 现在统一是
spring-ai-starter-model-xxx/spring-ai-starter-vector-store-xxx。 QuestionAnswerAdvisor在独立制品里 —— 必须额外引入spring-ai-vector-store-advisor。dimensions必须和 embedding 模型对齐 —— 表结构建好后改维度需要重建表,上线前一定确认。
8.2 流式相关
proxy_buffering off+X-Accel-Buffering no—— 不配就是"转圈十秒然后整段蹦出来"。- SSE 超时要三层对齐:Nginx
proxy_read_timeout> 后端SseEmitter超时 > 模型单次超时。 - 必须处理客户端主动断开 ——
emitter.send()抛IOException是常态,要complete()而不是打印 error 堆栈。 - 中间插入心跳 —— 长耗时的工具调用期间没有输出,部分网关会判定空闲断开。可以每 15 秒发一个注释行
:heartbeat。 - 不要在 Controller 里
blockLast()—— 会把 Tomcat 线程池拖垮。
8.3 Agent 与检索
- 工具方法自己兜异常 —— 一次工具报错会中断整个调用循环,用户体验断崖式下跌。
- 给 RAG 设
similarityThreshold—— 只有topK等于"不管多不相干都硬塞",是幻觉最大的来源。 - 工具数量控制在 20 个以内 —— 工具越多,模型选错的概率越高。超过之后应该按业务域拆分 Agent 或引入路由。
- 会话 ID 必须由服务端从登录态派生,不要信前端传什么就用什么 —— 否则存在越权读取他人会话的风险。
8.4 成本与稳定性
- 必须有 Token 审计表 —— 没有它,成本失控和效果变差都无从定位。
- 给单用户做限流 —— AI 接口的调用成本是普通接口的几百倍,没有限流等于对外裸奔。
- 日志脱敏 —— Prompt 里经常带手机号、订单号、内部地址,落库前要过一遍脱敏。
- 给模型调用加超时和降级 —— 模型服务抖动时,返回"服务繁忙"远好于把请求线程全部挂死。
九、小结
把这次的关键点压缩成四句话:
- AI 应用层的热度已经从"模型"转移到"Agent 工程化",GitHub 周榜/月榜超过一半的仓库都属于这一类。
- Java 侧存在明确的内容供给缺口,而 Spring AI 2.x 已经把生产需要的抽象(Advisor、工具调用、记忆、RAG、MCP)全部做好了。
- "能上线的 AI 服务"和"能跑通的 Demo"差别在五件事上:Token 审计、异常兜底、连接治理、检索阈值、限流降级。
- 这套架构的价值在于可演进:向量库能换、模型能换、工具能通过 MCP 由别的团队维护,而业务代码不用动。
源码获取
完整工程(含 pom.xml、application.yml、建表 SQL、Dockerfile、docker-compose.yml)已经整理好,评论区回复「AI Agent」我直接发你,另外附一份可导入的 Postman 集合。
下一篇预告(已经排好,欢迎投票)
- 《MCP 实战:把公司内部系统封装成 MCP Server,让 Agent 直接调用》 —— 把第 5.8 节展开成完整工程
- 《RAG 检索效果调优:从 60 分到 90 分,我改了这 7 个参数》 —— 切分策略、混合检索、重排序的实测对比
- 《pgvector 还是 Milvus?我用 100 万条真实语料压测了一遍》 —— 性能、成本、运维复杂度三维对比
想看哪篇?评论区告诉我序号,票数最高的先写。 如果这套架构对你有帮助,点个赞、收个藏,关注我,下一篇第一时间推送给你。有任何实现上的问题,评论区直接问,我都会逐条回复。
本文基于 Spring Boot 4.1.1 / Spring AI 2.0.1 编写,文中所有 API 与配置项均对照官方文档核实。技术版本迭代较快,落地前建议再确认一次当前版本。
转载自 CSDN-专业IT技术社区
原文链接:https://blog.csdn.net/wkt1105436760/article/details/166792200




