1. 为什么要在本地联调里用自然语言查 RocketMQ 消息
RocketMQ 的消息查询一直是联调阶段最费时间的环节之一。平时排查一条消息有没有发出去、内容对不对,通常要在控制台里翻 Topic、填 MessageId、选时间范围,或者干脆写一段DefaultMQAdminExt的临时测试代码。消息量一多,控制台翻页慢,脚本又要反复改参数,效率很低。
Spring AI 的 Tool Calling 加上 MCP 协议,正好能把这件事变简单。MCP 可以理解成 AI 应用和外部工具之间的标准接口,它规定了工具怎么描述、参数怎么传、结果怎么回。Spring AI 负责把 Java 方法注册成模型可调用的工具,模型根据你的自然语言描述去决定调用哪个方法、传什么参数。两者结合后,你只要说一句“帮我查一下 topic 是 xxx、messageId 是 xxx 的消息”,链路就会自动走到 RocketMQ 的查询接口上。
这套方案适合谁?适合正在做 RocketMQ 本地开发、联调、测试的同学,也适合想把日常运维动作接进 AI 客户端的团队。本文聚焦本地开发与联调场景,给出可复制的 MCP 服务端骨架、Spring AI 客户端配置,以及一次完整的自然语言查询验证。跑通之后你会清楚每个配置项在干什么,而不是只复制一堆代码。
需要提前说明的是,本文的 MCP 服务端只做消息查询这类只读操作,不涉及生产库直连,也不建议把管理类写操作直接暴露给模型。联调环境用本地或测试 Nameserver 即可。
2. TaoToken 前置:给 Spring AI 客户端准备模型入口
Spring AI 客户端要调用模型,需要一个兼容 OpenAI 接口的模型服务地址和 API Key。我这边用的是 TaoToken,它的接口路径和 OpenAI 风格一致,Spring AI 的OpenAiChatModel可以直接对接,不用改太多配置。
先到官网注册并进入控制台,在 API Keys 页面创建一个 Key。地址如下:
- 官网入口:https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content=
- 控制台(创建 Key):https://taotoken.net/console?utm_source=taotoken_aicg_blog_end&utm_content=console&utm_campaign=rewrite
- API Keys 管理:https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_content=api-keys&utm_campaign=rewrite
创建完 Key 之后,建议先在模型对话页面做一次连通性测试,确认 Key 和模型名可用,再去配 Spring AI。模型对话入口:
https://taotoken.net/model-chat?utm_source=taotoken_aicg_blog_end&utm_content=model-chat&utm_campaign=rewrite
API 的基础地址是https://taotoken.net/api,注意这个地址不带 UTM 参数,配置里直接写它就行。Key 的格式通常是sk-开头,复制后先放到环境变量里,别硬编码进代码。
注意:API Key 只放在本地环境变量或配置中心,不要提交到 Git 仓库。联调阶段可以先用测试 Key,跑通后再换成正式 Key。
如果你后面要做长期的编码或 Agent 场景,可以了解下 Coding Plan,它更适合持续性的开发任务:
https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=coding-plan&utm_campaign=rewrite
3. 可复制配置:MCP 服务端与 Spring AI 客户端骨架
这一节分两部分:先写 MCP 服务端,把 RocketMQ 查询封装成工具;再写 Spring AI 客户端,把模型和 MCP 工具接起来。
3.1 MCP 服务端依赖与工具注册
服务端用 Spring Boot 3.4.x + Spring AI MCP Server Starter,JDK 21。核心依赖如下:
<properties> <maven.compiler.source>21</maven.compiler.source> <maven.compiler.target>21</maven.compiler.target> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> <spring-boot.version>3.4.1</spring-boot.version> <spring-ai.version>1.0.0-M6</spring-ai.version> <rocketmq.version>5.1.0</rocketmq.version> </properties> <dependencies> <dependency> <groupId>org.springframework.ai</groupId> <artifactId>spring-ai-mcp-server-webmvc-spring-boot-starter</artifactId> <version>${spring-ai.version}</version> </dependency> <dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-client</artifactId> <version>${rocketmq.version}</version> </dependency> <dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-tools</artifactId> <version>${rocketmq.version}</version> </dependency> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <version>1.18.30</version> </dependency> </dependencies>rocketmq-tools里包含DefaultMQAdminExt,这是查询消息用的管理端扩展类。spring-ai-mcp-server-webmvc-spring-boot-starter负责把带@Tool注解的方法暴露成 MCP 工具。
工具注册用一个配置类完成,把MessageService里的方法注册进ToolCallbackProvider:
@Configuration public class MCPAutoConfiguration { @Bean public ToolCallbackProvider rocketmqTools(MessageService messageService) { return MethodToolCallbackProvider.builder() .toolObjects(messageService) .build(); } }MethodToolCallbackProvider会扫描messageService里所有带@Tool注解的方法,把方法名、描述、参数类型转成 MCP 工具定义。模型看到的就是这些描述,所以描述写得越清楚,模型选工具越准。
3.2 核心查询工具方法
MessageService里定义查询方法,用@Tool标注。这里用DefaultMQAdminExt按 MessageId 查询,并缓存每个 Nameserver 对应的管理端实例,避免每次查询都重新启动客户端:
@Service @RequiredArgsConstructor public class MessageServiceImpl implements MessageService { private final Map<String, DefaultMQAdminExt> adminExtCache = new ConcurrentHashMap<>(); private static final int QUERY_MESSAGE_MAX_NUM = 64; @Tool(description = "通过 nameserver、topic 和 messageId 查询 RocketMQ 消息内容", name = "查询消息") @Override public MessageView queryMessageById(String nameserver, String topic, String messageId, String accessKey, String secretKey) { DefaultMQAdminExt adminExt = adminExtCache.computeIfAbsent(nameserver, ns -> { DefaultMQAdminExt ext = (accessKey != null && secretKey != null) ? new DefaultMQAdminExt(new AclClientRPCHook( new SessionCredentials(accessKey, secretKey))) : new DefaultMQAdminExt(); ext.setNamesrvAddr(ns); ext.setInstanceName(Long.toString(System.currentTimeMillis())); try { ext.start(); } catch (MQClientException e) { throw new RuntimeException("启动 MQAdminExt 失败: " + e.getMessage(), e); } return ext; }); try { long begin = MessageClientIDSetter.getNearlyTimeFromID(messageId).getTime() - 1000 * 60 * 60 * 13L; QueryResult result = adminExt.queryMessageByUniqKey( topic, messageId, QUERY_MESSAGE_MAX_NUM, begin, Long.MAX_VALUE); if (result.getMessageList().isEmpty()) { return null; } MessageExt ext = result.getMessageList().get(0); return new MessageView( ext.getTopic(), messageId, new String(ext.getBody(), StandardCharsets.UTF_8), JSON.toJSONString(ext.getProperties())); } catch (MQClientException | InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("查询消息失败: " + e.getMessage(), e); } } }几个关键点。queryMessageByUniqKey的第三个参数是单次最大返回条数,这里设 64,避免一次拉太多。第四个参数是起始时间,MessageId 里本身带时间戳,用getNearlyTimeFromID反推出来再往前推 13 小时,是为了覆盖时区差异和消息堆积的情况。MessageView是个简单 DTO,包含 topic、messageId、body 和 properties。
3.3 MCP 服务端 application.yml
服务端配置主要是端口和 MCP 的传输方式:
server: port: 8081 spring: application: name: rocketmq-mcp-server ai: mcp: server: name: rocketmq-mcp version: 1.0.0 type: SYNC sse-endpoint: /sse sse-message-endpoint: /mcp/messagetype: SYNC表示同步工具调用,适合查询这类短耗时操作。sse-endpoint是客户端建立 SSE 长连接的路径,sse-message-endpoint是客户端发送工具调用请求的路径。这两个路径客户端要和服务端保持一致。
3.4 Spring AI 客户端配置
客户端是另一个 Spring Boot 应用,引入 Spring AI 的 OpenAI Starter 和 MCP Client Starter:
<dependency> <groupId>org.springframework.ai</groupId> <artifactId>spring-ai-openai-spring-boot-starter</artifactId> <version>1.0.0-M6</version> </dependency> <dependency> <groupId>org.springframework.ai</groupId> <artifactId>spring-ai-mcp-client-spring-boot-starter</artifactId> <version>1.0.0-M6</version> </dependency>客户端的application.yml要同时配模型和 MCP 服务端地址:
server: port: 8080 spring: ai: openai: base-url: https://taotoken.net/api api-key: ${TAOTOKEN_API_KEY} chat: options: model: gpt-4o-mini temperature: 0.2 mcp: client: enabled: true name: rocketmq-client version: 1.0.0 type: SYNC sse: connections: rocketmq-server: url: http://127.0.0.1:8081 sse-endpoint: /ssebase-url指向 TaoToken 的 API 地址,api-key从环境变量读取。temperature设低一点,让模型在选工具和填参数时更稳定。mcp.client.sse.connections下面配的是 MCP 服务端的连接信息,url是服务端地址,sse-endpoint要和 3.3 里的一致。
3.5 客户端调用入口
客户端里注入ChatClient,把 MCP 工具挂上去,然后接收自然语言输入:
@RestController @RequiredArgsConstructor public class ChatController { private final ChatClient.Builder chatClientBuilder; private final ToolCallbackProvider mcpToolCallbackProvider; @PostMapping("/chat") public String chat(@RequestBody String userInput) { ChatClient chatClient = chatClientBuilder .defaultToolCallbacks(mcpToolCallbackProvider) .build(); return chatClient.prompt() .user(userInput) .call() .content(); } }defaultToolCallbacks把 MCP 客户端拉取到的工具注册进对话上下文。模型在生成回复前会先判断是否需要调用工具,需要的话就按工具定义生成参数,MCP 客户端再把调用转发给服务端。
4. 验证请求:一次完整的自然语言查询
配置写完,先启动 MCP 服务端,再启动客户端。服务端启动日志里应该能看到 MCP Server 注册的工具列表,包含“查询消息”这一项。
4.1 用 MCP Inspector 先验证服务端
在正式接模型之前,建议先用 MCP Inspector 单独验证服务端工具是否可用。安装并启动:
npx @modelcontextprotocol/inspector浏览器打开http://127.0.0.1:6274,在连接配置里填服务端地址http://127.0.0.1:8081,SSE 路径/sse。连接成功后左侧会列出工具,点“查询消息”,填入参数:
{ "nameserver": "127.0.0.1:9876", "topic": "xiaozou-batch-topic", "messageId": "AC1400010D3E068DE1451D546AEE0173", "accessKey": null, "secretKey": null }如果返回了消息体,说明服务端到 RocketMQ 的链路是通的。这一步能排除掉 RocketMQ 连接、ACL、MessageId 格式等问题,把问题范围缩小到模型侧。
4.2 用自然语言发起查询
服务端验证通过后,向客户端发请求:
curl -X POST http://127.0.0.1:8080/chat \ -H "Content-Type: text/plain" \ -d "帮我查询 nameserver 地址为 127.0.0.1:9876,topic 为 xiaozou-batch-topic,messageId 为 AC1400010D3E068DE1451D546AEE0173 的消息,accessKey 和 secretKey 为空"模型会解析这句话,识别出要调用“查询消息”工具,并从自然语言里抽取 nameserver、topic、messageId 三个参数,accessKey 和 secretKey 按“为空”处理成 null。调用结果返回后,模型会把消息内容组织成一段可读的回复。
4.3 成功结果长什么样
正常情况下,你会看到类似这样的返回:
{ "topic": "xiaozou-batch-topic", "messageId": "AC1400010D3E068DE1451D546AEE0173", "body": "{\"orderId\":\"20250101001\",\"status\":\"PAID\"}", "properties": "{\"KEYS\":\"20250101001\",\"TAGS\":\"pay\"}" }模型侧可能会把它转述成“这条消息的 topic 是 xiaozou-batch-topic,消息体里 orderId 是 20250101001,状态是 PAID”。到这一步,整条链路就跑通了:自然语言 → 模型解析 → MCP 工具调用 → RocketMQ 查询 → 结果回传。
5. 本篇常见错排查
跑不通的时候,按下面几个方向排查,基本能覆盖大部分问题。
工具列表为空。客户端启动后如果模型一直说“没有可用工具”,先看 MCP 客户端有没有成功连上服务端。检查spring.ai.mcp.client.sse.connections里的url和sse-endpoint是否和服务端一致。服务端如果是SYNC类型,客户端也要配SYNC,类型不匹配会导致握手失败。
模型不调用工具,直接编答案。这种情况通常是工具描述不够清楚,或者temperature太高。把@Tool的description写具体,比如“通过 nameserver、topic 和 messageId 查询 RocketMQ 消息内容”,而不是只写“查询消息”。temperature降到 0.2 以下,模型会更倾向于按工具定义走。
查询返回 null。先确认 MessageId 是否正确,RocketMQ 的 MessageId 区分大小写。再确认时间范围,如果消息是很久以前发的,getNearlyTimeFromID反推的起始时间可能不够早,可以适当把往前推的时长调大。另外确认 topic 和 nameserver 是否匹配,跨集群查是查不到的。
ACL 报错。如果 RocketMQ 开了 ACL,accessKey和secretKey必须传,且要有对应 topic 的查询权限。联调环境如果没开 ACL,传 null 即可,但要注意DefaultMQAdminExt的构造方式会因此不同,代码里已经按 null 做了分支。
SSE 连接超时。检查服务端端口是否被占用,防火墙是否放行。本地联调一般用127.0.0.1,如果服务端和客户端不在同一台机器,要把127.0.0.1换成实际 IP,并确认 MCP 服务端的sse-endpoint路径没有被网关改写。
模型返回乱码或截断。消息体如果是二进制或压缩内容,直接按 UTF-8 解码会出问题。联调阶段建议只查文本消息,二进制消息先跳过。另外queryMessageByUniqKey的返回条数上限设太大可能拖慢响应,64 是个比较稳的值。
如果排查过程中需要确认模型侧是否正常,可以到模型对话页面单独测一下模型连通性:
https://taotoken.net/model-chat?utm_source=taotoken_aicg_blog_end&utm_content=model-chat&utm_campaign=rewrite
接入相关的文档和 Key 管理入口:
- 接入文档:https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc&utm_campaign=rewrite
- API Keys:https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_content=api-keys&utm_campaign=rewrite
6. 把查询链路接进日常联调
跑通之后,这套东西最实用的地方是把它接进你日常用的 AI 客户端。比如在支持 MCP 的编辑器或 Agent 工具里配置远程 MCP 服务端,之后查消息就不用再切控制台了。配置方式和 3.4 里的客户端类似,填服务端地址和 SSE 路径即可。
如果你打算长期用这套链路做编码和联调,Coding Plan 会更合适,它面向持续性的开发任务:
https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=coding-plan&utm_campaign=rewrite
实际用下来,有几个经验值得记一下。工具方法尽量保持单一职责,一个方法只做一件事,模型选起来不容易错。参数名用英文,描述用中文,模型对参数名的匹配更准。查询类工具返回结构化的 DTO,比返回裸字符串更好,模型能更稳定地转述。另外,MCP 服务端不要暴露删除、重置这类写操作,联调环境也尽量只读,避免误操作。
这套骨架目前只实现了按 MessageId 查询,扩展方向很直接:加一个按 Topic 和时间范围拉取消息列表的工具,再加一个按 Key 查询的工具,基本就能覆盖日常排查的大部分场景。每个工具就是一个带@Tool注解的方法,注册逻辑不用改。