
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_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content控制台创建 Keyhttps://taotoken.net/console?utm_sourcetaotoken_aicg_blog_endutm_contentconsoleutm_campaignrewriteAPI Keys 管理https://taotoken.net/api-keys?utm_sourcetaotoken_aicg_blog_endutm_contentapi-keysutm_campaignrewrite创建完 Key 之后建议先在模型对话页面做一次连通性测试确认 Key 和模型名可用再去配 Spring AI。模型对话入口https://taotoken.net/model-chat?utm_sourcetaotoken_aicg_blog_endutm_contentmodel-chatutm_campaignrewriteAPI 的基础地址是https://taotoken.net/api注意这个地址不带 UTM 参数配置里直接写它就行。Key 的格式通常是sk-开头复制后先放到环境变量里别硬编码进代码。注意API Key 只放在本地环境变量或配置中心不要提交到 Git 仓库。联调阶段可以先用测试 Key跑通后再换成正式 Key。如果你后面要做长期的编码或 Agent 场景可以了解下 Coding Plan它更适合持续性的开发任务https://taotoken.net/coding-plan?utm_sourcetaotoken_aicg_blog_endutm_contentcoding-planutm_campaignrewrite3. 可复制配置MCP 服务端与 Spring AI 客户端骨架这一节分两部分先写 MCP 服务端把 RocketMQ 查询封装成工具再写 Spring AI 客户端把模型和 MCP 工具接起来。3.1 MCP 服务端依赖与工具注册服务端用 Spring Boot 3.4.x Spring AI MCP Server StarterJDK 21。核心依赖如下properties maven.compiler.source21/maven.compiler.source maven.compiler.target21/maven.compiler.target project.build.sourceEncodingUTF-8/project.build.sourceEncoding spring-boot.version3.4.1/spring-boot.version spring-ai.version1.0.0-M6/spring-ai.version rocketmq.version5.1.0/rocketmq.version /properties dependencies dependency groupIdorg.springframework.ai/groupId artifactIdspring-ai-mcp-server-webmvc-spring-boot-starter/artifactId version${spring-ai.version}/version /dependency dependency groupIdorg.apache.rocketmq/groupId artifactIdrocketmq-client/artifactId version${rocketmq.version}/version /dependency dependency groupIdorg.apache.rocketmq/groupId artifactIdrocketmq-tools/artifactId version${rocketmq.version}/version /dependency dependency groupIdorg.projectlombok/groupId artifactIdlombok/artifactId version1.18.30/version /dependency /dependenciesrocketmq-tools里包含DefaultMQAdminExt这是查询消息用的管理端扩展类。spring-ai-mcp-server-webmvc-spring-boot-starter负责把带Tool注解的方法暴露成 MCP 工具。工具注册用一个配置类完成把MessageService里的方法注册进ToolCallbackProviderConfiguration 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 MapString, 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 Starterdependency groupIdorg.springframework.ai/groupId artifactIdspring-ai-openai-spring-boot-starter/artifactId version1.0.0-M6/version /dependency dependency groupIdorg.springframework.ai/groupId artifactIdspring-ai-mcp-client-spring-boot-starter/artifactId version1.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:8081SSE 路径/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:9876topic 为 xiaozou-batch-topicmessageId 为 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 开了 ACLaccessKey和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_sourcetaotoken_aicg_blog_endutm_contentmodel-chatutm_campaignrewrite接入相关的文档和 Key 管理入口接入文档https://taotoken.net/doc?utm_sourcetaotoken_aicg_blog_endutm_contentdocutm_campaignrewriteAPI Keyshttps://taotoken.net/api-keys?utm_sourcetaotoken_aicg_blog_endutm_contentapi-keysutm_campaignrewrite6. 把查询链路接进日常联调跑通之后这套东西最实用的地方是把它接进你日常用的 AI 客户端。比如在支持 MCP 的编辑器或 Agent 工具里配置远程 MCP 服务端之后查消息就不用再切控制台了。配置方式和 3.4 里的客户端类似填服务端地址和 SSE 路径即可。如果你打算长期用这套链路做编码和联调Coding Plan 会更合适它面向持续性的开发任务https://taotoken.net/coding-plan?utm_sourcetaotoken_aicg_blog_endutm_contentcoding-planutm_campaignrewrite实际用下来有几个经验值得记一下。工具方法尽量保持单一职责一个方法只做一件事模型选起来不容易错。参数名用英文描述用中文模型对参数名的匹配更准。查询类工具返回结构化的 DTO比返回裸字符串更好模型能更稳定地转述。另外MCP 服务端不要暴露删除、重置这类写操作联调环境也尽量只读避免误操作。这套骨架目前只实现了按 MessageId 查询扩展方向很直接加一个按 Topic 和时间范围拉取消息列表的工具再加一个按 Key 查询的工具基本就能覆盖日常排查的大部分场景。每个工具就是一个带Tool注解的方法注册逻辑不用改。