Spring AI 本质上是 Spring 生态里的一套开发框架,它帮我们解决的问题从“访问数据库、访问 Redis、访问 MQ”,变成了“访问大模型”。
Spring AI 官方对它的定位可以理解为三句话:
Spring AI项目的目标,是让开发者可以更简单地开发包含人工智能能力的应用,而不需要引入不必要的复杂度。- 如果大家了解 Python 生态里的
LangChain,可以把 Spring AI 理解为 Java / Spring 生态中对标 LangChain 的一套 AI 应用开发框架。 - Spring AI 要解决的核心问题是:
Connecting your enterprise Data and APIs with AI Models,也就是把企业自己的数据、业务接口和 AI 模型连接起来。

使用Spring AI 所需的依赖和配置如下:
```
<properties>
<spring-ai.version>1.0.0</spring-ai.version>
</properties>
<dependencyManagement>
<dependencies>
<!-- Spring AI BOM:统一管理 Spring AI 相关依赖版本 -->
<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>
<!-- Spring AI OpenAI 兼容接口依赖 -->
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-starter-model-openai</artifactId>
</dependency>
</dependencies>
```
这里的 spring-ai-openai-spring-boot-starter 不一定只能访问 OpenAI 官方模型。很多国内模型平台、本地模型网关、公司内部模型服务,只要提供的是 OpenAI 兼容接口,通常都可以使用这一套配置方式接入。
```
spring:
ai:
openai:
# 模型接口地址。可以是 OpenAI 官方地址,也可以是兼容 OpenAI 协议的第三方平台地址
base-url: https://dashscope.aliyuncs.com/compatible-mode
# 模型平台提供的访问密钥,用来证明“你是谁、有没有权限调用模型”
api-key: your-api-key
chat:
options:
# 具体使用哪个聊天模型,例如 gpt-4o-mini、qwen-plus、deepseek-chat 等
model: qwen-plus
```
通常会把API-KEY配置到环境变量里面,配置yml的时候写api-key:${API-KEY}
这里重点解释两个配置:
| 配置 | 含义 | 从哪里获取 |
|---|---|---|
base-url | 模型服务的接口地址 | 模型平台控制台、公司内部模型网关文档、本地模型服务地址 |
api-key | 调用模型时使用的密钥 | 模型平台控制台创建,或由公司统一分配 |
大家可以类比以前连数据库:数据库要有 url,模型也要有 base-url;数据库要有用户名密码,模型也要有 api-key;数据库要指定库名和驱动,模型也要指定 model。
Prompt
在使用Spring AI 前,接下来我们先讲 Prompt,也就是提示词。因为我们使用模型时,模型到底怎么回答,很大程度由提示词决定。
这里要强调一下,本课程中所说的“模型”,主要指的是大语言模型,也就是 LLM。大语言模型有几个非常明显的特点:
- 它擅长理解自然语言,也擅长生成文本、总结文本、改写文本、抽取信息。
- 它不是数据库,不保证每次都严格返回固定结果。
- 它不是业务系统,不天然知道我们公司的数据和接口。
- 它的输出会受到上下文和提示词的强烈影响。
所以,我们和大模型交互时,不能只问一句“帮我做一下”,而要把任务、背景、要求、输出格式尽量讲清楚。
Message 角色
现在的大模型提示词通常不是一整段纯文本,而是由多条 Message 组成。每条 Message 都有自己的角色。
| 角色 | 作用 | 类比 |
|---|---|---|
System | 设定 AI 的人设、行为准则、背景知识和回答边界 | 幕后导演 |
User | 用户提出的问题、指令或需求 | 提问的人 |
Assistant | 大模型之前的回复,常用于维持多轮对话上下文 | 聊天记录里的 AI 回复 |
Tool / Function | 工具执行后的结果,用来告诉模型工具返回了什么 | 外部系统查询结果 |
比如:
``` System: 你是一名经验丰富的 Java 讲师,回答要适合初学者。 User: 请解释一下什么是 Spring Bean。 Assistant: Spring Bean 就是交给 Spring 容器管理的对象…… User: 那它和 new 出来的对象有什么区别? ```
这里每一行都可以看作一条消息,多条消息组合起来,才构成一次完整的 Prompt。
Prompt 核心要点
Spring AI 官方文档中提到,设计 Prompt 时通常需要考虑几个组成部分。
Instructions:明确指令。告诉 AI 要做什么,就像我们给一个同事安排任务一样,越清晰越好。
External Context:外部上下文。把相关背景、资料、约束条件补充进去,让 AI 知道当前任务发生在什么场景里。User Input:用户输入。也就是用户真正提出的问题或需求。Output Indicator:输出格式要求。告诉 AI 希望以什么格式返回,比如 JSON、表格、列表等。 // 拿到的是对象 “姓名张三,”
如果我们只是问“帮我分析一下这个病人”,模型不知道病人信息在哪里,也不知道分析目标是什么。更好的 Prompt 是:
``` 你是一名智能分诊助手。 请根据下面的患者描述,判断患者可能需要挂哪个科室。 要求: 1. 只能从 内科、外科、儿科、急诊科 中选择。 2. 给出 1 句简短理由。 3. 使用 JSON 格式返回。 患者描述:发热 3 天,咳嗽,咽痛,无明显胸痛。 ```
这个 Prompt 就明显更容易得到稳定结果。
简单一轮对话
真正访问大模型时,Spring AI 中最核心的对象之一就是 ChatClient。大家可以把 ChatClient 类比成以前使用的 RedissionClient 或者 RestTemplate:RestTemplate 用来访问普通 HTTP 服务,ChatClient 用来访问聊天模型。
我们先来看下如何得到一个ChatClient对象
```
import org.springframework.ai.chat.client.ChatClient;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class SpringAIConfig {
/**
* 创建 ChatClient 对象。
* ChatClient.Builder 由 Spring AI 自动提供,它会读取 application.yml 中的模型配置。
*/
@Bean
public ChatClient chatClient(ChatClient.Builder builder) {
return builder.build();
}
}
```
接着定义一个 Controller 测试最简单的一轮对话:
```
import org.springframework.ai.chat.client.ChatClient;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
@RestController
public class ChatController {
private final ChatClient chatClient;
public ChatController(ChatClient chatClient) {
this.chatClient = chatClient;
}
@GetMapping("/ai")
public String generation(@RequestParam String userInput) {
return this.chatClient
.prompt()
.user(userInput)
.call()
.content();
}
}
```
这段代码的调用链非常重要:chatClient.prompt().user(userInput).call().content()。
prompt():开始构建一次提示词请求。user(userInput):添加一条用户消息。call():向模型发起同步调用。content():取出模型返回的文本内容。
写完后,可以使用 ApiFox 或浏览器访问:
``` GET http://localhost:8080/ai?userInput=请用一句话介绍Spring AI ```
系统消息
接下来,我们在上一节的基础上再加入系统消息。在ChatClient中添加系统消息有三种方式:
- 局部添加系统消息
- 添加默认系统消息
- 使用动态参数
局部系统消息只对当前这一次调用生效。
```
@GetMapping("/ai/java-teacher")
public String javaTeacher(@RequestParam String question) {
return chatClient.prompt()
.system("你是一名经验丰富的Java讲师,回答要适合初学者,尽量使用大白话。")
.user(question)
.call()
.content();
}
```
如果一个 ChatClient 在整个项目中都希望使用同一套系统设定,可以在构建 ChatClient 时设置默认系统消息。
```
@Bean
public ChatClient javaTeacherChatClient(ChatClient.Builder builder) {
return builder
.defaultSystem("你是一名经验丰富的Java讲师,回答要清晰、准确、适合初学者。")
.build();
}
```
有时候系统消息中有一部分内容是动态的。此时,我们可以在系统消息中添加占位符,在执行时动态传递参数填充系统 消息中的占位符
```
@Bean
public ChatClient dynamicTeacherChatClient(ChatClient.Builder builder) {
return builder
// 这里注意,默认情况下占位符的格式是{占位符名称}
.defaultSystem("你是一名{language}讲师,请使用{style}的方式回答问题。")
.build();
}
```
```
@GetMapping("/ai/dynamic-teacher")
public String dynamicTeacher(@RequestParam String question,
@RequestParam(defaultValue = "Java") String language,
@RequestParam(defaultValue = "大白话") String style) {
return chatClient.prompt()
.system(prompt -> prompt
.param("language", language)
.param("style", style))
.user(question)
.call()
.content();
}
```
注意:如果系统消息中写了 {language}、{style} 这种参数占位符,但是调用时没有提供对应参数,Spring AI 会在模板渲染阶段报错。
结构化输出
在和模型交互的过程中,模型给我们返回的通常都是普通字符串。但是在真实项目中,我们经常不希望模型返回一段自然语言,而是希望它返回结构化的数据。
比如让模型生成某个演员的电影作品列表,如果只返回一段普通文本,后端还要自己解析;如果直接返回 Java 对象,后续代码处理就会方便很多。
Spring AI 的 ChatClient 提供了一个非常直接的方法:
``` .call() .entity(目标类型) ```
这里我们不展开讲各种 xxxConverter,只围绕官方文档中 entity() 的用法看三种情况:普通对象、List、Map。
普通对象
官方文档中第一个例子是:生成某个演员的 5 部电影作品,并把结果映射成一个对象。
先定义一个普通 Java Bean。官方文档里使用的是 record,但我们课程中为了和大家前面学习的 Java Bean 保持一致,这里改成普通类。
```
@Data
public class ActorsFilms {
private String actor;
private List<String> movies;
}
```
然后使用 entity(ActorsFilms.class) 接收结果:
```
import org.springframework.ai.chat.client.ChatClient;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
@RestController
public class StructuredOutputController {
private final ChatClient chatClient;
public StructuredOutputController(ChatClient chatClient) {
this.chatClient = chatClient;
}
/**
* 示例请求:GET /ai/actor-films?actor=Tom Hanks
*/
@GetMapping("/ai/actor-films")
public ActorsFilms actorFilms(@RequestParam(defaultValue = "Tom Hanks") String actor) {
return chatClient.prompt()
.user(user -> user
.text("Generate the filmography of 5 movies for {actor}.")
.param("actor", actor))
.call()
.entity(ActorsFilms.class);
}
}
```
这里最关键的是:
``` .entity(ActorsFilms.class) ```
它表示让 Spring AI 尝试把模型输出映射成一个 ActorsFilms 对象。
List 对象
第二种情况:如果最外层返回值本身就是一个 List,比如一次生成 Tom Hanks 和 Bill Murray 两个演员的电影作品,就需要使用 ParameterizedTypeReference。
```
import org.springframework.core.ParameterizedTypeReference;
import org.springframework.web.bind.annotation.GetMapping;
import java.util.List;
@GetMapping("/ai/actor-films-list")
public List<ActorsFilms> actorFilmsList() {
return chatClient.prompt()
.user("Generate the filmography of 5 movies for Tom Hanks and Bill Murray.")
.call()
.entity(new ParameterizedTypeReference<List<ActorsFilms>>() {
});
}
```
这里的重点是:
```
.entity(new ParameterizedTypeReference<List<ActorsFilms>>() {
})
```
为什么不能直接写 List.class?
因为 Java 泛型存在类型擦除。如果只写 List.class,程序只能知道结果是一个 List,但不知道 List 里面的元素到底是什么类型。所以,当最外层结果是 List 时,要用 ParameterizedTypeReference> 把完整泛型类型告诉 Spring AI。
Map 对象
第三种情况:如果我们不想提前定义 Java Bean,也可以让模型返回一个 Map。
下面这个例子:让模型介绍一座城市,返回一个包含若干字段的 JSON 对象。因为城市信息的字段并不固定(有的城市可能多一个字段,有的少一个),用 Map 接收最省事。
```
import org.springframework.core.ParameterizedTypeReference;
import org.springframework.web.bind.annotation.GetMapping;
import java.util.Map;
@GetMapping("/ai/city-info")
public Map<String, Object> cityInfo() {
return chatClient.prompt()
.user("用 JSON 介绍一下北京,包含城市名称、所属国家、人口数量这几个字段")
.call()
.entity(new ParameterizedTypeReference<Map<String, Object>>() {
});
}
```
这里的重点是:
```
.entity(new ParameterizedTypeReference<Map<String, Object>>() {
})
```
Map 的特点是灵活,不需要提前定义类;但是缺点也很明显:字段名和字段类型没有 Java 类约束,后续取值时更容易写错,也可能需要手动类型转换。
所以实际项目中可以这样选择:
| 返回类型 | 写法 | 适合场景 |
|---|---|---|
| 普通对象 | .entity(ActorsFilms.class) | 最外层是一个结构稳定的对象 |
List<对象> | .entity(new ParameterizedTypeReference>() {}) | 最外层是一组同类型对象 |
Map | .entity(new ParameterizedTypeReference>() {}) | 临时验证、字段不稳定、结构比较动态 |
最后再强调一次对象中包含 List 的情况:
ActorsFilms里面有List movies,但最外层是ActorsFilms,所以用.entity(ActorsFilms.class)。- 只有当最外层本身就是
List时,才使用ParameterizedTypeReference>。
多轮对话
前面的接口都是一问一答。用户问一句,模型答一句。问题是:真实聊天通常是多论对话。
比如:
``` 用户:我发烧三天了。 模型:还有咳嗽、咽痛吗? 用户:有咳嗽。 ```
第三句话“有咳嗽”本身没有完整含义。如果模型不知道前面聊过“发烧三天”,就无法正确理解,但是我们知道模型只有短期记忆。
所以,多轮对话的核是让模型在当前请求中看到之前的聊天记录,而这就是会话记忆。实现会话记忆,本质上要解决三个问题:
- 存在哪?
聊天记录要存在哪里?Spring AI 中使用 ChatMemoryRepository接口来定义会话记忆的保存行为,根据具体存储位置的不同,比如内存、JDBC、Cassandra、Neo4j 等,Spring AI 提供了ChatMemoryRepository接口的不同实现。
- 存多少?
聊天记录不能无限制全部塞给模型。原因很简单:模型上下文长度有限,聊天记录越多,请求越慢、费用越高,太久之前的消息可能也没有价值。
Spring AI 中 ChatMemory 接口表示会话记忆,它会决定如何管理消息,比如保留最近 N 条消息。
可以这样理解二者关系:存储组件负责“消息放在哪里”,ChatMemory 负责“怎么读、怎么写、保留多少”。
- 如何隔离不同会话?
不同用户、不同窗口、不同业务单据的聊天记录必须隔离。核心思路就是使用 conversationId。
比如用户 A 的会话 ID 是 user-a-session-001,用户 B 的会话 ID 是 user-b-session-001。每次调用模型时带上对应的会话 ID,Spring AI 就能把不同会话的记忆区分开。
会话记忆的实现
Spring AI 实际使用中,不我们每次手动调用方法保存会话中的消息,也不用把历史消息查出来再拼到 Prompt 里发送给模型。而是通过内置的Advisor自动保存会话中的消息,自动拼接会话中的消息到提示词中。
Advisor 可以理解为 ChatClient 调用过程中的拦截器,它们的执行原理如下:
- 用户代码通过
ChatClient发起请求。 - 请求不会立刻到模型,而是先经过 Advisor 链。
- Advisor 可以在请求发送前修改 Prompt,比如加入历史消息、加入检索到的知识、记录日志。
- 修改后的请求发送给模型。
- 模型返回响应。
- 响应也会经过 Advisor 链,Advisor 可以记录响应日志、保存消息等。
- 最终结果返回给 Controller。
这和大家熟悉的拦截器、过滤器思想非常像:Spring MVC 的拦截器可以在 Controller 前后做增强,Advisor 可以在大模型调用前后做增强。

所以,Advisor在ChatClient发消息前在提示词中添加会话历史消息,并保存用户发送的用户消息等。在收到模型响应时自动保存模型返回的助手消息。如果想要实现会话记忆,只需要使用会话记忆相关的Advisor即可。目前针对会话记忆主要使用的是MessageChatMemoryAdvisor
MessageChatMemoryAdvisor 的使用非常简单,代码如下:先配置一个带记忆能力的 ChatClient,
```
import org.springframework.ai.chat.client.ChatClient;
import org.springframework.ai.chat.client.advisor.MessageChatMemoryAdvisor;
import org.springframework.ai.chat.memory.ChatMemory;
import org.springframework.ai.chat.memory.InMemoryChatMemoryRepository;
import org.springframework.ai.chat.memory.MessageWindowChatMemory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class ChatMemoryConfig {
@Bean
public ChatMemory chatMemory() {
return MessageWindowChatMemory.builder()
// 指定ChatMemory的底层存储使用纯内存存储
.chatMemoryRepository(new InMemoryChatMemoryRepository())
.maxMessages(20)
.build();
}
@Bean
public ChatClient memoryChatClient(ChatClient.Builder builder, ChatMemory chatMemory) {
return builder
// 指定默认的advisor,当然也可以局部指定就像局部指定系统消息一样
.defaultAdvisors(MessageChatMemoryAdvisor.builder(chatMemory).build())
.build();
}
}
```
然后定义 Controller:
```
import org.springframework.ai.chat.client.ChatClient;
import org.springframework.ai.chat.memory.ChatMemory;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
@RestController
public class MemoryChatController {
private final ChatClient memoryChatClient;
public MemoryChatController(ChatClient memoryChatClient) {
this.memoryChatClient = memoryChatClient;
}
/**
* 示例请求:
* GET /ai/memory?conversationId=1001&message=我叫张三
* GET /ai/memory?conversationId=1001&message=我叫什么名字
*/
@GetMapping("/ai/memory")
public String memory(@RequestParam String conversationId,
@RequestParam String message) {
return memoryChatClient.prompt()
.user(message)
// 实现会话记忆的隔离,指定会话id
.advisors(advisor -> advisor
.param(ChatMemory.CONVERSATION_ID, conversationId))
.call()
.content();
}
}
```
测试时要注意:同一个 conversationId 才会共享记忆。换一个 conversationId,模型就看不到之前的聊天记录了。
SimpleLoggerAdvisor
在引入一个新的Advisor,SimpleLoggerAdvisor 它可以帮助我们打印请求和响应日志。
```
import org.springframework.ai.chat.client.advisor.SimpleLoggerAdvisor;
@Bean
public ChatClient loggerChatClient(ChatClient.Builder builder, ChatMemory chatMemory) {
return builder
.defaultAdvisors(new SimpleLoggerAdvisor())
.build();
}
```
它的作用不是改变模型能力,而是方便我们观察最终发送给模型的请求是什么、模型返回的响应是什么、Advisor 链是否按预期工作。
流式响应
如果仔细观察我们访问大模型的结果,就会发现我们每次获取的都是大模型完整的回复内容。但是大家回忆一下平时我们在使用DeepSeek,ChatGPT等大模型产品时,每次的回复是一点一点显示出来,而不是一次就显示完整的回复。
这种将响应内容分成多个数据块(chunk)逐步发送给客户端,客户端可以边接收边处理,而无需等待一次接收整个响应内容的响应数据传输方式,我们称之为流式响应。
之所以大模型产品普遍采用流式响应,是因为:
- LLM 是逐个 token 生成文本的,如果大模型回复的文本内容比较多,那么用户可能需要等待一个未知的时长才能看到完整响应
- 但是使用流式响应可以改善用户体验,用户几乎可以立即开始阅读大模型的响应,它大大缩短了用户所看到的首次响应的时间
SSE 协议
简单思考一下就会发现,流式响应的实现和普通的HTTP响应一定是不一样的,因为:
- 普通的HTTP响应式一次请求,对应一次响应,一次通信过程结束
- 但是,流式响应式一次请求,对应多次响应
流式响应的实现需要遵循一套特殊的协议 SSE,全称 Server-Sent Events:
- SSE协议基于 HTTP
- 允许服务端向客户端持续推送事件流,每一个事件流中可以包含事件类型和数据。
- 每个事件必须以两个换行符(\n\n)结尾,它表示“一个 chunk”和“下一个chunk”的物理边界
- 每个事件中的数据就代表流式响应中的一个chunk,事件类型可以作为逻辑标识,标志流式响应的结束。但是流式响应中事件类型是可选的,如果流式响应中没有事件类型也没关系,服务端在流逝响应结束后可以物理断开连接
SSE 返回给前端的数据格式大致如下:
``` id: 1 event: message data: 你好 id: 2 event: message data: ,我是 id: 3 event: message data: Spring AI ```
| 字段 | 是否必须 | 说明 |
|---|---|---|
data | 必须 | 当前事件携带的数据 |
event | 可选 | 事件类型,前端可按类型监听 |
id | 可选 | 事件 ID,可用于断线重连 |
retry | 可选 | 客户端重连间隔 |
Spring AI 流式接口
后端可以返回 Flux:
```
import org.springframework.ai.chat.client.ChatClient;
import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import reactor.core.publisher.Flux;
@RestController
public class StreamChatController {
private final ChatClient chatClient;
public StreamChatController(ChatClient chatClient) {
this.chatClient = chatClient;
}
@GetMapping(value = "/ai/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> stream(@RequestParam String message) {
return chatClient.prompt()
.user(message)
// 调用call获取的是整个响应,而调用stream获取的是流式响应
.stream()
// 获取流式响应中的chunk数据
.content()
.doOnNext(chunk -> System.out.println("模型片段:" + chunk))
// 流式响应的末尾拼接字符串"[END]"作为自定义结束标志
.concatWith(Flux.just("[END]"));
}
}
```
这里重点看三个点:stream() 表示使用流式调用,content() 把模型输出转换成文本片段流,concatWith(Flux.just("[END]")) 在流结束后追加结束标志。
让模型使用你的知识
到目前为止,模型主要依赖它自己训练时学到的知识来回答问题。但是在真实业务里,光靠这点知识远远不够,因为
- 不了解最新动态。模型训练完成之后再发生的事情,它根本没”见过”。比如 Spring AI 官方就提到,GPT-3.5/4 的训练数据只到 2021 年 9 月,问它之后的新闻、新政策、新文档,它只会告诉你”不知道”。
- 不了解你们公司。企业内部的规章制度、项目文档、订单数据、私有业务规则,不可能出现在模型的训练数据里,模型自然也答不上来。
- 不会调用外部系统。比如查当前数据库里的库存、调支付接口、读实时天气,这些都不是脑子里”记住”的知识,而是需要主动去外部系统里拿。
所以光让模型”自己想”是不够的,我们必须想办法把自己的数据和外部 API 接到模型上去用。
针对这件事,Spring AI 官方在给出了三条主流的技术路线:
| 方式 | 核心思想 | 适合场景 |
|---|---|---|
| Fine Tuning(微调) | 拿你的数据去重新训练模型,改动模型内部的权重 | 想固定回答风格、专精某类任务,且团队有足够的算力和样本 |
| Prompt Stuffing(提示词塞数据) | 不动模型,回答前直接把相关资料塞进 Prompt,让模型”边看边答”。其典型工程化实现就是 RAG | 企业知识库问答、文档问答(最常用) |
| Tool Calling(工具调用) | 不动模型,给模型一份”工具清单”,让它在需要时主动调用我们写好的服务/接口,去拿实时数据或执行操作 | 调用外部 API、查询数据库、操作业务系统 |
简单对比一下这三条路线:
Fine Tuning是去改造模型本身,效果最深入,但成本最高,对大多数业务团队来说门槛过高。Prompt Stuffing / RAG不改模型,只是答题前先”翻一下小抄”,是企业知识库场景下最常用的方案。Tool Calling同样不改模型,但相当于额外给了模型”动手”的能力,让它能主动调外部系统拿数据、做事情。
后两种方式都不需要重新训练模型,部署和迭代都很轻量,是企业落地最常见的两种方式。
RAG
RAG 全称是 Retrieval Augmented Generation,翻译成中文叫“检索增强生成”。
简单理解:RAG 就是给大模型配一个外挂知识库。模型回答前,先去知识库查资料,再结合查到的资料生成答案。
可以类比考试:大模型像一个学习能力很强的学生,但它不能记住所有企业内部资料。RAG 就像允许它在答题前翻我们准备好的“资料小抄”。
RAG 能解决几个核心问题:
- 解决大模型知识过期问题,知识库可以随时更新。
- 让回答更可追溯,答案可以对应到原始文档。
- 无需重新训练模型,也能让 AI 掌握企业内部知识。
RAG 基本原理如下:

流程可以拆成三步:
- 用户提问后,先根据问题去知识库检索相关内容。
- 将检索到的内容和用户问题一起放进 Prompt。
- 大模型参考这些内容生成最终回答。
接下来,我们以做一个”企业员工手册问答助手”为例,构建一个员工手册知识库,让用户可以询问考勤、请假、报销、远程办公等制度问题。
知识的表示
企业知识可能来自各种文档:HTML、PDF、Word、Markdown、TXT 等。那是不是直接把这些文档原文存起来就行?当然不行。
如果直接存原文,检索时大概率只能做关键词匹配。关键词匹配有两个明显问题:
词汇鸿沟:用户说“在家办公”,文档写“远程办公”,关键词可能匹配不到。缺乏语义理解:只看词有没有出现,不理解句子真正含义。
所以 RAG 通常会把文本转成向量。
向量可以理解为文本在数学空间中的坐标:一段文本可以表示成一个向量,含义相近的文本向量距离更近,含义差异大的文本向量距离更远。
把文本转成向量的过程叫 Embedding。实现 Embedding 的模型叫 EmbeddingModel。我们把文本输入给 EmbeddingModel,它输出对应向量。

总结一句:为了更好地按语义检索知识,我们需要先把知识文本通过 EmbeddingModel 转成向量。
这些向量要存在哪里?就需要向量数据库。常见向量数据库包括 Milvus、Chroma、Pinecone、RedisSearch、Elasticsearch。本课程中我们使用 Elasticsearch 作为向量存储。
知识库的构建
知识库构建本质上是一个 ETL 流程。
| 阶段 | 英文 | 作用 |
|---|---|---|
| E | Extract | 读取原始文档 |
| T | Transform | 清洗、切分、转换文档 |
| L | Load | 写入目标存储,比如向量数据库 |
RAG 的 ETL 流程通常是离线处理的。也就是说,并不是用户每次提问时才临时读取所有文档、切分、向量化、写入数据库。那样太慢了。
正确流程是:提前把企业文档读取出来,切分成适合检索的小段,调用 EmbeddingModel 转成向量,写入向量数据库,用户提问时只做检索和生成。
文档读取
Spring AI 中读取文档的核心接口可以简化理解为 DocumentReader。
| Reader | 适合读取 |
|---|---|
TextReader | TXT、普通文本 |
JsonReader | JSON 文件 |
MarkdownDocumentReader | Markdown 文档 |
PagePdfDocumentReader | PDF 文档 |
TikaDocumentReader | 多种复杂文件格式 |
本课程中我们使用 TextReader,因为它最简单,适合先理解完整流程。
接下来我们把一份”员工手册节选”做成一个文档文件。在 src/main/resources/document 目录下新建 employee-handbook.txt,内容如下(节选,可按需扩充):
``` 员工手册节选 考勤与工作时间:公司实行标准工作制,工作时间为周一至周五 9:00 至 18:00,中午 12:00 至 13:30 为午休时间。员工因交通、天气或其他原因无法按时到岗时,应提前通过企业微信向直属主管说明情况。未经说明且超过规定时间未到岗的,按迟到或旷工处理。 请假与审批:员工请假应提前提交申请,并说明请假类型、起止时间和请假原因。病假需在返岗后补充病假证明或就诊记录。年假、事假、调休等申请由直属主管审批。连续请假超过三天的,还需要部门负责人确认。 报销与发票:员工因公产生的交通、住宿、会议等费用,可以在费用发生后 30 天内提交报销申请。申请时应上传真实、完整、合规的发票和费用说明。报销单据应与实际业务一致。个人消费、无关招待、无法说明业务用途的费用,不属于公司报销范围。 远程办公:员工因特殊情况需要远程办公时,应提前向直属主管申请,并说明远程办公日期、工作安排和联系方式。远程办公期间,员工应保持工作时间内在线,及时响应会议、消息和任务协作。涉及客户资料、合同、财务数据等敏感信息时,应使用公司批准的设备和网络环境。 信息安全:员工不得将公司账号、密码、验证码提供给他人,也不得使用私人网盘保存公司内部资料。离职、转岗或项目结束时,应按要求归还设备、移交文档,并删除个人设备中保存的公司敏感信息。 ```
读取代码如下:
```
import org.springframework.ai.document.Document;
import org.springframework.ai.reader.TextReader;
import org.springframework.core.io.Resource;
import org.springframework.core.io.ResourceLoader;
import org.springframework.stereotype.Service;
import java.util.List;
@Service
public class DocumentReadService {
private final ResourceLoader resourceLoader;
public DocumentReadService(ResourceLoader resourceLoader) {
this.resourceLoader = resourceLoader;
}
public List<Document> readEmployeeHandbook() {
Resource resource = resourceLoader.getResource("classpath:document/employee-handbook.txt");
TextReader textReader = new TextReader(resource);
textReader.getCustomMetadata().put("source", "employee-handbook.txt");
return textReader.get();
}
}
```
Document 可以理解为 Spring AI 对”知识片段”的统一抽象。它通常包含 text 文本内容和 metadata 元数据,比如来源文件、页码、分类、时间等。在我们这个案例里,TextReader 读出来后就是一个 Document,text 就是整份员工手册节选,metadata 里带上了 source=employee-handbook.txt,方便后面追溯每段答案来自哪个文件。
文档转化
读取出来的文档往往比较长,不适合直接写入向量库,需要先做切分。
但是这里有个非常关键的问题先要想清楚:为什么必须切?整篇直接灌进去不行吗?
我们结合刚才那份员工手册文档看,至少有三个原因:
- 一个向量只能”代表”一段语义,整篇文档算出的向量会”平均化”
- 大模型有上下文窗口上限
- 不切分还会带来噪声和成本问题
Spring AI 中负责文档转换的统一接口是 DocumentTransformer:
```
public interface DocumentTransformer extends Function<List<Document>, List<Document>> {
default List<Document> transform(List<Document> transform) {
return apply(transform);
}
}
```
Spring AI 官方 ETL Pipeline 文档里提供的切分文档的DocumentTransformer只有一个 TokenTextSplitter,它包含如下参数:
chunkSize(默认 800):每个 chunk 的目标 token 数。切分时以这个值作为参考长度。minChunkSizeChars(默认 350):chunk 的最小字符数。切分时只会在大于这个长度之后再去找标点切分点,避免切出过短的小碎片。minChunkLengthToEmbed(默认 5):chunk 被保留的最小长度。切完之后如果一段比这个值还短,就直接丢弃,不会进入最终结果。maxNumChunks(默认 10000):一份文档最多切成多少段。超过就停止切分,主要起兜底保护作用,防止异常文本切出几十万段。keepSeparator(默认 true):是否保留分隔符(主要是换行\n)。true时保留,文档段落格式更接近原文;false时去掉,文本更紧凑。
这里要注意:Spring AI 1.x 的 TokenTextSplitter 没有 withPunctuationMarks 配置方法,句子边界字符在源码中是固定的:英文句号 .、英文问号 ?、英文感叹号 ! 和换行 \n。也就是说,它不能直接通过 builder 配置中文句号 。、中文问号 ?、中文感叹号 !、中文分号 ;。
了解了参数,再大致说一下它的切分思路:先按 chunkSize 切出候选 chunk;只有当剩余 token 数超过 chunkSize 时,才会在 minChunkSizeChars 之后、最靠后的一个内置边界字符处截断,让 chunk 尽量落在句子边界上;最后过滤掉长度小于 minChunkLengthToEmbed 的小碎片。完整的规则细节可以参考附录的[TokenTextSplitter 完整切分规则](#tokentextsplitter-完整切分规则)。
针对我们的”员工手册”案例,可以这样来定义TokenTextSplitter:
```
import org.springframework.ai.document.Document;
import org.springframework.ai.transformer.splitter.TokenTextSplitter;
import org.springframework.stereotype.Service;
import java.util.List;
@Service
public class DocumentSplitService {
public List<Document> split(List<Document> documents) {
// 使用 builder 方式构造 TokenTextSplitter,可读性更好
TokenTextSplitter splitter = TokenTextSplitter.builder()
// 员工手册每个小节都不长,这里把目标长度调小,方便按小节附近切开
.withChunkSize(160)
// 小节段落通常在百字以上,80 个字符后再找边界,避免切得太碎
.withMinChunkSizeChars(80)
// 切完之后,长度小于 5 个字符的 chunk 会被直接丢弃(基本是噪声)
.withMinChunkLengthToEmbed(5)
// 示例文档较短,100 段已经足够作为保护上限
.withMaxNumChunks(100)
// 保留换行等分隔符,让段落格式更接近原文(方便后续阅读和调试)
.withKeepSeparator(true)
.build();
return splitter.apply(documents);
}
}
```
切分参数没有”万能值”,但对中文知识库类 RAG 项目,有几条建议:
- chunk 大小要跟文档粒度匹配
像员工手册这种每个小节本来就不长的文档,可以先从 100~300 token 开始;长篇文章或报告再考虑 300~800 token。太小语义不完整,太大则容易把多个制度主题混在一个向量里。
- chunk 之间保留 10%~20% 的重叠(chunk overlap)
切分时让相邻 chunk 共享一小段内容(比如 50~150 token),可以避免一句关键描述刚好被切到两段之间、导致两段都检索不命中。重叠比例过大则会浪费存储和 token 成本。
- 优先按”语义边界”切,而不是按字符硬切
能按标题切就按标题切,能按段落切就按段落切,再退而求其次按句子切。TokenTextSplitter 会优先在内置的英文标点和换行处做边界对齐。
- 不同类型文档应使用不同切分策略
- 结构化文档(Markdown、HTML):优先按章节、标题切。
- 长篇文章(PDF、Word):按段落 + 句子切,加 overlap。
- 问答、FAQ 类:以”一条 QA”为最小单元,不强切。
- 代码:按函数/类切,避免切断函数体。
- 针对中文,优先保留换行和段落结构
Spring AI 1.x 不能配置中文标点边界,因此中文文档不要只依赖 TokenTextSplitter 自动按中文句号切句。更稳妥的做法是:在原始文档中保留清晰的段落换行,或者在进入 TokenTextSplitter 之前先按标题、段落、FAQ 条目等业务语义边界做一层预切分。
- 切分不是一次性工作,要持续监控
线上跑一段时间后,看哪些用户问题召回不到、或者召回了无关内容,回过头调整 chunk 大小、overlap、预切分策略和段落结构,是一个迭代过程。
文档写入
文档写入的核心接口是 DocumentWriter:
```
public interface DocumentWriter extends Consumer<List<Document>> {
default void write(List<Document> documents) {
accept(documents);
}
}
```
在 RAG 中,我们通常写入的是向量数据库,所以更常用的是 VectorStore。
```
public interface VectorStore extends DocumentWriter, VectorStoreRetriever {
void add(List<Document> documents);
// 其他方法暂时先不关注
}
```
VectorStore 可以理解为 Spring AI 对向量数据库的统一抽象。底层可以是 Elasticsearch、Redis、Chroma、Milvus 等,不同的向量数据库有不同的VectorStore实现。
我们使用 Elasticsearch 存储向量,需要引入依赖:
```
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-starter-vector-store-elasticsearch</artifactId>
</dependency>
```
这里没有单独写 `,是因为前面已经通过 spring-ai-bom` 统一管理了 Spring AI 相关依赖的版本。
同时还需要配置 Embedding 模型。因为写入向量库时,文本必须先转成向量。
```
spring:
ai:
openai:
base-url: https://api.example.com
api-key: your-api-key
chat:
options:
model: your-chat-model
embedding:
options:
model: your-embedding-model
vectorstore:
elasticsearch:
index-name: spring-ai-docs
dimensions: 1536
elasticsearch:
uris: http://ip地址:9200
```
| 参数 | 含义 | |
|---|---|---|
index-name | Elasticsearch 索引名 | |
dimensions | 向量维度,必须和 Embedding 模型输出维度一致 | |
spring.elasticsearch.uris | Elasticsearch 服务地址 |
VectorStore 在保存文档时,并不是直接把文本原样写入向量库。它会先调用配置好的 Embedding 模型,把文档文本转成向量,然后再把向量和文档内容一起保存到向量数据库。
```
import org.springframework.ai.document.Document;
import org.springframework.ai.vectorstore.VectorStore;
import org.springframework.stereotype.Service;
import java.util.List;
@Service
public class DocumentWriteService {
private final VectorStore vectorStore;
public DocumentWriteService(VectorStore vectorStore) {
this.vectorStore = vectorStore;
}
public void write(List<Document> documents) {
// add方法会先完成Embedding向量化,再写入向量数据库
vectorStore.add(documents);
}
}
```
完整 ETL Controller
把读取、切分、写入串起来:
```
import org.springframework.ai.document.Document;
import org.springframework.ai.reader.TextReader;
import org.springframework.ai.transformer.splitter.TokenTextSplitter;
import org.springframework.ai.vectorstore.VectorStore;
import org.springframework.core.io.Resource;
import org.springframework.core.io.ResourceLoader;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;
@RestController
public class RagEtlController {
private final ResourceLoader resourceLoader;
private final VectorStore vectorStore;
public RagEtlController(ResourceLoader resourceLoader, VectorStore vectorStore) {
this.resourceLoader = resourceLoader;
this.vectorStore = vectorStore;
}
@PostMapping("/ai/rag/etl")
public String etl() {
Resource resource = resourceLoader.getResource("classpath:document/employee-handbook.txt");
TextReader textReader = new TextReader(resource);
textReader.getCustomMetadata().put("source", "employee-handbook.txt");
List<Document> documents = textReader.get();
TokenTextSplitter splitter = TokenTextSplitter.builder()
.withChunkSize(160)
.withMinChunkSizeChars(80)
.withMinChunkLengthToEmbed(5)
.withMaxNumChunks(100)
.withKeepSeparator(true)
.build();
List<Document> splitDocuments = splitter.apply(documents);
vectorStore.add(splitDocuments);
return "ETL完成,写入文档片段数量:" + splitDocuments.size();
}
}
```
知识的检索及融合
知识检索理论
文档中的知识在经过了ETL流程后,就以向量的形式存储到了向量数据库中。当用户提出问题后,如何从向量数据库中找到和用户问题相关的知识呢?其实,思路也很简单,就是将用户的问题也用同一个向量模型转化为向量,然后计算用户的问题向量与知识库中向量的相似度,相似度越高,则说明改向量所代表的知识与用户问题越相关。
如何判断向量之间的相关性呢?通过余弦相似度,其值基于向量点积和模长推导,对应的公式如下:
$$
\cos\theta = \frac{\vec{a} \cdot \vec{b}}{|\vec{a}| \cdot |\vec{b}|}
$$
其中$\vec{a}$和 $\vec{b}$是两个非零向量,$\theta$ 是它们的夹角,$\vec{a} \cdot \vec{b}$是点积,$|\vec{a}|$、$|\vec{b}|$是向量模长。
- 当 cos*θ*=1 时,*θ*=0∘,两向量方向完全相同(相似度最高)。
- 当 cos*θ*=0 时,*θ*=90∘,两向量垂直(无任何方向关联)。
- 当 cos*θ*=−1 时,*θ*=180∘,两向量方向完全相反(相似度最低)。
知识检索实现
首先需要引入如下依赖:
```
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-advisors-vector-store</artifactId>
</dependency>
```
知识入库后,用户提问时就要做两件事:
- 从向量数据库检索相关知识。
- 把检索到的知识放进 Prompt,让模型参考后回答。
Spring AI 提供了 QuestionAnswerAdvisor 帮我们完成这个流程。
简单理解:QuestionAnswerAdvisor 是一个专门用于 RAG 问答的 Advisor,它会在调用模型前,根据用户问题去 VectorStore 检索相关文档,并把文档内容加入 Prompt。
```
import org.springframework.ai.chat.client.ChatClient;
import org.springframework.ai.chat.client.advisor.vectorstore.QuestionAnswerAdvisor;
import org.springframework.ai.vectorstore.SearchRequest;
import org.springframework.ai.vectorstore.VectorStore;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
@RestController
public class RagChatController {
private final ChatClient chatClient;
private final VectorStore vectorStore;
public RagChatController(ChatClient chatClient, VectorStore vectorStore) {
this.chatClient = chatClient;
this.vectorStore = vectorStore;
}
@GetMapping("/ai/rag/chat")
public String ragChat(@RequestParam String question) {
QuestionAnswerAdvisor advisor = QuestionAnswerAdvisor.builder(vectorStore)
.searchRequest(SearchRequest.builder()
.topK(5)
.similarityThreshold(0.7)
.build())
.build();
return chatClient.prompt()
.advisors(advisor)
.user(question)
.call()
.content();
}
}
```
这里重点理解两个参数:
| 参数 | 作用 |
|---|---|
topK | 最多检索几个相关片段 |
similarityThreshold | 相似度低于阈值的片段不要 |
如果 topK 太小,可能漏掉关键信息;如果太大,Prompt 会变长,噪声也会增加。如果阈值太高,可能查不到内容;如果太低,可能把不相关内容塞给模型。
RAG 调优很多时候就是在调这些参数,以及调文档切分策略。
Tool Calling
前面 RAG 解决的是“让模型看我们的知识”。但是很多业务场景中,模型不只是要看资料,还要调用系统能力。
比如查询今天北京天气、查询某个订单状态、创建一条待办事项、调用挂号接口、根据患者 ID 查询病历摘要。
这就需要 Tool Calling,也叫工具调用。
Spring AI 官方提到,Tools 主要用于两类场景:
Information Retrieval:信息检索。让模型调用工具查询外部信息,例如天气、库存、订单、病历。Taking Actions:执行动作。让模型调用工具完成某个操作,例如创建工单、发送消息、提交表单。
代码层面什么是 Tool
在 Spring AI 中,代码层面可以把一个 Java 方法暴露成 Tool。
也就是说:一个带有工具描述的方法,就是模型可以选择调用的能力。
最常用的是声明式定义,也就是使用 @Tool 注解。
```
import org.springframework.ai.tool.annotation.Tool;
import org.springframework.ai.tool.annotation.ToolParam;
import org.springframework.stereotype.Component;
@Component
public class WeatherTools {
@Tool(description = "根据城市名称查询当前天气")
public String getWeather(
@ToolParam(description = "城市名称,例如北京、上海、广州") String city) {
if ("北京".equals(city)) {
return "北京今天晴,气温 18 到 28 摄氏度。";
}
return city + "今天多云,气温 20 到 26 摄氏度。";
}
}
```
| 注解 | 作用 |
|---|---|
@Tool | 标记当前方法可以作为工具暴露给模型,并描述工具用途 |
@ToolParam | 描述方法参数含义,帮助模型正确生成工具入参 |
这里要特别注意:模型不是直接读 Java 代码。模型主要根据工具名称、工具描述、参数描述来判断是否调用工具。所以描述一定要清楚。
工具调用过程
工具调用的整体流程如下:
```
sequenceDiagram
participant User as 👤 用户
participant AIApp as 🤖 AI应用
participant LLM as 💡 基础大模型
participant Tools as 🔧 工具集
User->>AIApp: 🟪 1. 提出问题(如:北京今天天气?)
AIApp->>LLM: 🟦 2.1 工具信息 + Prompt提示词
alt 根据大模型返回的响应判断,是否需要工具调用
LLM->>AIApp: 🟨 2.2 判断无需工具调用,返回最终响应
else 根据大模型返回的响应判断,需要工具调用
AIApp->>Tools: 🟦 2.3 调用指定工具(传入入参)
Tools->>AIApp: 🟩 2.4 返回工具执行结果(如天气数据/计算结果)
AIApp->>LLM: 🟦 2.5 整合Prompt(Prompt提示词+工具结果)
LLM->>AIApp: 🟦 2.6 返回结合工具执行结果的最终响应
end
```
用大白话讲:应用先把“有哪些工具可用”告诉模型;模型根据用户问题判断要不要调用工具;如果需要调用,模型返回工具名和参数;Spring AI 根据工具名和参数调用对应 Java 方法;工具返回结果;Spring AI 再把工具结果交给模型;模型基于工具结果组织最终回答。
所以真正执行工具的不是模型本身,而是我们的 Java 应用。
如果让我们定义的Tool生效呢?只需要把定义好的tool设置给ChatClient对象即可
```
import org.springframework.ai.chat.client.ChatClient;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
@RestController
public class ToolChatController {
private final ChatClient chatClient;
private final WeatherTools weatherTools;
public ToolChatController(ChatClient chatClient, WeatherTools weatherTools) {
this.chatClient = chatClient;
this.weatherTools = weatherTools;
}
@GetMapping("/ai/tool/weather")
public String weather(@RequestParam String question) {
return chatClient.prompt()
.user(question)
// tools 方法可以接收一个或多个工具对象
// Spring AI 会自动扫描其中标了 @Tool 的方法并注册给本次调用
.tools(weatherTools)
.call()
.content();
}
}
```
当用户问”北京今天天气怎么样”时,模型会发现普通语言模型本身不知道实时天气,于是选择调用 getWeather 工具。
附录
Prompt 优化技巧
下面补充一些常见 Prompt 技巧。大家先不用死记,先理解它们分别解决什么问题。
给 AI 提供清晰的任务描述和角色定位,帮助模型理解背景和期望。
``` 系统:你是一位经验丰富的 Python 教师,擅长向初学者解释编程概念。 用户:请解释 Python 中的列表推导式,包括基本语法和 2-3 个实用示例。 ```
提供足够的上下文信息和期望输出格式,减少模型不确定性。
``` 请提供一个社交媒体营销计划,针对一款新上市的智能手表。计划应包含: 1. 目标受众描述 2. 三个内容主题 3. 每个平台的内容类型建议 4. 发布频率建议 示例格式: 目标受众:[描述] 内容主题:[主题1],[主题2],[主题3] 平台策略:[平台] - [内容类型] - [频率] ```
通过列表、表格等结构化格式,让模型输出更有条理。
``` 分析以下公司的优势和劣势: 公司:Tesla 请使用表格格式回答,包含以下列: - 优势(最少3项) - 每项优势的简要分析 - 劣势(最少3项) - 每项劣势的简要分析 - 应对建议 ```
引导模型分步骤思考,提高复杂问题的准确性。
``` 问题:一个商店售卖 T 恤,每件 15 元。如果购买 5 件以上可以享受 8 折优惠。小明买了 7 件 T 恤,他需要支付多少钱? 请一步步思考解决这个问题: 1. 首先计算 7 件 T 恤的原价 2. 确定是否符合折扣条件 3. 如果符合,计算折扣后的价格 4. 得出最终支付金额 ```
通过提供几个输入输出示例,帮助模型理解任务模式。
``` 我将给你一些情感分析的例子,然后请你按照同样的方式分析新句子的情感倾向。 输入:"这家餐厅的服务太差了,等了一个小时才上菜" 输出:负面,因为描述了长时间等待和差评服务 输入:"新买的手机屏幕清晰,电池也很耐用" 输出:正面,因为赞扬了产品的多个方面 现在分析这个句子: "这本书内容还行,但是价格有点贵" ```
将复杂任务拆成步骤,降低模型遗漏关键环节的概率。
``` 请帮我创建一个简单的网站落地页设计方案,按照以下步骤: 步骤1:分析目标受众(考虑年龄、职业、需求等因素) 步骤2:确定页面核心信息(主标题、副标题、价值主张) 步骤3:设计页面结构(至少包含哪些区块) 步骤4:制定视觉引导策略(颜色、图像建议) 步骤5:设计行动召唤按钮和文案 ```
让模型检查自己的输出,提高复杂任务的可靠性。
``` 解决以下概率问题: 从一副标准扑克牌中随机抽取两张牌,求抽到至少一张红桃的概率。 首先给出你的解答,然后: 1. 检查你的推理过程是否存在逻辑错误 2. 验证你使用的概率公式是否正确 3. 检查计算步骤是否有误 4. 如果发现任何问题,提供修正后的解答 ```
当答案需要依据材料时,要求模型说明依据,降低胡编乱造概率。
``` 请解释光合作用的过程及其在植物生长中的作用。在回答中: 1. 提供光合作用的科学定义 2. 解释主要的化学反应 3. 描述影响光合作用效率的关键因素 4. 说明其对生态系统的重要性 对于任何可能需要具体数据或研究支持的陈述,请明确指出这些信息的来源,并说明这些信息的可靠性。 ```
要求模型从不同角色或专业视角分析问题,让答案更全面。
``` 分析“城市应该禁止私家车进入市中心”这一提议: 请从以下 4 个不同角度分析: 1. 环保专家视角 2. 经济学家视角 3. 市中心商户视角 4. 通勤居民视角 对每个视角: - 提供支持该提议的 2 个论点 - 提供反对该提议的 2 个论点 - 分析可能的折中方案 ```
TokenTextSplitter 完整切分规则
- 这里把 Spring AI 1.x 中
TokenTextSplitter的完整切分过程列出来,方便课后回顾。日常使用只需要知道”按 token 切、尽量在内置英文标点或换行处切开”即可,但当你遇到切分结果不符合预期时,这份规则可以帮你定位原因。
完整流程如下:
- 先把整段文本用
CL100K_BASE编码转成 token 序列。 - 每轮先判断当前还剩多少 token:
- 如果剩余 token 数不超过
chunkSize,直接把剩余内容作为最后一个候选 chunk。 - 如果剩余 token 数超过
chunkSize,取前chunkSize个 token 作为候选窗口,并把它解码回文本。 - 对候选窗口尝试做边界对齐:
- 切分点只会从源码内置的 4 个边界字符中找:英文句号
.、英文问号?、英文感叹号!、换行\n。 - 如果找到了边界字符,并且它的位置大于
minChunkSizeChars,就在这个边界字符后提前截断 chunk。 - 如果没有找到边界字符,或者边界字符位置不大于
minChunkSizeChars,就不按标点截断,直接使用当前候选窗口。 - 提前截断时,候选窗口里没进入当前 chunk 的 token 不会丢,会留到下一轮继续处理。
- 去除首尾空白;根据
keepSeparator决定是否把系统换行符替换为空格。 - 如果最终长度大于
minChunkLengthToEmbed,就把它加入结果,否则丢弃。 - 从剩余 token 中移除已经加入当前 chunk 的 token,继续处理后面的文本。
- 直到所有 token 处理完或达到
maxNumChunks上限。 - 如果因为达到
maxNumChunks上限而仍有剩余文本,剩余文本长度满足minChunkLengthToEmbed时,会作为最后一个 chunk 加入结果。
这里要特别注意两点:
- 只有当前剩余 token 数超过
chunkSize时,才会尝试按边界字符截断。如果一份文档本来就比chunkSize小,它会被作为单独一个 chunk 返回,不会被多余地切碎,避免对小文档造成不必要的碎片化。 - Spring AI 1.x 的
TokenTextSplitter里中文标点。 ? ! ;不在内置边界字符里。如果中文文档没有换行,它通常会更接近按 token 长度硬切,而不是严格按中文句子切。
工业级 RAG 常见增强方案
课堂中我们只展示了最基础的 RAG:用户问题进入系统后,直接到向量数据库检索相关片段,然后把片段塞进 Prompt 交给大模型回答。但在真实工业级场景中,RAG 往往会做很多增强,常见方案如下:
- 查询改写(Query Rewrite):先让模型把用户口语化、模糊、不完整的问题改写成更适合检索的标准问题,提高召回质量。
- 多查询扩展(Multi-Query):针对同一个用户问题生成多个不同表达方式的问题,分别检索后合并结果,减少“问法不同导致搜不到”的问题。
- HyDE(Hypothetical Document Embeddings):先让模型根据问题生成一段“假想答案/假想文档”,再用这段文本去做向量检索,适合用户问题很短但语义不充分的场景。
- 混合检索(Hybrid Search):同时使用向量检索和关键词检索,向量检索擅长语义相似,关键词检索擅长精确匹配专有名词、编号、术语,两者互补。
- 元数据过滤(Metadata Filter):检索时根据文档来源、时间、部门、权限、业务类型等元数据先过滤范围,避免把不该看的或不相关领域的知识召回。
- 上下文窗口扩展(Context Window Expansion):检索命中某个片段后,把它前后的相邻片段也一起带上,解决单个片段上下文不完整的问题。
- ReRank(重排序/精排):向量库相似度通常是“粗排”,它只比较“问题向量”和“文档向量”在向量空间里的距离,速度快但判断较粗;ReRank 会把用户问题和候选文档成对输入给专门的重排模型,让模型更细致地判断“这段文档是否真的能回答这个问题”,再重新排序,所以叫“精排”。
- 查询路由(Query Routing):先判断用户问题属于哪个知识库、哪个业务域或是否需要查数据库,再路由到对应的检索器或工具。
- 多路召回(Multi-Retriever):同时从向量库、数据库、搜索引擎、知识图谱等多个来源召回信息,再统一融合。
MCP
我们先回到第一个问题。在前面学习 Tool Calling 的时候,我们写过类似这样的工具:
```
public class WeatherTools {
@Tool(description = "根据城市名查询天气")
public String getWeather(String city, String date) {
// 调天气接口
return "晴,25度";
}
}
```
这个工具写得很好,但它有几个明显的局限:
- 语言限制:这个工具只能在 Java 项目里被大模型调用。如果另一个团队用 Python 写 AI 应用,他们没法直接用我们写好的这个工具。
- 难以复用:换一个新的 SpringBoot 工程,这个工具的代码就要拷贝一份过去;如果工具内部依赖了一堆 jar 包,整个项目结构也要跟着搬。
- 重复造轮子:今天我要查天气,自己实现一遍;明天另一个团队也要查天气,他们也要从头实现一遍。全世界都在重复写同样的工具代码,浪费时间。
那有没有一种方式,让工具的开发和使用变成”标准化”的?任何人开发的工具,任何人都可以拿来直接给自己的 AI 用 —— 就像 USB 接口一样,谁的设备都能插上去 , MCP协议就是为此而诞生的。
MCP 的全称是 Model Context Protocol(模型上下文协议),是由 Anthropic 公司(也就是 Claude 的开发公司)在 2024 年提出的一套开放协议。
它的核心作用是:统一 AI 应用和外部工具/数据源的交互方式。
打个比方:在 MCP 出现之前,每家的 AI 工具就像每家的电源插头都不一样 —— 有圆头的、有扁头的、有三孔的,互相之间没法通用。MCP 就相当于规定了”全世界统一用一种插头”,从此以后所有的 AI 工具都按这个标准来开发,所有的 AI 应用也都能直接用这些工具。
它的存在意义在于:
- 跨语言:MCP 是一个协议(类似 HTTP),任何语言都可以实现,不再受 Java、Python 等语言的限制
- 能复用:一个 MCP 工具开发好之后,任何 AI 应用都可以接入使用
- 生态丰富:因为是开放协议,全世界的开发者都在贡献 MCP 工具,我们用别人写好的就行
MCP 的 C/S 架构
MCP 协议采用的是 C/S 架构(客户端/服务端架构)。这里有两个核心角色:
- MCP Server(服务端):负责把工具”对外暴露”出来。具体来说,谁实现了一个工具(比如查天气),就把这个工具包在一个 MCP Server 里面,然后启动这个 Server。
- MCP Client(客户端):负责”调用”工具。我们的 AI 应用(比如基于 Spring AI 写的应用)就充当 MCP Client,去连接 MCP Server,发现 Server 上有哪些工具可用,然后让大模型按需调用。
```
graph LR
A[AI应用<br/>MCP Client] -->|调用工具| B[MCP Server<br/>暴露工具]
B -->|返回结果| A
```
那 MCP Client 和 MCP Server 之间到底通过什么方式通信呢?这就涉及到 MCP 的两种通信方式。

通信方式一:Stdio(标准输入输出)
Stdio 通信方式是基于操作系统的标准输入和标准输出来传递消息的。它的工作过程是:
- MCP Client 在本地把 MCP Server 作为子进程启动起来
- MCP Client 通过往这个子进程的标准输入(stdin)写数据来发送请求
- MCP Server 通过往自己的标准输出(stdout)写数据来返回响应
简单说,就是 Client 和 Server 之间通过”管道”通信,两者必须运行在同一台机器上。
易混淆点:标准输入不就是键盘、标准输出不就是控制台吗?
键盘和屏幕只是 stdin/stdout 的默认连接对象,不是它们的本质。stdin 本质上是”进程读数据的入口”,stdout 是”进程写数据的出口”,至于这两头接的是什么,是可以被替换的。比如命令行里的管道
cat data.txt | sort,前一个程序的标准输出直接成了后一个程序的标准输入,全程没有键盘和屏幕参与。
MCP 的 Stdio 模式正是利用了这一点:Client 把 Server 作为子进程启动时,操作系统把 Server 的 stdin/stdout 接到了与 Client 之间的管道上,而不是接到键盘和屏幕。于是 Client 把请求写进 Server 的 stdin,Server 把响应写到自己的 stdout 流回 Client —— 这是一条进程间通信的专用数据通道,并非给人看的控制台。也正因为 stdout 被协议数据占用了,这类 Server 的日志通常改打到 stderr,以免污染数据。
还有一个小问题:为什么是 Client 启动 Server,而不是 Server 自己独立运行?
原因是:Stdio 本质上是基于操作系统的管道通信,而管道通信有一个硬性约束 —— 必须是父子进程关系才能建立。所以在 Stdio 模式下,必然是其中一方把另一方拉起来。MCP 协议规定:由 Client 把 Server 作为子进程启动。
适用场景:本地工具的调用。比如读写本地文件、操作本地数据库、调用本地 CLI 工具等。
通信方式二:SSE / Streamable HTTP(基于 HTTP 的网络通信)
这种方式是基于 HTTP 协议进行通信的,Client 和 Server 可以分别部署在不同的机器上,通过网络访问。
这种模式下,MCP Server 是独立部署、独立启动的一个 HTTP 服务(比如就是一个普通的 Spring Boot 应用),它一直在某个端口上监听请求。MCP Client 不需要去启动 Server,只需要知道 Server 的 URL 地址,通过 HTTP 协议发送请求即可。这就和我们平时调用一个普通的 Web 服务非常类似。
适用场景:远程工具的调用。比如远端服务器上部署了一个 MCP Server,全世界的 AI 应用都可以连过来用它。
两种通信方式的对比:
| 通信方式 | 部署形式 | Server 由谁启动 | Server 生命周期 | 适用场景 |
|---|---|---|---|---|
| Stdio | Client 和 Server 在同一台机器 | Client 把 Server 作为子进程拉起 | Client 退出,Server 也退出 | 本地工具:读写文件、操作本地数据库等 |
| SSE / HTTP | Client 和 Server 通过网络通信 | Server 独立部署、独立启动 | Server 长期运行,与 Client 无关 | 远程工具:远端服务、跨网络的工具调用 |
发现MCP 服务
我们刚才说了,MCP 一大好处就是”别人写的工具我们可以直接用”。那这些”别人写好的 MCP Server”去哪里找呢?
主要有两个渠道:
渠道一:MCP 服务市场
服务市场就像一个”应用商店”,里面收录了大量开发者贡献的 MCP Server。比较有名的有:
- MCP.so:社区维护的 MCP 服务市场,收录了大量 MCP Server
- Smithery:另一个非常流行的 MCP 服务市场
- modelscope MCP 广场:阿里魔搭社区维护的 MCP 服务市场,国内访问较快
需要特别说明的是:绝大多数 MCP 服务市场,提供的都是”本地下载 MCP 服务端代码并运行”的方式。
也就是说,你在市场上找到一个心仪的 MCP Server,市场上给你的不是一个现成的远程地址,而是一行安装命令(比如 npm install),你需要把这个 Server 下载到自己的本地机器上。然后你的 AI 应用就用 Stdio 方式连接它。
渠道二:云服务平台提供的云端 MCP
MCP 的实现
Spring AI 本身就支持MCP协议,接下来我们基于Spring AI用两个案例分别演示:
- 案例一:作为 MCP Client,调用别人实现的 MCP Server —— 高德地图 MCP Server
- 案例二:作为 MCP Server,自己实现一个查天气的 MCP Server,并让另一个项目去调用它
案例一:调用别人的 MCP Server
我们这个案例要实现的功能是:根据一个公网 IP,让大模型告诉我们这个 IP 对应哪个城市。
大模型自己肯定不知道某个 IP 对应哪儿,但是高德地图官方提供了一个 MCP Server,里面就有”IP 定位”的工具,我们直接用它就行。
第一步:本地安装并运行高德地图 MCP Server
因为我们前面说过,绝大多数 MCP Server 都是本地运行的,受限于上课环境(没有可用的远端 MCP),我们也使用本地运行的方式。
打开命令行,执行下面这条命令安装高德地图官方 MCP Server:
``` npm install -g @amap/amap-maps-mcp-server --registry=https://registry.npmmirror.com ```
这条命令的含义是:通过 npm 全局安装高德地图官方提供的 MCP Server 包。--registry 参数指定从国内的 npm 镜像源下载,速度更快。
安装成功之后,这个 MCP Server 就在我们本机上准备好了(注意此时还没启动,等下我们的 AI 应用启动时会通过 Stdio 方式自动把它拉起来)。
第二步:在 Spring Boot 项目中引入 MCP Client 依赖
```
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-starter-mcp-client</artifactId>
</dependency>
```
这个依赖就是 Spring AI 提供的 MCP Client Starter,引入后我们就具备了连接 MCP Server 的能力。
第三步:添加 MCP Client 的 yml 配置
```
spring:
ai:
mcp:
client:
stdio:
servers-configuration: classpath:mcp-servers/servers.json # MCP Server 配置文件路径
```
每个配置项的含义都已经在注释中说明。这里最关键的是最后一个 servers-configuration,它指向一个文件 servers.json,这个文件里到底是什么呢?
第四步:编写 servers.json 配置文件
servers.json 这个文件的作用,就是告诉 MCP Client:我要连接哪些 MCP Server,每个 Server 怎么启动。
这个文件可以从两个地方获取:
- 从 MCP 服务市场:在 MCP.so、modelscope 等市场上,每个 MCP Server 的详情页会直接给出可以拷贝的 servers.json 配置片段
- 从 MCP Server 官方文档:比如高德地图的 MCP Server 文档里也提供了对应的配置
对于我们案例里的高德地图 MCP Server,在 src/main/resources/mcp-servers/ 目录下创建 servers.json,内容如下:
```
{
"mcpServers": {
"amap-maps": {
"command": "npx.cmd",
"args": [
"-y",
"@amap/amap-maps-mcp-server"
],
"env": {
"AMAP_MAPS_API_KEY": "你自己的高德地图API Key"
}
}
}
}
```
简单解释一下:
mcpServers下可以配置多个 MCP Server,每个 Server 有一个名字(比如这里的amap-maps)command+args:通过什么命令启动这个 MCP Server。这里就是用npx跑刚才我们 npm 安装好的包env:环境变量。高德地图的 MCP Server 需要一个 API Key 才能调用高德的接口,我们要把自己的 Key 填在这里(去高德开放平台 https://lbs.amap.com 免费申请即可)
到这一步,配置就完成了。Spring AI 启动时会自动读 yml → 读 servers.json → 通过 Stdio 方式启动高德 MCP Server → 建立连接 → 获取 Server 上的所有工具。
第五步:编写配置类,把 MCP 的工具注册到 ChatClient
```
@Configuration
public class McpClientConfig {
/**
* 构建 ChatClient,并注入 MCP Client 提供的所有工具
*
* @param builder ChatClient 的构建器
* @param toolCallbackProvider Spring AI 自动注入的工具回调提供者,
* 它会把 MCP Server 上的所有工具自动转换为 ToolCallback
*/
@Bean
public ChatClient chatClient(ChatClient.Builder builder,
ToolCallbackProvider toolCallbackProvider) {
return builder
// 把 MCP Server 上的工具,全部注入给 ChatClient
.defaultToolCallbacks(toolCallbackProvider)
.build();
}
}
```
这里有一个关键点:ToolCallbackProvider 是 Spring AI 自动帮我们注入的 Bean,它内部包含了从所有 MCP Server 上发现到的工具列表。我们把它注入到 ChatClient 之后,大模型就可以自动选择并调用这些 MCP 工具了,效果就跟我们前面学过的 @Tool 注解定义的本地工具一模一样。
第六步:编写 Controller 测试
```
@RestController
public class McpController {
private final ChatClient chatClient;
public McpDemoController(ChatClient chatClient) {
this.chatClient = chatClient;
}
@GetMapping("/ip2city")
public String ip2city(String ip) {
// 用户消息,让大模型查这个 IP 所在的城市
// 大模型会自动判断需要调用高德 MCP 的 IP 定位工具
return chatClient.prompt()
.user("请帮我查一下 IP 地址 " + ip + " 对应的城市是哪个")
.call()
.content();
}
}
```
启动项目,访问 http://localhost:8080/ip2city?ip=114.247.50.2,就能看到大模型返回的城市信息。
到这里你可能会有疑问:我就写了几行配置和代码,为什么大模型就能调用高德的工具了?整个过程中数据是怎么流转的?
我们用一个时序图来梳理:
```
sequenceDiagram
participant User as 👤 用户
participant App as 🤖 AI应用<br/>(MCP Client)
participant LLM as 💡 大模型
participant Server as 🔧 高德MCP Server
Note over App,Server: 启动阶段
App->>Server: 1. 通过Stdio启动MCP Server进程
App->>Server: 2. 询问可用的工具列表
Server->>App: 3. 返回工具列表(含IP定位工具等)
Note over User,Server: 调用阶段
User->>App: 4. 提问:IP 114.247.50.2 是哪里?
App->>LLM: 5. 发送用户消息 + 工具列表(MCP 提供)
LLM->>App: 6. 返回:需要调用 IP 定位工具
App->>Server: 7. 通过 Stdio 调用 IP 定位工具
Server->>App: 8. 返回 IP 对应的城市信息
App->>LLM: 9. 把工具结果再次发给大模型
LLM->>App: 10. 整合后返回最终答复
App->>User: 11. 返回答复给用户
```
整个流程的关键点:
- 启动阶段:MCP Client 通过 Stdio 把 MCP Server 拉起来,并查询有哪些工具
- 调用阶段:MCP 工具和我们前面学的
@Tool定义的工具用法完全一样 —— 大模型决定要不要用,决定了就让 MCP Client 去调
案例二:实现自己的 MCP Server
学完了”调用别人的 MCP Server”,我们再来学”自己实现一个 MCP Server”。
我们要实现的 MCP Server 功能很简单:根据时区名查询当前时间。
⚠️ 注意:MCP Server 要放在一个新的 SpringBoot 工程里(和案例一的 MCP Client 工程区分开),因为它们是两个独立的进程。
这里我们要让自己实现的 MCP Server 以 SSE/HTTP 方式 对外提供服务,这样它就是一个独立运行的 Web 服务,任何机器上的 MCP Client 都能通过 HTTP 连过来调用它,更接近企业级真实场景。
第一步:在新工程中引入 MCP Server 依赖
```
<!-- WebMVC 版本的 MCP Server Starter,对外提供基于 HTTP 的 MCP 服务 -->
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-starter-mcp-server-webmvc</artifactId>
</dependency>
```
注意:这里用的是 spring-ai-starter-mcp-server-webmvc,不是 Stdio 版本的 spring-ai-starter-mcp-server。webmvc 版本会自动帮我们暴露 SSE / Streamable HTTP 的端点。
第二步:添加 MCP Server 的 yml 配置
```
spring:
ai:
mcp:
server:
enabled: true # 启用 MCP Server
name: time-mcp-server # MCP Server 的名称
server:
port: 8081 # MCP Server 监听的端口,Client 通过这个端口连过来
```
```
public class TimeService {
/**
* 返回当前时间
*/
@Tool(description = "获取当前指定时区的时间,返回精确都按秒的日期字符串")
public String getTime(@ToolParam(description = "指定的时区名称, 比如Asia/Shanghai") String zone) {
System.out.println("方法被调用了!");
DateTimeFormatter dateTimeFormatter = null;
ZoneId zoneId = null;
try {
dateTimeFormatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
zoneId = ZoneId.of(zone);
return ZonedDateTime.now(zoneId).format(dateTimeFormatter);
} catch (Exception e) {
e.printStackTrace();
return "获取" + zone + "时区的时间失败";
}
}
}
```
这里用到的 @Tool 和 @ToolParam 注解,跟我们之前学 Tool Calling 时一模一样,完全没有任何特殊的写法。这一点正是 MCP 的优点:写工具的方式完全没变,只是部署形态变了而已。
第四步:编写配置类,把工具注册为 MCP Server 提供的工具
```
@Configuration
public class WeatherMcpServerConfig {
/**
* 把 WeatherService 中的工具方法,注册为 MCP Server 对外暴露的工具
*
* 注意:这里必须返回 List<ToolCallback> 或者 ToolCallbackProvider,
* 才能被 MCP Server Starter 自动识别并发布。
*/
@Bean
public List<ToolCallback> weatherToolCallbacks() {
// 通过 ToolCallbacks.from() 工具方法,把 WeatherService 中所有的 @Tool 方法
// 自动转换为 ToolCallback 对象,并放进 List
return List.of(ToolCallbacks.from(new WeatherService()));
}
}
```
第五步:启动 MCP Server
第六步:在 MCP Client 工程中追加这个新的 Server 配置
现在回到案例一的 MCP Client 工程,我们要让它同时连接两个 MCP Server:
- 高德地图 MCP Server(Stdio 方式)
- 我们自己的时间 MCP Server(SSE/HTTP 方式)
由于这两个 Server 的通信方式不同,配置位置也不同:
- Stdio 方式的高德 MCP Server 继续放在
servers.json - SSE 方式的时间 MCP Server 直接在 yml 中配置
修改 MCP Client 工程的 yml:
```
spring:
ai:
mcp:
client:
stdio:
servers-configuration: classpath:mcp-servers/servers.json
# SSE/HTTP 方式的 MCP Server,直接在 yml 中配置
sse:
connections:
time-mcp-server: # 给这个 Server 起个名字
url: http://localhost:8081 # Server 的访问地址(如果 Server 在远程机器上就改成对应 IP)
```
简单解释 SSE 配置:
connections下可以配置多个 SSE 方式的 MCP Serverurl:MCP Server 的访问基地址,由于这里我们 Server 也部署在本机,所以是localhost:8081。生产环境改成实际的远程地址即可
至此,MCP Client 启动之后,通过 Stdio 启动并连接高德 MCP Server,通过 HTTP 连接我们自己的天气 MCP Server。所有 Server 上的工具都会被汇总到 ToolCallbackProvider 中,对大模型而言看到的就是”一堆可用的工具”,不区分来源。
第七步:测试调用
写一个新的接口测试:
```
@GetMapping("/time")
public String weather(@RequestParam String zone) {
return chatClient.prompt()
.user("请帮我查询 " + zone + "当前的时间")
.call()
.content();
}
```
访问 http://localhost:8080/time?zone=GMT+8,可以看到大模型自动调用了我们自己实现的工具。
⚠️ 测试前需要保证 MCP Server 工程已经启动(监听 8081 端口),否则 Client 启动时会连接失败。
—
多模态
我们之前学的所有内容,无论是 Prompt、还是 RAG、还是 Tool Calling,输入和输出的核心都是”文字”。
但是在现实世界里,信息的载体远远不止文字。你随便打开微信聊天界面看看,里面有:文字、图片、语音、视频、文件 …… 每一种都是不同形式的信息。我们把这些不同形式的信息,称为不同的“模态”(Modality)。
那什么是多模态呢?多模态就是指模型能够同时处理多种形式的信息。
我们之前学的大模型(比如通义千问标准版、GPT-4 文本版)严格来说叫“大语言模型”(LLM),它们只能处理文字输入、产生文字输出。
而“多模态大模型” 在大语言模型的基础上,增加了对图片、音频、视频等其他模态的理解能力。比较有名的多模态模型有:
- 阿里 qwen-vl 系列、qwen-omni 系列:图像、视频、语音理解
- OpenAI GPT-4o:图、文、音的全模态模型
- Google Gemini:同样支持多模态输入
打个比方:
- 大语言模型 就像一个只会读字、只会写字的学者,你给他文字他懂、他能写文字给你
- 多模态模型 就像一个眼耳口手脑全能的人,你给他图他能看、给他音频他能听、给他文字他能读,输出形式也可以更丰富
下面我们用阿里的 qwen-omni-turbo 模型,演示一个”上传图片让模型描述图片内容”的案例。
qwen-omni-turbo 是通义千问的一个多模态模型,支持图像、音频、视频的输入,输出文本。
第一步:在 yml 中配置多模态模型
```
spring:
ai:
openai:
base-url: https://dashscope.aliyuncs.com/compatible-mode
api-key: ${DASHSCOPE_API_KEY}
chat:
options:
model: qwen-omni-turbo # 指定多模态模型
stream-usage: false # qwen-omni-turbo 要求 stream 调用时 stream-usage 必须为 false
```
第二步:编写 Controller
```
@RestController
public class MultiModalController {
private final ChatClient chatClient;
public MultiModalController(ChatClient.Builder builder) {
this.chatClient = builder.build();
}
/**
* 多模态接口:接收用户消息 + 上传的图片,返回模型对图片的理解(流式)
*
* @param message 用户文字消息,比如"这张图片里有什么?"
* @param image 用户上传的图片文件
* @return Flux<String> 流式响应,逐字返回模型的回答
*/
@PostMapping(value = "/image-understand",
produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> understandImage(@RequestParam String message,
MultipartFile image) throws IOException {
// 1. 把上传的图片包装为 Spring AI 中的 Media 对象
// 第一个参数是图片的 MIME 类型(image/jpeg、image/png 等)
// 第二个参数是图片的内容资源
Media imageMedia = Media.builder()
.mimeType(MimeType.valueOf(image.getContentType()))
.data(image.getResource())
.build();
// 2. 调用 ChatClient,注意 user() 方法的参数:
// 可以同时传入文字内容和 Media 对象(也就是图片)
// 这就是"多模态输入"的体现 —— 输入既有文字又有图片
return chatClient.prompt()
.user(promptUserSpec -> promptUserSpec
.text(message) // 文字部分
.media(imageMedia)) // 图片部分
.stream() // 使用流式响应
.content();
}
}
```
代码的核心在于 user() 这一段:通过 media() 方法把图片附加到用户消息里。大模型在拿到这个请求后,会同时理解用户的文字描述和图片内容,给出综合的回答。
测试时,可以用 ApiFox 或者 Postman 发起 multipart/form-data 请求:
message:填一段文字,比如”请描述这张图片”image:上传一张本地图片文件
返回值是流式响应,前端会逐字看到模型对图片的理解。
—
模型私有化部署
到目前为止,我们用的所有模型都是”在线 API”形式:写好 API Key,调用阿里云、OpenAI 等厂商的远程接口。
但是在实际项目中,这种”在线 API”在很多场景下是不能用的,因为:
1. 数据安全和隐私要求高的场景
像金融、政府、军工这些行业,业务数据通常都是高度敏感的。如果调用外部 API,意味着每次请求都要把用户数据通过公网发给第三方厂商,这是绝大多数企业的合规体系所不允许的。
2. 网络隔离的场景
有些企业的内网是和外网完全物理隔离的(俗称”内外网物理隔离”),根本没法访问 OpenAI、阿里云的接口。这种情况下,云端 API 是天然用不了的,必须用本地部署的模型。
3. 成本控制的场景
云端 API 通常按 token 计费,一旦应用规模上去(比如每天百万级调用),成本会非常高。而私有化部署的模式是”一次性买/租 GPU,长期使用,调用次数不限”,对于调用量极大的场景,长期下来成本反而更低。
4. 定制和稳定性要求高的场景
云端 API 偶尔会有限流、不稳定的情况,且模型版本可能会被厂商更新,导致同样的 prompt 突然出现行为变化。私有化部署的模型,版本完全由自己控制,不会出现一觉醒来模型行为变了的情况。
5. 离线场景
比如车载、工业设备、边缘计算场景,设备本身可能没有稳定的网络连接,必须把模型部署在设备上才能用。
以上场景就适合把模型部署到自己的服务器上,也就是模型的私有化部署。
3.2 Ollama 介绍
知道了为什么要私有化部署,下一个问题就是:怎么部署?
如果让我们自己从零搭建 GPU 服务、加载模型、写推理引擎,难度很大。好在已经有现成的工具帮我们做这件事,其中最流行的一个就是 —— Ollama。
Ollama 是一个开源的本地大模型运行框架,它把”下载模型 + 加载模型 + 提供 API”这一整套流程封装到了一个命令行工具里。你只需要执行一条命令,就能在自己的电脑上跑起来一个大模型。
Ollama 支持很多主流的开源模型,比如:Llama 系列(Meta 开源),Qwen 系列(阿里通义千问开源版),DeepSeek 系列,Mistral 系列
Gemma(Google 开源)….
官网地址:https://ollama.com,在官网首页,你可以看到所有支持的模型的完整列表,每个模型都有不同的规模(比如 0.5b、7b、72b,代表参数量),数字越小越省资源、越快,但能力也越弱。
通过 Ollama 部署运行本地模型
下面以 Linux 系统 为例(生产环境一般都是 Linux 服务器),讲解部署的完整过程。
第一步:下载 Ollama
打开 Ollama 官网的下载页 https://ollama.com/download,注意选择正确的平台:
- macOS:下载 .zip 包
- Windows:下载 .exe 安装包
- Linux:下载 tgz 压缩包(我们这里用 Linux)
在我们当前的课程中,我们使用 Linux 版本。
第二步:安装 Ollama
1)解压
``` tar -xzf ollama-linux-amd64.tgz -C /usr/local/ ```
将下载好的 tgz 包解压到 /usr/local/ 目录,解压后 ollama 可执行文件就放在 /usr/local/bin/ 下了。
2)验证安装
``` ollama -v ```
如果输出版本号,说明安装成功。
⚠️ 特别注意:如果输入 ollama -v 提示”command not found”,是因为 /usr/local/bin/ 没有加入 PATH 环境变量。这种情况下需要配置 PATH 环境变量:
``` export PATH=$PATH:/usr/local/bin ```
把这一行加到 ~/.bashrc 文件末尾,然后 source ~/.bashrc 让它生效。
3)配置 Ollama 相关的环境变量
除了 PATH 之外,我们还要配置三个专门给 Ollama 用的环境变量:
``` export OLLAMA_HOST=0.0.0.0:11434 export OLLAMA_MODELS=$HOME/ollama ```
分别解释下这三个环境变量的含义:
- OLLAMA_MODELS:指定 Ollama 把下载下来的模型文件存放在哪里。模型文件通常都很大(小则几 GB,大则几十 GB),默认会放在用户主目录下,磁盘可能不够。这里我们指定到
$HOME/ollama,可以根据自己服务器实际的磁盘情况自由调整。 - OLLAMA_HOST:指定 Ollama 服务监听的地址和端口。
0.0.0.0表示监听所有网卡(即允许其他机器来访问,不只是本机),11434是 Ollama 的默认端口。如果只填127.0.0.1:11434,那就只允许本机访问。
配置完后,启动 Ollama 服务:
``` ollama serve ```
到这一步,Ollama 本身就部署完成了,但是还没有任何模型可以用。启动成功后可以通过浏览器访问http://ip:11434/api/tags

第三步:下载并运行模型
我们这里下载一个轻量级模型 qwen2.5:0.5b 作为演示(只有约 400MB,跑起来不需要 GPU)。
打开 Ollama 官网,搜索 qwen2.5,进入模型详情页,选择 0.5b 规格,可以看到对应的命令:
``` ollama run qwen2.5:0.5b ```
执行这条命令,会发生两件事:
- 第一次执行:Ollama 检测到本地没有这个模型,会先去远程仓库下载(耗时取决于网速和模型大小),下载完后自动运行

- 后续执行:模型已经在本地了,直接运行,不会再下载

运行成功后,命令行进入一个交互式对话界面,你可以直接和模型聊天。按 Ctrl+D 退出对话,但模型服务仍在 11434 端口上运行,等着别人调用。
Spring AI 与 Ollama 整合
模型在本地跑起来了,下一步就是用 Spring AI 去调用它。
第一步:引入 Maven 依赖
```
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-starter-model-ollama</artifactId>
</dependency>
```
这个 starter 专门用来对接 Ollama,里面封装了和 Ollama 接口通信的所有细节。
第二步:添加 yml 配置
```
spring:
ai:
ollama:
base-url: http://192.168.1.100:11434 # Ollama 服务地址,改为你部署 Ollama 的机器 IP
chat:
options:
model: qwen2.5:0.5b # 指定使用的模型,必须是已经在 Ollama 中下载过的
```
注意:
base-url就是 Ollama 服务的访问地址。如果 Ollama 部署在本机就写http://localhost:11434,部署在其他机器就写远程 IP。model必须是你前面用ollama run下载过的模型,否则会报”模型不存在”。
第三步:编写一个简单的 Controller 测试
```
@RestController
public class OllamaController {
private final ChatClient chatClient;
public OllamaController(ChatClient.Builder builder) {
// 直接基于注入的 builder 构建 ChatClient
// 框架会根据 yml 中的 spring.ai.ollama.* 配置自动把 OllamaChatModel 装好
this.chatClient = builder.build();
}
/**
* 简单一轮对话,测试本地部署的 Ollama 模型是否可用
*/
@GetMapping("/ollama/chat")
public String chat(@RequestParam String message) {
return chatClient.prompt()
.user(message)
.call()
.content();
}
}
```
启动 Spring Boot 项目,访问 http://localhost:8080/ollama/chat?message=你好,请介绍一下你自己,就能看到本地 Ollama 模型返回的内容。
Ollama常用命令
学会了如何把 Ollama 跑起来还不够,在实际使用过程中我们还经常需要做这些事情:看看本地都装了哪些模型、删掉一个不用的模型、查看模型详情、看看现在哪些模型正在跑等。这些都通过 Ollama 自带的命令行就能完成。
下面分类介绍最常用的几个命令。
1. 服务相关
``` # 启动 Ollama 服务(前台启动,会一直占着当前终端) ollama serve # 查看 Ollama 版本 ollama -v ```
ollama serve 是 Ollama 服务的启动命令,启动后默认监听 11434 端口,所有的模型调用都要通过它。生产环境一般会把它注册为 systemd 服务在后台运行。
2. 模型下载与运行
``` # 仅下载模型,不运行(适合提前把模型准备好) ollama pull qwen2.5:0.5b # 下载并运行模型(如果本地没有会自动先 pull) ollama run qwen2.5:0.5b ```
简单解释:
pull和run的区别:pull只下载不启动,run下载完会立刻进入对话;- 模型名称的格式是
模型名:版本/规模,比如qwen2.5:0.5b表示通义千问 2.5 的 0.5B 参数版本,不同规模文件大小差别非常大;
3. 模型管理
``` # 查看本地已安装的所有模型 ollama list # 查看当前正在运行(已加载到内存)的模型 ollama ps # 查看某个模型的详细信息(参数量、模板、license 等) ollama show qwen2.5:0.5b # 删除一个本地模型,释放磁盘空间 ollama rm qwen2.5:0.5b ```
这几个命令对应学生熟悉的 Linux/Docker 思维:
ollama list类似docker images,看本地装了什么ollama ps类似docker ps,看现在跑了什么ollama rm类似docker rmi,删除占空间的模型
4. 一份速查表
| 命令 | 作用 | 类比 |
|---|---|---|
ollama serve | 启动 Ollama 服务 | 启动数据库服务 |
ollama pull | 下载模型 | docker pull |
ollama run | 下载并运行 / 直接对话 | docker run |
ollama list | 列出本地所有模型 | docker images |
ollama ps | 查看正在运行的模型 | docker ps |
ollama show | 查看模型详情 | docker inspect |
ollama rm | 删除本地模型 | docker rmi |
熟记这张表,日常 Ollama 的运维基本都能搞定。
Spring AI Alibaba
Spring AI Alibaba 是专门用来构建 Agent 智能体应用的框架。
在学习 Spring AI 的时候,我们已经知道了如何通过 Spring AI 与大模型进行基本的交互。但是,仅仅能和大模型对话是远远不够的。在实际开发中,我们需要构建的是一个能够自主思考、自主行动的”智能体”(Agent)。
什么是 Agent?简单来说,站在应用程序的角度,Agent 就是一个能够自主决策并执行任务的程序。它不仅仅是和大模型聊天,而是能够:
- 自主决策:根据用户的需求,自己分解任务,自己判断应该做什么
- 落地执行:调用各种工具(比如查数据库、调接口)来完成任务
而 Spring AI Alibaba,就是帮助我们快速构建这种 Agent 应用的框架。
Spring AI Alibaba 分层架构
Spring AI Alibaba 从架构上包含三层:
[缺失图片] image/Spring%20AI%20Alibaba课件制作说明/high-level.png
我们从下往上来看这三层:
第一层:Augmented LLM(增强大模型层):其实对应的基本就是Spring AI
第二层:Graph Runtime(图运行时层):graph 是一个低级别的工作流和多代理协调框架,能够帮助开发者实现复杂的应用程序编排,Graph 是 Agent Framework 的底层运行时基座
第三层:Agent Framework(智能体框架层):是一个以 ReactAgent 为核心的 Agent 开发框架,使开发者能够构建具备自动上下文工程和人机交互等核心能力的Agent。
代码层面的Agent理解
这里需要再次解释一个概念——Agent,这个概念需要分场景来理解,在应用程序这个层面,Agent表示能够自主决策并执行任务的程序,在代码层面我们也会使用Agent,但是代码层面的Agent的含义就完全不同了。
怎么理解代码层面的一个Agent?
- 在代码层面它就是一个独立的执行单元,该单元具有与模型交互的能力,可以有自己的工具集,自己的会话记忆,自己所访问的知识库
- 同时,这个执行单元只会扮演单一角色,实现单一功能
如果我们要实现一个功能很简单的Agent应用程序,代码层面的一个Agent就可以实现,比如你要做一个”订单查询 Agent”,它能根据用户的自然语言、查询用户订单
但是,如果我们想实现一些具有复杂功能的Agent应用程序,在代码实现层面可能就需要多个代码层面的Agent实现不同角色的分工与协作。举一个例子来说明,比如说我想实现一个能够自主制定约会行程的Agent应用程序,在代码实现层面我可能需要扮演不同角色的多个独立执行单元来分工协作:
- 收集用户信息的执行单元:负责接收用户信息,并通过LLM提取出用户的约会时间,出发地点,日期,以及偏好等信息
- 专门负责查询天气的执行单元:负责根据约会日期查询天气,并给出穿搭建议
- 专门负责出行路线的执行单元:根据用户的出发地点和约会目标区域,给出合适的出行建议
- 负责收集活动信息的执行单元: 根据用户的偏好,查询目标区域,在约会日期合适的活动地点,就餐地点等
- 负责根据活动信息指定约会计划的执行单元:根据收集到的信息,整理,筛选,排列,得到完整的约会行程安排
以上的每一个独立的执行单元,都对应一个代码层面的Agent,多个代码层面Agent相互分工协作,而这些代码层面的Agent分工协作最终就实现了整个应用程序层面的Agent功能!
这些代码层面的多Agent分工协作,有一个专业的词来描述——Multi-agent,同时这多个Agent分工协作是有固定的先后执行顺序即执行流程的,这套多个代码层面的Agent的执行流程我们就称之为多Agent工作流(Multi-agnet workflow)
Spring AI Alibaba 分层使用场景
在理解了代码层面的Agent,以及Multi-agent, Multi-agent workflow,我们就可以理解Agentic Framework这一层和Graph Runtime这一层的使用场景了:
- Agentic Framework这一层提供了ReactAgent,一个ReactAgent对象就可以被当做是我们前面所说的一个代码层面的Agent,因此就已经可以具备构建一个基本的Agent应用的能力,同时,如果我们需要使用Multi-Agent,Agentic Framework也为我们提供了基础的多智能体实现方式
``` 比如"写一篇文章然后翻译成英文",这里面涉及两个角色:写作者和翻译者,两个角色可以用两个ReactAgent实现,然后可以利用Agentic Framework提供的基础workflow的实现,让写作ReactAgent先执行,翻译ReactAgent后执行,针对写作ReactAgent写出的内容来翻译 ```
- 但是,在企业级开发中,Multi-Agent有更为复杂灵活的协作场景,大量自定义逻辑,更复杂的状态控制,此时Agentic Framework这一层可能就满足不了我们的需求了,就只能使用Graph Runtime这一层的功能了
``` 如果你的多 Agent 协作场景非常复杂,有大量自定义逻辑、复杂的状态控制、需要精确控制执行流程,那就需要用 Graph Runtime。 举个例子:你要做一个"医疗问诊系统",流程是这样的: 1. 先识别患者是谁 2. 然后提取症状信息,如果信息不够还要追问 3. 追问完了要做症状标准化,让用户确认 4. 然后做急诊评估 5. 最后检索疾病证据,给出分诊建议 这种复杂度,就只能用 Graph Runtime 来自定义工作流了。 ```
所以接下来的课件内容,分成两个大的部分:
- Agent Framework:学习 ReactAgent 的使用,以及基于 ReactAgent 实现的基本 Multi-Agent 工作流
- Graph Runtime:学习复杂多 Agent 工作流的自定义实现,包括节点定义、状态定义、图的运行时状态流转(边的定义)等
—
Agent Framework
学习 Agent Framework ,主要包含包括两块内容:
- ReactAgent 的使用
- 基于 ReactAgent 实现的基本 Multi-Agent 工作流
在使用前需要引入如下依赖:
```
<dependency>
<groupId>com.alibaba.cloud.ai</groupId>
<artifactId>spring-ai-alibaba-agent-framework</artifactId>
<version>1.1.2.2</version>
</dependency>
<dependency>
<groupId>com.alibaba.cloud.ai</groupId>
<artifactId>spring-ai-alibaba-starter-dashscope</artifactId>
<version>1.1.2.2</version>
</dependency>
```
以及配置如下:
```
spring:
application:
name: alibaba
ai:
dashscope:
api-key: 你的apikey
base-url: https://dashscope.aliyuncs.com
chat:
options:
model: qwen-plus
```
ReactAgent
ReactAgent 是 Spring AI Alibaba 中构建 Agent 的核心类。接下来我们逐步学习 ReactAgent 的各项功能。
基本功能
最基本的 ReactAgent,只需要配置一个大模型,就可以实现与模型的简单交互(一轮对话)。
ReactAgent 的构建:
```
@Autowired
ChatModel chatModel;
// 构建 ReactAgent
ReactAgent agent = ReactAgent.builder()
.name("my_agent") // Agent 的名称
.model(chatModel) // 指定使用的大模型
.build();
```
ReactAgent 的调用:
```
// 调用 Agent,传入用户消息,返回模型的回复
AssistantMessage response = agent.call("你好,请介绍一下你自己");
System.out.println(response.getText());
```
就这么简单,两步就完成了:构建 Agent → 调用 Agent。
System Prompt & Instruction
在实际开发中,我们通常需要给 Agent 设定一个”角色”或者”行为规范”。这就需要用到 System Prompt 和 Instruction。
它们的区别是什么?
- System Prompt(系统提示词):定义 Agent 的角色和整体行为框架。比如”你是一个专业的医疗助手”。
- Instruction(指令):定义 Agent 在具体任务中的行为指导。比如”请用中文回答,回答控制在100字以内”。它更偏向具体的任务要求。
打个比方:System Prompt 就像是给员工定的”岗位职责”,而 Instruction 就像是给员工布置的”具体任务要求”。
ReactAgent 的构建:
```
ReactAgent agent = ReactAgent.builder()
.name("medical_agent")
.model(chatModel)
// 系统提示词:定义角色
.systemPrompt("你是一个专业的医疗健康助手,擅长回答健康相关的问题。")
// 指令:定义具体行为要求
.instruction("请用通俗易懂的语言回答用户的健康问题,回答控制在200字以内。" +
"如果问题超出你的能力范围,请建议用户去医院就诊。")
.build();
```
ReactAgent 的调用:
```
AssistantMessage response = agent.call("最近总是头疼是怎么回事?");
System.out.println(response.getText());
```
会话记忆
默认情况下,Agent 是”没有记忆”的。每次调用都是独立的,它不会记得之前说过什么。
但在实际场景中,我们需要 Agent 能记住对话上下文。比如用户说”我叫张三”,下次问”我叫什么”时,Agent 应该能回答出来。
要实现会话记忆,需要两个东西:
- BaseCheckpointSaver:负责保存对话状态(可以理解为”记忆存储器”), 但是BaseCheckPointSaver本身是个接口,他可以有不同的实现类比如MemorySaver可以将会话记忆保存在内存中
- threadId:会话 ID,用来区分不同的对话,从而实现会话记忆的隔离
ReactAgent 的构建:
```
ReactAgent agent = ReactAgent.builder()
.name("memory_agent")
.model(chatModel)
.instruction("你是一个友好的助手,请记住用户告诉你的信息。")
// 配置记忆存储器
.saver(new MemorySaver())
.build();
```
ReactAgent 的调用:
```
// 创建运行配置,指定会话 ID
RunnableConfig config = RunnableConfig.builder()
.threadId("conversation_001") // 同一个 threadId 表示同一个对话
.build();
// 第一次对话
agent.call("你好,我叫张三", config);
// 第二次对话(同一个 threadId,Agent 能记住之前的内容)
AssistantMessage response = agent.call("我叫什么名字?", config);
System.out.println(response.getText());
// 输出:你叫张三。
```
注意:如果换一个 threadId,Agent 就不会记得之前的对话了,因为那就相当于另一个会话”聊天窗口”。
工具调用
大模型本身是没有”动手能力”的,它只能生成文本。但是通过工具调用(Tool Calling),我们可以让 Agent 具备”执行操作”的能力。
比如:大模型不知道今天几号,但我们可以给它一个”获取当前时间”的工具,当用户问”今天几号”时,Agent 会自动调用这个工具获取答案。
定义工具:
使用 @Tool 注解将一个方法标记为工具:
```
public class DateTimeTools {
@Tool(description = "获取当前的日期和时间")
public String getCurrentDateTime() {
return LocalDateTime.now().toString();
}
}
```
注意:@Tool 注解的 description 非常重要,它告诉大模型这个工具是干什么的、什么时候该用。
ReactAgent 的构建:
```
ReactAgent agent = ReactAgent.builder()
.name("tool_agent")
.model(chatModel)
.instruction("你是一个有用的助手,可以帮用户查询时间和设置闹钟。")
// 传入工具实例
.tools(new DateTimeTools())
.saver(new MemorySaver())
.build();
```
ReactAgent 的调用:
```
RunnableConfig config = RunnableConfig.builder()
.threadId("tool_demo")
.build();
// Agent 会自动判断需要调用工具,获取当前时间后再回答
AssistantMessage response = agent.call("现在几点了?", config);
System.out.println(response.getText());
```
整个过程是这样的:用户提问 → Agent 判断需要调用工具 → 调用工具获取结果 → 将结果整合后回答用户。这一切都是自动完成的。
结构化输出
有时候我们希望模型 返回的不是一段自然语言文本,而是一个结构化的 JSON 对象,方便程序直接使用。
比如:让 Agent 从一段文字中提取联系人信息,返回 {name, email, phone} 格式的数据。
使用 outputType 方法即可实现:
ReactAgent 的构建:
```
// 定义输出结构(一个普通的 Java 类)
@Data
public class ContactInfo {
@JsonPropertyDescription("用户姓名")
private String name;
@JsonPropertyDescription("用户姓名")
private String email;
@JsonPropertyDescription("用户手机号")
private String phone;
// getter 和 setter 省略
}
// 构建 Agent,指定输出类型
ReactAgent agent = ReactAgent.builder()
.name("extractor_agent")
.model(chatModel)
.instruction("你是一个信息提取助手,请从用户提供的文本中提取联系人信息。")
// 指定结构化输出类型
.outputType(ContactInfo.class)
.saver(new MemorySaver())
.build();
```
ReactAgent 的调用:
```
RunnableConfig config = RunnableConfig.builder()
.threadId("extract_demo")
.build();
AssistantMessage response = agent.call(
"张三的邮箱是zhangsan@example.com,电话是13800138000", config);
System.out.println(response.getText());
// 输出:{"name": "张三", "email": "zhangsan@example.com", "phone": "13800138000"}
```
框架会自动将 Java 类转换为 JSON 格式说明,告诉大模型应该按什么格式输出。但是,模型返回的仍然只是一个JSON字符串,我们如何得到我们想要的对象呢?
``` BeanOutputConverter<T> converter = new BeanOutputConverter<>(outputType); ContactInfo contactInfo = converter.convert(message.getText()); ```
流式响应
ReactAgent也可以实现流式响应如下:
```
ReactAgent medicalAgent = ReactAgent.builder()
.name("medical_agent")
.model(chatModel)
// 系统提示词:定义角色
.systemPrompt("你是一个专业的医疗健康助手,擅长回答健康相关的问题。")
// 指令:定义具体行为要求
.instruction("请用通俗易懂的语言回答用户的健康问题,回答控制在200字以内。" +
"如果问题超出你的能力范围,请建议用户去医院就诊。")
.build();
RunnableConfig config = RunnableConfig.builder()
.threadId("extract_demo")
.build();
// 调用streamMesasages方法,Flux<Message> 中的每个Message对象都包含模型响应的一个chunk字符串
Flux<Message> messageFlux = medicalAgent.streamMessages("最近总是头疼是怎么回事?");
// 订阅消费每个Message输出其文本内容
messageFlux.map(message -> message.getText()).subscribe(System.out::println);
```
Multi-Agent(多智能体)
前面我们学了 ReactAgent,一个 ReactAgent 就是一个代码层面的 Agent,它扮演单一角色、实现单一功能。
但在实际开发中,很多任务不是一个角色能搞定的,会涉及 Multi-agent,以及Multi-agent workflow, 接下来,我们学习Spirng AI Alibaba Agent Framewok提供的四种基本workflow工作流实现:
- 顺序执行(Sequential Agent)
- 并行执行(Parallel Agent)
- 路由(LlmRoutingAgent)
状态对象与 Instruction 占位符
在讲具体的 Multi-Agent 之前,我们先要理解两个重要概念。
状态对象(OverAllState)
多智能体协作时,有一个核心问题:某个 Agent 产出的结果,另一个 Agent 怎么拿到?
思路很简单:产出结果的 Agent 把结果存入一个公共对象,另一个 Agent 从这个公共对象中取出它需要的结果。
这个被所有 Agent 共享的公共对象,就是所谓的状态对象(OverAllState)。你可以把它理解为一个公共的 Map,每个 Agent 都可以往里面存数据,也可以从里面取数据。
Instruction 占位符
理解了状态对象之后,占位符就很好理解了。
占位符的格式是 {keyName},我们在定义Instruction时,可以在Instruction中使用占位符,在运行时它会自动从状态对象中查找对应的值并替换占位符。
ReactAgent配合 outputKey() 方法可以将Agent执行的结果存入状态对象中:
outputKey("article"):表示这个 Agent 的输出结果会以 “article” 为 key 存入状态对象{article}:在另一个 Agent 的 instruction 中使用,表示从状态对象中取出 key 为 “article” 的值
```
// 第一个 Agent:输出结果存入状态对象,key 为 "article"
ReactAgent writerAgent = ReactAgent.builder()
.name("writer_agent")
.instruction("你是一个作家。请根据用户的提问进行回答:{input}。")
.outputKey("article") // 输出存入状态,key = "article"
.build();
// 第二个 Agent:通过 {article} 占位符获取第一个 Agent 的输出
ReactAgent reviewerAgent = ReactAgent.builder()
.name("reviewer_agent")
.instruction("请对文章进行评审修正:\n{article},最终返回修正后的文章")
.outputKey("reviewed_article")
.build();
```
其中 {input} 是一个特殊的占位符,它代表用户的原始输入。
顺序执行(Sequential Agent)
在顺序执行模式中,多个 ReactAgent 按预定义的顺序依次执行。每个 Agent 的输出成为下一个 Agent 的输入。
流程:Agent A 处理 → Agent A 的输出传给 Agent B → Agent B 处理 → 最终结果
```
graph LR
subgraph SequentialAgent
direction LR
A[sub_agent_1<br/>writerAgent] --> B[sub_agent_2<br/>reviewerAgent]
end
Input([Input]) --> A
B --> Output([Output])
State[(OverAllState)]
A -.->|"写入 article"| State
State -.->|"读取 {article}"| B
B -.->|"写入 reviewed_article"| State
```
```
import com.alibaba.cloud.ai.graph.agent.flow.agent.SequentialAgent;
// 创建写作 Agent
ReactAgent writerAgent = ReactAgent.builder()
.name("writer_agent")
.model(chatModel)
.description("专业写作Agent")
.instruction("你是一个知名的作家。请根据用户的提问进行回答:{input}。")
.outputKey("article") // 输出存入状态,key = "article"
.build();
// 创建评审 Agent
ReactAgent reviewerAgent = ReactAgent.builder()
.name("reviewer_agent")
.model(chatModel)
.description("专业评审Agent")
.instruction("你是一个评论家。请对文章进行评审修正:{article}。" +
"最终只返回修改后的文章。")
.outputKey("reviewed_article") // 输出存入状态,key = "reviewed_article"
.build();
// 创建顺序 Agent,可以让 subAgents 列表顺序执行
SequentialAgent blogAgent = SequentialAgent.builder()
.name("blog_agent")
.description("先写文章,再评审修改")
// 给顺序执行的Agent设置子Agnet
.subAgents(List.of(writerAgent, reviewerAgent)) // 按顺序执行
.build();
// 调用invoke执行SequentialAgent
Optional<OverAllState> result = blogAgent.invoke("帮我写一个100字左右的散文");
if (result.isPresent()) {
OverAllState state = result.get(); // 这就是状态对象
// 从状态对象中取出评审后的文章
state.value("reviewed_article").ifPresent(article -> {
System.out.println("评审后文章: " + ((AssistantMessage) article).getText());
});
}
```
注意三个关键点:
{input}占位符:获取用户原始输入outputKey("article"):第一个 Agent 的输出存入状态{article}占位符:第二个 Agent 从状态中读取第一个 Agent 的输出- 最终返回的
OverAllState对象就是状态对象,包含了所有 Agent 的输出
并行执行(Parallel Agent)
在并行执行模式中,多个 Agent 同时处理相同的输入,它们的结果被收集并合并。
流程:输入同时发给所有 Agent → 所有 Agent 并行处理 → 结果合并
```
graph TB
subgraph ParallelAgent
direction TB
A[sub_agent_1<br/>散文Agent]
B[sub_agent_2<br/>诗歌Agent]
C[sub_agent_3<br/>剧本Agent]
end
Input([Input]) --> A
Input --> B
Input --> C
A --> Output_1([Output_1])
B --> Output_2([Output_2])
C --> Output_3([Output_3])
State[(OverAllState)]
A -.->|"写入 prose_result"| State
B -.->|"写入 poem_result"| State
C -.->|"写入 script_result"| State
```
```
import com.alibaba.cloud.ai.graph.agent.flow.agent.ParallelAgent;
// 散文 Agent
ReactAgent proseWriterAgent = ReactAgent.builder()
.name("prose_writer_agent")
.model(chatModel)
.description("专门写散文的AI助手")
.instruction("你是一个散文作家。用户给你的主题是:{input},请创作一篇100字左右的散文。")
.outputKey("prose_result")
.build();
// 诗歌 Agent
ReactAgent poemWriterAgent = ReactAgent.builder()
.name("poem_writer_agent")
.model(chatModel)
.description("专门写现代诗的AI助手")
.instruction("你是一个现代诗人。用户给你的主题是:{input},请创作一首现代诗。")
.outputKey("poem_result")
.build();
// 创建并行 Agent
ParallelAgent parallelAgent = ParallelAgent.builder()
.name("parallel_creative_agent")
.description("并行执行多个创作任务")
.mergeOutputKey("merged_results") // 合并结果的 key
.subAgents(List.of(proseWriterAgent, poemWriterAgent))
.mergeStrategy(new ParallelAgent.DefaultMergeStrategy())
.build();
// 调用
Optional<OverAllState> result = parallelAgent.invoke("以'西湖'为主题");
if (result.isPresent()) {
OverAllState state = result.get();
state.value("prose_result").ifPresent(r -> System.out.println("散文: " + r));
state.value("poem_result").ifPresent(r -> System.out.println("诗歌: " + r));
}
```
路由(LlmRoutingAgent)
在路由模式中,使用大模型动态决定将请求路由到哪个子 Agent。适合需要智能选择不同专家的场景。
流程:路由 Agent 接收输入 → LLM 分析并选择最合适的子 Agent → 选中的 Agent 处理请求
```
graph TB
subgraph LlmRoutingAgent
direction TB
Router{LLM Router<br/>路由决策}
A[sub_agent_1<br/>数学Agent]
B[sub_agent_2<br/>文学Agent]
C[sub_agent_3<br/>编程Agent]
Router -->|"数学问题"| A
Router -->|"文学问题"| B
Router -->|"编程问题"| C
end
Input([Input]) --> Router
A --> Output([Output])
B --> Output
C --> Output
State[(OverAllState)]
Output -.->|"写入选中Agent的outputKey"| State
```
```
import com.alibaba.cloud.ai.graph.agent.flow.agent.LlmRoutingAgent;
// 写作 Agent
ReactAgent writerAgent = ReactAgent.builder()
.name("writer_agent")
.model(chatModel)
.description("擅长创作各类文章,包括散文、诗歌等文学作品")
.instruction("你是一个知名的作家,请根据用户的提问进行回答。")
.outputKey("writer_output")
.build();
// 翻译 Agent
ReactAgent translatorAgent = ReactAgent.builder()
.name("translator_agent")
.model(chatModel)
.description("擅长将文章翻译成各种语言")
.instruction("你是一个专业的翻译家,能够准确地将文章翻译成目标语言。")
.outputKey("translator_output")
.build();
// 创建路由 Agent
LlmRoutingAgent routingAgent = LlmRoutingAgent.builder()
.name("content_routing_agent")
.description("根据用户需求智能路由到合适的专家Agent")
.model(chatModel) // 路由决策也需要大模型
.subAgents(List.of(writerAgent, translatorAgent))
.build();
// 调用 - LLM 会自动选择最合适的 Agent
Optional<OverAllState> result1 = routingAgent.invoke("帮我写一篇关于春天的散文");
// LLM 会路由到 writerAgent
Optional<OverAllState> result2 = routingAgent.invoke("请将以下内容翻译成英文:春暖花开");
// LLM 会路由到 translatorAgent
```
路由的准确性取决于子 Agent 的 description 描述是否清晰明确。描述越清楚,LLM 越能准确判断该路由到哪个 Agent。
—
Graph Runtime
在学习了基于 ReactAgent 实现的 3种基本 Multi-Agent 工作流之后,够用了吗?
当然不是。因为在实际开发中,还有更为复杂的 Multi-Agent 协作场景。比如前面提到的医疗问诊系统,流程中有条件分支、循环、中断等待等复杂逻辑,Sequential、Parallel 这些基本模式根本搞不定。
那更复杂的 Multi-Agent 工作流怎么实现呢?没关系,我们可以自定义任何复杂的工作流,这是 Spring AI Alibaba 所支持的。
图的概念
假设我们有这样一个复杂的工作流:
用户提交一个问题 → 首先对问题进行预处理 → 然后判断问题类型:如果是简单问题,直接回答;如果是复杂问题,先做深度分析,分析完再生成回答 → 最后对回答进行质量验证:如果质量合格就输出,不合格就重新生成回答。
这个流程画成流程图大概是这样的:
```
graph LR
START([开始]) --> 预处理
预处理 --> 判断类型
判断类型 -->|简单| 直接回答
判断类型 -->|复杂| 深度分析
深度分析 --> 生成回答
直接回答 --> 质量验证
生成回答 --> 质量验证
质量验证 -->|合格| END([结束])
质量验证 -->|不合格| 生成回答
```
很显然,前面的顺序执行、并行执行、路由、监督者都无法直接实现这种流程。
但是,如果仔细观察这个流程图,你会发现:它其实就是一个有向图!
- 每个处理步骤就是图中的一个节点
- 步骤之间的箭头就是图中的边
- 有些边是无条件的(预处理完一定到判断类型)
- 有些边是有条件的(根据判断结果走不同的路)
于是,实现一个复杂的 Multi-Agent 工作流,就变成了实现一个有向图。而 Graph Runtime 中的 StateGraph 对象,就可以表示一个任意复杂度的有向图!
图的实现
可以通过 new StateGraph() 得到一个 StateGraph 对象:
``` StateGraph graph = new StateGraph(); ```
但很显然,这个对象现在什么都没有。要想让它真的能表示一个复杂的图,还需要给它添加三个关键要素:
- 状态(State):存储整个图运行过程中需要用到的数据(比如用户消息),或者图中产生的数据。具体实现上它是一个
Map。 - 节点(Node):图中的每个节点是一个执行逻辑单元,完成某一项特定的功能。它接受当前 State 作为输入,执行某些操作(如调用 LLM 或自定义逻辑),并将结果更新到 State 中,供其他节点使用。
- 边(Edge):定义节点间的控制流,实现分支逻辑。它决定了一个节点执行完之后,下一个要执行的节点是谁。
简而言之:Node 完成工作,Edge 告诉下一步该做什么,State 存储图执行过程中的数据。
接下来我们分别学习节点和边的实现。
节点(Node)
节点通过实现 NodeAction 接口来定义:
```
public interface NodeAction {
// 入参 state:当前图的状态对象,可以从中读取数据
// 返回值 Map:要更新到状态中的数据(key-value 形式)
Map<String, Object> apply(OverAllState state) throws Exception;
}
```
下面是一个 NodeAction 的实现示例:
```
public class TextProcessorNode implements NodeAction {
@Override
public Map<String, Object> apply(OverAllState state) throws Exception {
// 1. 从状态中获取数据
String input = state.value("input", "").toString();
// 2. 执行节点自己的逻辑(这里可以使用 ReactAgent、调接口等)
String processedText = input.toUpperCase().trim();
// 3. 返回要更新到状态中的数据
Map<String, Object> result = new HashMap<>();
result.put("processed_text", processedText);
return result; // 这些数据会被合并到状态对象中
}
}
public class AnalyzerNode implements NodeAction {
@Override
public Map<String, Object> apply(OverAllState state) throws Exception {
// 1. 从状态中获取数据
String text = state.value("processed_text", "").toString();
// 2. 执行节点自己的逻辑(这里可以使用 ReactAgent、调接口等)
String type = text.contains("JAVA") ? "JAVA" : "OTHER";
// 3. 返回要更新到状态中的数据
Map<String, Object> result = new HashMap<>();
result.put("type", type);
return result; // 这些数据会被合并到状态对象中
}
}
```
将节点加入图:
```
// 使用 node_async 包装后添加到图中
import static com.alibaba.cloud.ai.graph.action.AsyncNodeAction.node_async;
graph.addNode("processor", AsyncNodeAction.node_async(new TextProcessorNode()));
graph.addNode("analyzer", AsyncNodeAction.node_async(new AnalyzerNode()));
```
AsyncNodeAction.node_async 是一个工具方法,将同步的 NodeAction 包装为异步执行。这是添加节点的标准方式。
特殊节点 START 和 END:
一个图中本身就包含两个特殊节点:
StateGraph.START:图的起点,表示工作流从这里开始StateGraph.END:图的终点,表示工作流到这里结束
我们不需要手动创建它们,只需要在定义边的时候引用即可。
边(Edge)
普通边:
普通边表示无条件的流转,即”A 节点执行完,一定执行 B 节点”。
```
// 语法:addEdge(源节点名称, 目标节点名称)
graph.addEdge(StateGraph.START, "processor"); // 开始 → 预处理
graph.addEdge("processor", "analyzer"); // 预处理 → 分析
graph.addEdge("analyzer", StateGraph.END); // 分析 → 结束
```
条件边:
条件边表示根据某个条件决定下一步走哪个节点,即”A 节点执行完,根据条件决定走 B 还是走 C”。
```
import static com.alibaba.cloud.ai.graph.action.AsyncEdgeAction.edge_async;
// 语法:addConditionalEdges(源节点, 条件函数, 路由映射)
graph.addConditionalEdges(
"condition_node", // 源节点:从哪个节点出发
// 条件函数:根据状态返回一个字符串,表示走哪条路
AsyncEdgeAction.edge_async(state -> {
String type = state.value("question_type", "simple").toString();
return type; // 返回 "simple" 或 "complex"
}),
// 路由映射:条件函数返回值 → 目标节点
Map.of(
"simple", "direct_answer", // 返回 "simple" → 走 direct_answer 节点
"complex", "deep_analysis" // 返回 "complex" → 走 deep_analysis 节点
)
);
```
并行边:
一个节点可以有多个出边,表示同时执行多个后续节点:
```
// 一个节点同时连接多个后续节点(并行执行)
graph.addEdge("start_node", List.of("agent1", "agent2", "agent3"));
// 多个并行节点汇聚到一个节点
graph.addEdge(List.of("agent1", "agent2", "agent3"), "aggregator");
```
图的运行
图定义好之后,需要编译并运行。
编译图:
``` // 1. 创建编译配置 CompileConfig compileConfig = CompileConfig.builder().build(); // 2. 编译图,得到可执行的 CompiledGraph CompiledGraph compiledGraph = graph.compile(compileConfig); ```
对于一个图而言,在其执行过程中也有很多状态需要存储,比如当前执行到了哪个节点,下一个要执行的节点是谁,当前整个图共用的状态对象OverAllSate…., 在ReactAgent中我们使用BaseCheckpointSaver实现会话记忆的存储,同样在图中我们还可以使用它实现整个图状态的存储!
MemorySaver 是最简单的 CheckPointer 实现,它将状态保存在内存中(当然除了MemorySaver之外,还有RedisSaver,MySqlSaver)。
为了实现图状态的存储,我们需要在编译图的时候指定一个BaseCheckpointSaver对象:
```
// 构造saver参数
SaverConfig saverConfig = SaverConfig.builder()
// 基于内存保存图的运行状态
.register(new MemorySaver())
.build();
// 设置saver参数
CompileConfig compileConfig = CompileConfig.builder().saverConfig(saverConfig).build();
// 设置compile参数
CompiledGraph graph = stateGraph.compile(compileConfig);
// 创建图
CompiledGraph compiledGraph = graph.compile(compileConfig);
```
构建 RunnableConfig:
RunnableConfig 包含运行时的配置信息,最重要的是 threadId,它就相当于会话Id,用来区分在不同会话中执行的图,以及图所保存的状态:
```
// threadId 的含义:标识一次完整的工作流执行
RunnableConfig runnableConfig = RunnableConfig.builder()
.threadId("workflow-001")
.build();
```
运行图:
```
// 准备输入数据
Map<String, Object> input = Map.of("input", "用户的问题内容");
// 执行图
Optional<OverAllState> result = compiledGraph.invoke(input, runnableConfig);
// 获取结果
result.ifPresent(state -> {
// 从状态对象中获取图中节点存储到其中的数据,数据都是key-value形式存储,因此指定key即可获取对应的数据
System.out.println("输出: " + state.value("output").orElse("无"));
});
```
这里需要注意的是:
- 图的每一次执行,在图的内部都会自动产生一个状态对象(OverAllState对象),该状态对象被图中的各个节点所共用
- 图的每一次执行,所产生的状态对象(OverAllstate对象),都是不同的对象(除非是HIL的中断恢复)
综合案例
下面给出一个完整的 Web 版综合案例,包含三个节点:问题分类节点、技术回答节点和生活回答节点。
1. 定义节点:
```
// 问题分类节点:判断问题是技术类还是生活类
public class ClassifyNode implements NodeAction {
@Override
public Map<String, Object> apply(OverAllState state) throws Exception {
String input = state.value("input", "").toString();
// 简单判断逻辑(实际可以调用 LLM)
String category = input.contains("代码") || input.contains("编程")
? "tech" : "life";
return Map.of("category", category);
}
}
// 技术回答节点:使用技术领域 Agent 专业地回答
public class TechAnswerNode implements NodeAction {
private final ReactAgent techAgent;
public TechAnswerNode(ChatModel chatModel) {
this.techAgent = ReactAgent.builder()
.name("tech_agent")
.model(chatModel)
.instruction("你是一个技术专家,请专业地回答用户的技术问题:{input}")
.outputKey("answer")
.build();
}
@Override
public Map<String, Object> apply(OverAllState state)
throws Exception {
String input = state.value("input", "").toString();
AssistantMessage response = techAgent.call(input);
return Map.of("answer", response.getText());
}
}
// 生活回答节点:使用生活领域 Agent 亲切地回答
public class LifeAnswerNode implements NodeAction {
private final ReactAgent lifeAgent;
public LifeAnswerNode(ChatModel chatModel) {
this.lifeAgent = ReactAgent.builder()
.name("life_agent")
.model(chatModel)
.instruction("你是一个生活顾问,请亲切地回答用户的生活问题:{input}")
.outputKey("answer")
.build();
}
@Override
public Map<String, Object> apply(OverAllState state)
throws Exception {
String input = state.value("input", "").toString();
AssistantMessage response = lifeAgent.call(input);
return Map.of("answer", response.getText());
}
}
```
2. 构建图并配置为 Spring Bean:
```
@Configuration
public class WorkflowConfig {
@Bean
public CompiledGraph qaWorkflow(ChatModel chatModel) throws Exception {
// 构建图
StateGraph graph = new StateGraph();
// 添加节点
graph.addNode("classify", node_async(new ClassifyNode()));
graph.addNode("tech_answer", node_async(new TechAnswerNode(chatModel)));
graph.addNode("life_answer", node_async(new LifeAnswerNode(chatModel)));
// 添加边
graph.addEdge(StateGraph.START, "classify");
// 条件边:classify 执行完后,根据 category 决定走技术节点还是生活节点
graph.addConditionalEdges(
"classify",
AsyncEdgeAction.edge_async(state ->
state.value("category", "life").toString()), // 返回 "tech" 或 "life"
Map.of(
"tech", "tech_answer", // 技术问题 → 技术回答节点
"life", "life_answer" // 生活问题 → 生活回答节点
)
);
graph.addEdge("tech_answer", StateGraph.END);
graph.addEdge("life_answer", StateGraph.END);
// 编译(配置 MemorySaver 以支持后续 Human-In-Loop 扩展)
return graph.compile(CompileConfig.builder()
.saverConfig(new MemorySaver())
.build());
}
}
```
3. Controller 层:
```
@RestController
@RequestMapping("/api/qa")
public class QaController {
private final CompiledGraph qaWorkflow;
public QaController(CompiledGraph qaWorkflow) {
this.qaWorkflow = qaWorkflow;
}
@PostMapping
public Map<String, Object> ask(@RequestBody Map<String, String> request) {
String question = request.get("question");
String threadId = request.getOrDefault("threadId", UUID.randomUUID().toString());
RunnableConfig config = RunnableConfig.builder()
.threadId(threadId)
.build();
Optional<OverAllState> result = qaWorkflow.invoke(
Map.of("input", question), config);
return result.map(state -> Map.of(
"answer", state.value("answer").orElse("无法回答"),
"category", state.value("category").orElse("unknown"),
"threadId", threadId
)).orElse(Map.of("error", "执行失败"));
}
}
```
自定义 FlowAgent
上面的综合案例直接把 StateGraph 编译成了 CompiledGraph,然后在 Controller 中调用。这个写法没有问题,但它只是一个工作流对象,还不是一个 Agent。
如果我们希望把任意自定义工作流也封装成 Agent,就可以继承 FlowAgent,并在 buildSpecificGraph 方法中返回自己定义的 StateGraph。这样一来,这个工作流就具备了 Agent 的调用入口,可以像其他 Agent 一样被 使用
还是以前面的问答 workflow 为例,等效的 FlowAgent 可以这样实现:
```
package com.example.qa;
import com.alibaba.cloud.ai.graph.CompileConfig;
import com.alibaba.cloud.ai.graph.StateGraph;
import com.alibaba.cloud.ai.graph.action.AsyncEdgeAction;
import com.alibaba.cloud.ai.graph.action.AsyncNodeAction;
import com.alibaba.cloud.ai.graph.agent.Agent;
import com.alibaba.cloud.ai.graph.agent.flow.agent.FlowAgent;
import com.alibaba.cloud.ai.graph.agent.flow.builder.FlowGraphBuilder;
import com.alibaba.cloud.ai.graph.exception.GraphStateException;
import org.springframework.ai.chat.model.ChatModel;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
public class QaWorkflowAgent extends FlowAgent {
private final ChatModel chatModel;
public QaWorkflowAgent(ChatModel chatModel, CompileConfig compileConfig) {
super(
"qa_workflow_agent",
"根据问题类型路由到技术回答节点或生活回答节点",
compileConfig,
List.<Agent>of()
);
this.chatModel = chatModel;
}
@Override
protected StateGraph buildSpecificGraph(FlowGraphBuilder.FlowGraphConfig config) throws GraphStateException {
StateGraph graph = new StateGraph(config.getName(), HashMap::new);
graph.addNode("classify", AsyncNodeAction.node_async(new ClassifyNode()));
graph.addNode("tech_answer", AsyncNodeAction.node_async(new TechAnswerNode(chatModel)));
graph.addNode("life_answer", AsyncNodeAction.node_async(new LifeAnswerNode(chatModel)));
graph.addEdge(StateGraph.START, "classify");
graph.addConditionalEdges(
"classify",
AsyncEdgeAction.edge_async(state ->
state.value("category", "life").toString()),
Map.of(
"tech", "tech_answer",
"life", "life_answer"
)
);
graph.addEdge("tech_answer", StateGraph.END);
graph.addEdge("life_answer", StateGraph.END);
return graph;
}
}
```
注意:buildSpecificGraph 方法只负责定义图结构并返回 StateGraph,不要在这个方法里手动调用 graph.compile(...)。FlowAgent 在执行 invoke 或 stream 时,会自动完成图初始化和编译;编译时使用的 MemorySaver、中断配置等,则来自构造 QaWorkflowAgent 时传入的 CompileConfig。
然后把这个自定义 Agent 注册为 Spring Bean:
```
package com.example.qa;
import com.alibaba.cloud.ai.graph.CompileConfig;
import com.alibaba.cloud.ai.graph.checkpoint.config.SaverConfig;
import com.alibaba.cloud.ai.graph.checkpoint.savers.MemorySaver;
import org.springframework.ai.chat.model.ChatModel;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class QaAgentConfig {
@Bean
public QaWorkflowAgent qaWorkflowAgent(ChatModel chatModel) {
CompileConfig compileConfig = CompileConfig.builder()
.saverConfig(SaverConfig.builder().register(new MemorySaver()).build())
.build();
return new QaWorkflowAgent(chatModel, compileConfig);
}
}
```
调用方式也从调用 CompiledGraph 变成了调用 Agent:
```
package com.example.qa;
import com.alibaba.cloud.ai.graph.OverAllState;
import com.alibaba.cloud.ai.graph.RunnableConfig;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.Map;
import java.util.Optional;
import java.util.UUID;
@RestController
@RequestMapping("/api/qa-agent")
public class QaAgentController {
private final QaWorkflowAgent qaWorkflowAgent;
public QaAgentController(QaWorkflowAgent qaWorkflowAgent) {
this.qaWorkflowAgent = qaWorkflowAgent;
}
@PostMapping
public Map<String, Object> ask(@RequestBody Map<String, String> request) {
String question = request.get("question");
String threadId = request.getOrDefault("threadId", UUID.randomUUID().toString());
RunnableConfig config = RunnableConfig.builder()
.threadId(threadId)
.build();
Optional<OverAllState> result = qaWorkflowAgent.invoke(
Map.of("input", question), config);
return result.map(state -> Map.of(
"answer", state.value("answer").orElse("无法回答"),
"category", state.value("category").orElse("unknown"),
"threadId", threadId
)).orElse(Map.of("error", "执行失败"));
}
}
```
这两种写法的图结构是等效的:
```
graph TD
START["START"] --> CLASSIFY["classify"]
CLASSIFY --> TECH["tech_answer"]
CLASSIFY --> LIFE["life_answer"]
TECH --> END["END"]
LIFE --> END
```
区别在于抽象层次不同:
- 直接暴露
CompiledGraph:适合只在当前服务内部执行一个工作流 - 封装成
FlowAgent:适合把这个工作流当成一个 Agent 复用,或者作为更大 Multi-Agent 系统中的一个子 Agent
—
Human-In-Loop(人工介入)
Human-In-Loop(人工介入)是指在工作流执行过程中,暂停执行,等待人工确认或输入后再继续。
什么场景需要 Human-In-Loop?
- 敏感操作确认:比如 Agent 要删除数据,需要用户确认”你确定要删除吗?”
- 信息补充:比如 Agent 需要用户提供更多信息才能继续
- 质量审核:比如 Agent 生成了一个方案,需要人工审核通过后才执行
要实现 Human-In-Loop,核心是两个步骤:
- 断点(Interrupt):图执行到某个节点时暂停,保存当前状态,并返回当前图的状态对象(OverAllState对象)
- 恢复执行(Resume):收到用户输入后,从断点处恢复执行(接着之前保存的状态继续执行)
这就像你在玩游戏时按了暂停键(断点),等你准备好了再按继续(恢复)。底层”存档”机制其实是由BaseCheckPointer实现的,它保证了暂停后状态不会丢失。
我们可以基于InterruptableAction接口 实现 Human-In-Loop
具体实现
使用 InterruptableAction 方式实现 Human-In-Loop,核心思路是:专门定义一个”等待节点”,该节点同时实现 InterruptableAction 接口和 AsyncNodeAction 接口。
InterruptableAction接口提供interrupt()方法 —— 决定是否中断NodeAction接口提供apply()方法 —— 恢复后执行的逻辑
InterruptableAction 接口定义(源码):
```
public interface InterruptableAction {
// 返回 Optional.empty() → 不中断,继续执行 apply()
// 返回 InterruptionMetadata → 触发中断,图暂停
Optional<InterruptionMetadata> interrupt(String nodeId, OverAllState state, RunnableConfig config);
}
```
用来实现中断的这个专门的节点定义如下:
```
public class WaitNode implements AsyncNodeAction, InterruptableAction {
// 返回 Optional.empty() → 不中断,继续执行 apply()
// 返回 InterruptionMetadata → 触发中断,图暂停
@Override
public Optional<InterruptionMetadata> interrupt(String nodeId,
OverAllState state,
RunnableConfig config) {
// 有恢复标记 → 说明这是中断后的恢复执行
if (config.metadata(RunnableConfig.HUMAN_FEEDBACK_METADATA_KEY).isPresent()) {
// 关键:判断到恢复标记后,立刻删除这个标记
// 否则本次恢复流程里如果再次遇到 WaitNode,会被误判为恢复执行,从而跳过新的中断
config.metadata().ifPresent(metadata ->
metadata.remove(RunnableConfig.HUMAN_FEEDBACK_METADATA_KEY));
return Optional.empty();
}
// 无恢复标记 → 首次进入,触发中断
return Optional.of(InterruptionMetadata.builder(nodeId, state).build());
}
// 节点正常执行时的apply方法,返回的map中的键值对会被放入State对象
@Override
public CompletableFuture<Map<String, Object>> apply(OverAllState state) {
// 恢复后执行:传递状态即可,图会根据边定义流转到下一个节点,也可以什么都不做
return CompletableFuture.completedFuture(Map.of());
}
}
```
这个等待节点的职责很单一:首次进入时中断图的执行,恢复后由该节点出发,继续流转到下一个节点。
注意这里删除恢复标记非常重要。主要是为了不影响后续其他节点的中断!
定义完 WaitNode 之后,我们先看一个非常简单的例子,不通过 Web 接口,只在普通方法中调用图。
先定义三个节点,每个节点单独负责一件事:
```
public class AskNameNode implements NodeAction {
@Override
public Map<String, Object> apply(OverAllState state) {
System.out.println("请用户输入姓名");
return Map.of("question", "你的姓名是什么?");
}
}
public class SayHelloNode implements NodeAction {
@Override
public Map<String, Object> apply(OverAllState state) {
String name = state.value("name", "同学").toString();
return Map.of("answer", "你好," + name);
}
}
```
其中 wait_input 节点直接使用前面定义的 WaitNode。接着把这三个节点串成一张图:
```
StateGraph graph = new StateGraph();
graph.addNode("ask_name", AsyncNodeAction.node_async(new AskNameNode()));
graph.addNode("wait_input", AsyncNodeAction.node_async(new WaitNode()));
graph.addNode("say_hello", AsyncNodeAction.node_async(new SayHelloNode()));
graph.addEdge(StateGraph.START, "ask_name");
graph.addEdge("ask_name", "wait_input");
graph.addEdge("wait_input", "say_hello");
graph.addEdge("say_hello", StateGraph.END);
CompiledGraph compiledGraph = graph.compile(CompileConfig.builder()
.saverConfig(SaverConfig.builder().register(new MemorySaver()).build())
.build());
```
第一次调用时,还没有设置恢复标记(因为还未触发中断)。图执行到 wait_input 时,WaitNode.interrupt() 会返回 InterruptionMetadata,图在这里暂停:
```
String threadId = "demo-thread-001";
RunnableConfig firstConfig = RunnableConfig.builder()
.threadId(threadId)
.build();
// 调用invokeAndGetOutput方法获取本地图执行的最后一个节点信息,NodeOutput对象中也会包含OverAllState对象
Optional<NodeOutput> firstOutput = compiledGraph.invokeAndGetOutput(Map.of(), firstConfig);
if (firstOutput.isPresent() && !firstOutput.get().isEND()) {
System.out.println("发生中断,中断节点:" + firstOutput.get().node());
System.out.println("当前状态:" + firstOutput.get().state().data());
}
```
这里使用的是 invokeAndGetOutput,不是 invoke。二者的区别是:
invoke(...)返回Optional,你只能拿到状态对象invokeAndGetOutput(...)返回Optional,除了能拿到状态对象,还能知道最后停在哪个节点
NodeOutput 的核心字段包括:
node():最后输出来自哪个节点。如果中断发生在wait_input,这里就是wait_inputstate():当前图的状态对象,也就是中断或结束时保存下来的 OverAllStateisEND():是否已经到达END节点。如果isEND()为false,通常说明这次调用停在了某个业务节点或等待节点;对于 WaitNode 场景,就可以据此判断发生了中断
因此,第一次调用可以通过 !firstOutput.get().isEND() 判断发生了中断。
第一次调用的执行过程如下:
```
graph TD
START["开始"] --> ASK["ask_name 节点"]
ASK --> WAIT["wait_input 节点"]
WAIT --> OUTPUT["返回 NodeOutput"]
OUTPUT --> PAUSE["图暂停"]
```
此时最后一个 NodeOutput 来自 wait_input,并且 isEND() 为 false。这说明图没有跑到 END 节点,而是在 wait_input 这里发生了中断。
第二次调用时,使用同一个 threadId,并设置 resume() 标记,同时通过 addStateUpdate 把用户输入合并到状态对象中:
```
RunnableConfig resumeConfig = RunnableConfig.builder()
.threadId(threadId)
// 添加恢复标记
.resume()
// 在OverAllState中合并用户输入的信息
.addStateUpdate(Map.of("name", "张三"))
.build();
Optional<NodeOutput> resumeOutput = compiledGraph.invokeAndGetOutput(Map.of(), resumeConfig);
if (resumeOutput.isPresent() && resumeOutput.get().isEND()) {
System.out.println("执行完成:" + resumeOutput.get().state().value("answer", ""));
}
```
第二次恢复的执行过程如下:
```
graph TD
USER["用户输入 张三"] --> CONFIG["构造 resumeConfig"]
CONFIG --> WAIT["回到 wait_input 节点"]
WAIT --> APPLY["执行 wait_input.apply"]
APPLY --> HELLO["say_hello 节点"]
HELLO --> END["END 节点"]
```
第二次恢复调用时,图从 checkpoint 恢复到 wait_input,WaitNode 识别到恢复标记后先删除该标记,再返回 Optional.empty(),然后执行 apply(),继续流转到 say_hello 和 END。
综合案例:智能订单处理
下面用一个完整案例来演示上述原理。场景:
- 用户提交一个订单请求
- 系统先做风险评估(risk_assess 节点)
- 如果是低风险订单 → 不中断,直接进入处理
- 如果是高风险订单 → 中断,等待人工审核
- 人工审核后 → 进入订单处理节点(不同于风险评估节点)
关键点:
- 有选择的中断:低风险不中断,高风险才中断
- 恢复后跳转到其他节点:中断发生在 wait_review 节点,恢复后流转到 process_order 节点
案例流程图:
```
graph TD
START([开始]) --> risk[risk_assess<br/>风险评估节点]
risk -->|"低风险<br/>outcome=CONTINUE"| process[process_order<br/>订单处理节点]
risk -->|"高风险<br/>outcome=NEED_REVIEW"| wait[wait_review<br/>等待审核节点<br/>实现 InterruptableAction]
wait -->|"首次进入<br/>interrupt返回metadata<br/>图中断"| PAUSE([图暂停<br/>返回审核请求给前端])
PAUSE -->|"人工审核后恢复"| wait
wait -->|"恢复进入<br/>删除恢复标记<br/>interrupt返回empty<br/>执行apply透传状态"| process
process --> END([结束])
```
以下代码可以直接复制到 Spring Boot 项目中运行(需要已引入 spring-ai-alibaba-agent-framework 依赖)。
1. 风险评估结果类:
```
package com.example.order;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@NoArgsConstructor
public class RiskResult {
private String riskLevel; // "HIGH" 或 "LOW"
private String reason;
}
```
2. 风险评估节点:
```
package com.example.order;
import com.alibaba.cloud.ai.graph.OverAllState;
import com.alibaba.cloud.ai.graph.RunnableConfig;
import com.alibaba.cloud.ai.graph.action.NodeActionWithConfig;
import com.alibaba.cloud.ai.graph.agent.ReactAgent;
import org.springframework.ai.chat.model.ChatModel;
import java.util.HashMap;
import java.util.Map;
public class RiskAssessNode implements NodeAction {
private final ReactAgent riskAgent;
public RiskAssessNode(ChatModel chatModel) {
this.riskAgent = ReactAgent.builder()
.name("risk_assess_agent")
.model(chatModel)
.instruction("判断订单风险等级。金额>5000或包含敏感商品为高风险。" +
"返回JSON: {\"riskLevel\": \"HIGH\", \"reason\": \"金额超过5000\"}")
.outputType(RiskResult.class)
.build();
}
@Override
public Map<String, Object> apply(OverAllState state) throws Exception {
String orderInfo = state.value("order_info", "").toString();
var response = riskAgent.invoke(Map.of("input", orderInfo));
String content = response.orElseThrow().value("risk_assess_agent").orElse("").toString();
// 简单解析:如果包含HIGH则高风险
boolean isHigh = content.contains("HIGH");
Map<String, Object> output = new HashMap<>();
output.put("risk_level", isHigh ? "HIGH" : "LOW");
output.put("risk_reason", content);
output.put("flow_outcome", isHigh ? "NEED_REVIEW" : "CONTINUE");
return output;
}
}
```
3. 等待审核节点(实现 InterruptableAction):
```
package com.example.order;
import com.alibaba.cloud.ai.graph.OverAllState;
import com.alibaba.cloud.ai.graph.RunnableConfig;
import com.alibaba.cloud.ai.graph.action.AsyncNodeActionWithConfig;
import com.alibaba.cloud.ai.graph.action.InterruptableAction;
import com.alibaba.cloud.ai.graph.action.InterruptionMetadata;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.CompletableFuture;
public class WaitReviewNode implements AsyncNodeAction, InterruptableAction {
@Override
public Optional<InterruptionMetadata> interrupt(String nodeId,
OverAllState state,
RunnableConfig config) {
if (config.metadata(RunnableConfig.HUMAN_FEEDBACK_METADATA_KEY).isPresent()) {
config.metadata().ifPresent(metadata ->
metadata.remove(RunnableConfig.HUMAN_FEEDBACK_METADATA_KEY));
return Optional.empty(); // 恢复执行,不中断
}
return Optional.of(InterruptionMetadata.builder(nodeId, state).build()); // 首次,中断
}
@Override
public CompletableFuture<Map<String, Object>> apply(OverAllState state) {
// 恢复后执行:透传审核结果
String reviewResult = state.value("review_result", "approved").toString();
Map<String, Object> output = new HashMap<>();
output.put("review_result", reviewResult);
return CompletableFuture.completedFuture(output);
}
}
```
4. 订单处理节点:
```
package com.example.order;
import com.alibaba.cloud.ai.graph.OverAllState;
import com.alibaba.cloud.ai.graph.RunnableConfig;
import com.alibaba.cloud.ai.graph.action.NodeActionWithConfig;
import java.util.HashMap;
import java.util.Map;
public class ProcessOrderNode implements NodeAction {
@Override
public Map<String, Object> apply(OverAllState state) throws Exception {
String orderInfo = state.value("order_info", "").toString();
String riskLevel = state.value("risk_level", "LOW").toString();
String reviewResult = state.value("review_result", "auto_approved").toString();
Map<String, Object> output = new HashMap<>();
if (reviewResult.contains("拒绝") || reviewResult.toLowerCase().contains("reject")) {
output.put("process_status", "REJECTED");
output.put("process_result", "订单未通过人工审核,系统已取消后续履约。订单信息:" + orderInfo);
output.put("next_action", "通知用户订单审核未通过,并记录审核结果:" + reviewResult);
return output;
}
if ("HIGH".equals(riskLevel)) {
output.put("process_status", "MANUAL_APPROVED");
output.put("process_result", "高风险订单已通过人工审核,可以继续履约。订单信息:" + orderInfo);
output.put("next_action", "标记为人工审核通过,进入仓库拣货和支付确认流程");
output.put("audit_note", reviewResult);
return output;
}
output.put("process_status", "AUTO_APPROVED");
output.put("process_result", "低风险订单自动通过,系统直接进入履约流程。订单信息:" + orderInfo);
output.put("next_action", "进入仓库拣货和支付确认流程");
return output;
}
}
```
5. 图的构建与 Controller(完整可运行):
```
package com.example.order;
import com.alibaba.cloud.ai.graph.*;
import com.alibaba.cloud.ai.graph.action.AsyncEdgeAction;
import com.alibaba.cloud.ai.graph.action.AsyncNodeActionWithConfig;
import com.alibaba.cloud.ai.graph.checkpoint.config.SaverConfig;
import com.alibaba.cloud.ai.graph.checkpoint.savers.MemorySaver;
import org.springframework.ai.chat.model.ChatModel;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.web.bind.annotation.*;
import java.util.Map;
import java.util.Optional;
import java.util.UUID;
@Configuration
public class OrderWorkflowConfig {
@Bean
public CompiledGraph orderWorkflow(ChatModel chatModel) throws Exception {
StateGraph graph = new StateGraph();
graph.addNode("risk_assess", AsyncNodeActionWithConfig.node_async(new RiskAssessNode(chatModel)));
graph.addNode("wait_review", new WaitReviewNode());
graph.addNode("process_order", AsyncNodeActionWithConfig.node_async(new ProcessOrderNode()));
graph.addEdge(StateGraph.START, "risk_assess");
graph.addConditionalEdges("risk_assess",
AsyncEdgeAction.edge_async(state -> {
String outcome = state.value("flow_outcome", "CONTINUE").toString();
return "NEED_REVIEW".equals(outcome) ? "wait_review" : "process_order";
}),
Map.of("wait_review", "wait_review", "process_order", "process_order"));
graph.addEdge("wait_review", "process_order");
graph.addEdge("process_order", StateGraph.END);
return graph.compile(CompileConfig.builder()
.saverConfig(SaverConfig.builder().register(new MemorySaver()).build())
.build());
}
}
@RestController
@RequestMapping("/api/order")
public class OrderController {
private final CompiledGraph orderWorkflow;
public OrderController(CompiledGraph orderWorkflow) {
this.orderWorkflow = orderWorkflow;
}
@PostMapping("/chat")
public Map<String, Object> chat(@RequestBody Map<String, String> request) {
String input = request.getOrDefault("input", "");
String resume = request.get("resume");
boolean isResume = resume != null && !resume.isBlank();
String threadId = request.get("threadId");
if (isResume && (threadId == null || threadId.isBlank())) {
return Map.of("type", "error", "message", "恢复执行必须传入上一次中断返回的 threadId");
}
if (!isResume && (threadId == null || threadId.isBlank())) {
threadId = UUID.randomUUID().toString();
}
RunnableConfig config = isResume
? RunnableConfig.builder()
.threadId(threadId)
.resume()
.addStateUpdate(Map.of("review_result", resume))
.build()
: RunnableConfig.builder()
.threadId(threadId)
.build();
Map<String, Object> graphInput = isResume
? Map.of()
: Map.of("order_info", input);
Optional<NodeOutput> output = orderWorkflow.invokeAndGetOutput(graphInput, config);
return output.map(nodeOutput -> {
OverAllState state = nodeOutput.state();
if (!nodeOutput.isEND()) {
return Map.<String, Object>of(
"type", "interrupt",
"threadId", threadId,
"node", nodeOutput.node(),
"message", "该订单风险较高,请人工审核。回复 approved 或 rejected。",
"risk_level", state.value("risk_level", "").toString(),
"risk_reason", state.value("risk_reason", "").toString()
);
}
return Map.<String, Object>of(
"type", "normal",
"threadId", threadId,
"process_status", state.value("process_status", "").toString(),
"process_result", state.value("process_result", "").toString(),
"next_action", state.value("next_action", "").toString()
);
}).orElse(Map.of("type", "error", "message", "图执行失败"));
}
}
```
这个 Controller 只有一个 /api/order/chat 方法。前端根据上一次响应的 type 决定下一次请求字段:
- 如果后端返回
type=normal,说明图已经正常执行完成,前端展示process_result即可 - 如果后端返回
type=interrupt,说明图在wait_review节点中断了,前端需要展示审核提示,并保留后端返回的threadId - 用户第一次提交订单时,把订单内容放在
input字段中 - 如果上一次响应是
type=interrupt,用户再次发送审核意见时,把审核意见放在resume字段中,并必须带上同一个threadId
Controller 的判断逻辑也很直接:
- 请求中有
resume:说明这是恢复执行,必须使用上一次中断返回的threadId,并构造RunnableConfig.builder().threadId(threadId).resume().addStateUpdate(...) - 请求中没有
resume:说明这是首次执行,如果前端没有传threadId,后端就生成一个新的threadId,并把input写入图输入的order_info - 调用统一使用
invokeAndGetOutput,然后通过nodeOutput.isEND()判断是否中断 - 如果
isEND()为false,返回type=interrupt - 如果
isEND()为true,返回type=normal
两条执行路径总结:
路径一:低风险订单(不中断)
```
前端请求:{ "input": "订单金额300,商品:普通鼠标" }
→ Controller 判断没有 resume,正常启动图
→ risk_assess 判定 LOW
→ 条件边直接到 process_order
→ ProcessOrderNode 返回 AUTO_APPROVED、process_result、next_action
→ Controller 返回 type=normal
```
路径二:高风险订单(中断 + 恢复后跳转到其他节点)
```
前端请求:{ "input": "订单金额8000,商品:敏感设备" }
→ Controller 判断没有 resume,正常启动图
→ risk_assess 判定 HIGH
→ 条件边到 wait_review
→ wait_review.interrupt() 返回 metadata,图中断
→ Controller 发现 nodeOutput.isEND() 为 false
→ 返回 type=interrupt、threadId、risk_reason 给前端
前端再次请求:{ "threadId": "上次返回的threadId", "resume": "approved" }
→ Controller 判断有 resume,构造 resumeConfig
→ 框架恢复到 wait_review,并把 review_result 合并到 state
→ wait_review.interrupt() 删除恢复标记并返回 empty
→ wait_review.apply() 透传审核结果
→ 通过边流转到 process_order
→ ProcessOrderNode 根据 review_result 返回 MANUAL_APPROVED 或 REJECTED
→ Controller 返回 type=normal
```