首屏导读 · 本教程配套付费专栏: 大模型工程师修炼手记
19.9 元(AI 编程 / Agent 实战 | 本文同主题系统课程)· AI时代程序员的自我提升49.9 元(AI 时代成长方法论)。单篇不过瘾?订阅解锁全量源码、实战与答疑;文末附资料包领取方式 ↓
引言
上一篇文章中,我们实现了基于stdio的MCP天气查询服务器,让本地运行的Claude Desktop能够获取模拟天气数据。但stdio模式有一个天然局限:服务器必须和客户端运行在同一台机器上,由Claude直接拉起进程。如果你的AI应用部署在云端,或者需要为多个远程客户端提供服务,stdio就鞭长莫及了。这时候,基于HTTP+SSE(Server-Sent Events)的MCP服务器才是王道——它让AI模型可以通过网络实时调用你的Java服务,无论服务部署在云上、本地还是内网。本文将带你用Spring Boot和WebFlux构建一个支持SSE的天气查询MCP服务器,让Claude变身真正的云端气象专家!
为什么选择SSE?
SSE(Server-Sent Events)是一种基于HTTP的轻量级推送技术,允许服务器主动向客户端发送数据。在MCP协议中,SSE传输模式具有以下优势:
- 远程访问:基于HTTP,天然支持跨网络调用,适合微服务架构。
- 实时推送:服务器可以随时将工具调用结果推送给客户端,无需客户端轮询。
- 自动重连:浏览器和HTTP客户端通常内置SSE重连机制,增强稳定性。
- 简单轻量:相比WebSocket,SSE只需要服务端实现单向推送,客户端通过EventSource即可接收,实现成本更低。
MCP的SSE握手流程如下:
- 客户端连接
/sse端点,建立SSE长连接。 - 服务器为该连接生成唯一
sessionId,并通过endpoint事件返回一个消息接收URL(如/messages?sessionId=xxx)。 - 客户端通过该URL发送JSON-RPC请求(POST)。
- 服务器处理请求后,通过最初的SSE连接将响应推回客户端。
项目准备
技术栈
- Java 17
- Spring Boot 3.2 + WebFlux
- Lombok
- Jackson
创建项目
可以使用Spring Initializr(https://start.spring.io/)生成基础项目,选择依赖:
- Spring WebFlux
- Lombok
Maven依赖(pom.xml)
<?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.example</groupId> <artifactId>weather-mcp-sse</artifactId> <version>1.0.0</version> <packaging>jar</packaging> <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>3.2.3</version> </parent> <properties> <java.version>17</java.version> </properties> <dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-webflux</artifactId> </dependency> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <optional>true</optional> </dependency> <!-- 如果需要调用真实天气API,可以添加HttpClient依赖 --> <!-- <dependency> <groupId>org.apache.httpcomponents.client5</groupId> <artifactId>httpclient5</artifactId> </dependency> --> </dependencies> <build> <plugins> <plugin> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-maven-plugin</artifactId> <configuration> <excludes> <exclude> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> </exclude> </excludes> </configuration> </plugin> </plugins> </build> </project>项目结构
weather-mcp-sse/ ├── src/main/java/com/example/weather/ │ ├── WeatherMcpSseApplication.java │ ├── controller/ │ │ └── McpController.java │ ├── handler/ │ │ └── McpMessageHandler.java │ ├── model/ │ │ ├── JsonRpcMessage.java │ │ └── JsonRpcError.java │ └── service/ │ └── SessionManager.java └── src/main/resources/ └── application.yml代码实现
1. 主应用类
package com.example.weather; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; @SpringBootApplication public class WeatherMcpSseApplication { public static void main(String[] args) { SpringApplication.run(WeatherMcpSseApplication.class, args); } }2. 配置类(可选,这里不额外配置,直接使用默认)
3. 模型类
JsonRpcMessage.java:通用的JSON-RPC消息封装。
package com.example.weather.model; import com.fasterxml.jackson.annotation.JsonInclude; import com.fasterxml.jackson.databind.JsonNode; import lombok.AllArgsConstructor; import lombok.Builder; import lombok.Data; import lombok.NoArgsConstructor; @Data @Builder @NoArgsConstructor @AllArgsConstructor @JsonInclude(JsonInclude.Include.NON_NULL) public class JsonRpcMessage { private String jsonrpc; private JsonNode id; private String method; private JsonNode params; private JsonNode result; private JsonRpcError error; }JsonRpcError.java:错误对象。
package com.example.weather.model; import com.fasterxml.jackson.databind.JsonNode; import lombok.AllArgsConstructor; import lombok.Builder; import lombok.Data; import lombok.NoArgsConstructor; @Data @Builder @NoArgsConstructor @AllArgsConstructor public class JsonRpcError { private int code; private String message; private JsonNode data; }4. 会话管理器
管理每个SSE连接的会话,维护会话ID到消息发射器的映射,用于向特定客户端推送消息。
package com.example.weather.service; import org.springframework.stereotype.Component; import reactor.core.publisher.FluxSink; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; @Component public class SessionManager { private final Map<String, FluxSink<String>> sessions = new ConcurrentHashMap<>(); public void register(String sessionId, FluxSink<String> sink) { sessions.put(sessionId, sink); } public void remove(String sessionId) { sessions.remove(sessionId); } public void sendMessage(String sessionId, String message) { FluxSink<String> sink = sessions.get(sessionId); if (sink != null) { sink.next("event: message\n"); sink.next("data: " + message + "\n\n"); } } public void sendEndpoint(String sessionId, String endpoint) { FluxSink<String> sink = sessions.get(sessionId); if (sink != null) { sink.next("event: endpoint\n"); sink.next("data: " + endpoint + "\n\n"); } } }5. MCP消息处理器
这是核心业务类,负责处理SSE连接的建立、JSON-RPC消息的分发、工具定义和执行。
package com.example.weather.handler; import com.example.weather.model.JsonRpcMessage; import com.example.weather.model.JsonRpcError; import com.example.weather.service.SessionManager; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import reactor.core.publisher.Flux; import reactor.core.publisher.FluxSink; import java.util.Random; import java.util.UUID; @Slf4j @Component @RequiredArgsConstructor public class McpMessageHandler { private final ObjectMapper objectMapper; private final SessionManager sessionManager; private final Random random = new Random(); /** * 创建SSE连接流 */ public Flux<String> createSseConnection() { String sessionId = UUID.randomUUID().toString(); return Flux.create(sink -> { sessionManager.register(sessionId, sink); // 立即发送endpoint事件,告知客户端后续消息发送地址 String endpoint = "/messages?sessionId=" + sessionId; sessionManager.sendEndpoint(sessionId, endpoint); log.info("SSE connected: {}", sessionId); // 连接关闭时清理会话 sink.onCancel(() -> { sessionManager.remove(sessionId); log.info("SSE disconnected: {}", sessionId); }); }, FluxSink.OverflowStrategy.BUFFER); } /** * 处理客户端通过/messages发送的JSON-RPC请求 */ public void handleMessage(String sessionId, String rawMessage) { log.info("Received from session {}: {}", sessionId, rawMessage); try { JsonNode request = objectMapper.readTree(rawMessage); JsonRpcMessage response = processRequest(request); if (response != null) { String responseStr = objectMapper.writeValueAsString(response); sessionManager.sendMessage(sessionId, responseStr); } } catch (Exception e) { log.error("Error handling message", e); } } /** * 处理具体的JSON-RPC请求 */ private JsonRpcMessage processRequest(JsonNode request) { String method = request.has("method") ? request.get("method").asText() : null; JsonNode params = request.get("params"); JsonNode id = request.get("id"); // 忽略通知(没有id的请求) if (id == null || id.isNull()) { return null; } try { JsonNode result = switch (method) { case "initialize" -> handleInitialize(params); case "tools/list" -> handleToolsList(); case "tools/call" -> handleToolsCall(params); default -> throw new IllegalArgumentException("Unknown method: " + method); }; return JsonRpcMessage.builder() .jsonrpc("2.0") .id(id) .result(result) .build(); } catch (Exception e) { return JsonRpcMessage.builder() .jsonrpc("2.0") .id(id) .error(JsonRpcError.builder() .code(-32603) .message("Internal error: " + e.getMessage()) .build()) .build(); } } /** * 初始化:返回协议版本、服务器能力、服务器信息 */ private JsonNode handleInitialize(JsonNode params) { ObjectNode result = objectMapper.createObjectNode(); result.put("protocolVersion", "0.1.0"); ObjectNode capabilities = objectMapper.createObjectNode(); capabilities.put("tools", true); result.set("capabilities", capabilities); ObjectNode serverInfo = objectMapper.createObjectNode(); serverInfo.put("name", "java-weather-mcp-sse"); serverInfo.put("version", "1.0.0"); result.set("serverInfo", serverInfo); return result; } /** * 工具列表:返回天气查询工具的定义 */ private JsonNode handleToolsList() { ObjectNode result = objectMapper.createObjectNode(); var toolsArray = objectMapper.createArrayNode(); ObjectNode tool = objectMapper.createObjectNode(); tool.put("name", "get_weather"); tool.put("description", "获取指定城市的实时天气信息"); // 定义输入参数:城市名称 ObjectNode inputSchema = objectMapper.createObjectNode(); inputSchema.put("type", "object"); ObjectNode properties = objectMapper.createObjectNode(); ObjectNode citySchema = objectMapper.createObjectNode(); citySchema.put("type", "string"); citySchema.put("description", "城市名称,例如:北京、上海、纽约"); properties.set("city", citySchema); inputSchema.set("properties", properties); inputSchema.put("required", objectMapper.createArrayNode().add("city")); tool.set("inputSchema", inputSchema); toolsArray.add(tool); result.set("tools", toolsArray); return result; } /** * 调用工具:执行天气查询 */ private JsonNode handleToolsCall(JsonNode params) { String name = params.get("name").asText(); JsonNode arguments = params.get("arguments"); if ("get_weather".equals(name)) { String city = arguments.get("city").asText(); String weatherInfo = getWeatherByCity(city); ObjectNode result = objectMapper.createObjectNode(); ObjectNode content = objectMapper.createObjectNode(); content.put("type", "text"); content.put("text", weatherInfo); result.set("content", objectMapper.createArrayNode().add(content)); return result; } else { throw new IllegalArgumentException("Unknown tool: " + name); } } /** * 模拟天气查询(可替换为真实API调用) */ private String getWeatherByCity(String city) { // 这里模拟返回天气数据,实际使用中可替换为HTTP请求调用天气API String[] conditions = {"晴朗", "多云", "小雨", "大雨", "阴天", "雾霾", "雷阵雨"}; String condition = conditions[random.nextInt(conditions.length)]; int temperature = 5 + random.nextInt(30); // 5~34度 int humidity = 40 + random.nextInt(50); // 40~89% return String.format("%s天气:%s,温度%d℃,湿度%d%%", city, condition, temperature, humidity); } // 如果需要接入真实天气API,可以参考以下代码(需添加HTTP客户端依赖) /* private String getWeatherByCity(String city) { String apiKey = "your_api_key"; String url = String.format( "https://api.openweathermap.org/data/2.5/weather?q=%s&appid=%s&units=metric&lang=zh_cn", city, apiKey ); // 使用WebClient或HttpClient发送请求并解析响应 // 返回格式化后的天气信息 } */ }6. 控制器
暴露两个端点:/sse用于建立SSE连接,/messages用于接收JSON-RPC请求。
package com.example.weather.controller; import com.example.weather.handler.McpMessageHandler; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.http.MediaType; import org.springframework.http.codec.ServerSentEvent; import org.springframework.web.bind.annotation.*; import reactor.core.publisher.Flux; @Slf4j @RestController @RequiredArgsConstructor public class McpController { private final McpMessageHandler messageHandler; @GetMapping(value = "/sse", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<ServerSentEvent<String>> sse() { return messageHandler.createSseConnection() .map(data -> ServerSentEvent.<String>builder() .data(data) .build()); } @PostMapping("/messages") public void handleMessage(@RequestParam("sessionId") String sessionId, @RequestBody String message) { messageHandler.handleMessage(sessionId, message); } }7. 配置文件 application.yml
spring: application: name: weather-mcp-sse server: port: 8080 address: 0.0.0.0 # 允许外部访问 logging: level: com.example.weather: DEBUG file: name: logs/weather-mcp.log max-size: 10MB max-history: 7打包与运行
打包为可执行JAR
mvn clean package在target目录下生成weather-mcp-sse-1.0.0.jar。
运行
java -jar target/weather-mcp-sse-1.0.0.jar如果希望后台运行(Linux/Mac):
nohup java -jar target/weather-mcp-sse-1.0.0.jar > weather.log 2>&1 &服务启动后,访问http://localhost:8080/sse即可建立SSE连接(可用curl测试)。
集成到Claude Desktop
Claude Desktop支持通过HTTP方式连接远程MCP服务器,无需本地拉起进程。只需在配置文件中指定SSE端点URL即可。
配置文件位置
- macOS:
~/Library/Application Support/Claude/claude_desktop_config.json - Windows:
%APPDATA%\Claude\claude_desktop_config.json
添加MCP服务器配置
{ "mcpServers": { "java-weather-sse": { "url": "http://your-server-ip:8080/sse" } } }如果服务器运行在本地,可以使用http://localhost:8080/sse。
注意:Claude Desktop需要能够访问该URL。如果服务器部署在远程云主机,请确保防火墙开放8080端口,并且Claude所在网络可以访问。
重启Claude Desktop
完全退出Claude(包括系统托盘图标),然后重新启动。此时Claude会自动连接SSE端点,获取工具列表。在对话界面中应该能看到“锤子”图标,表明天气工具可用。
测试效果
在Claude中输入:“北京今天天气怎么样?”或者“用get_weather工具查询上海的天气”。Claude会调用远程SSE服务器,返回模拟的天气信息。
示例对话:
用户:北京今天天气怎么样?
Claude(调用工具后):根据查询,北京天气:多云,温度21℃,湿度67%。
如果接入真实天气API,返回的数据会更加准确和详细。
进阶扩展
- 接入真实天气API:替换
getWeatherByCity方法,使用Spring WebClient调用OpenWeatherMap或和风天气API,并解析返回的JSON数据。 - 增加更多工具:例如
get_forecast获取未来一周天气预报,get_air_quality查询空气质量。 - 安全认证:为SSE连接添加token验证,防止未授权访问。
- 负载均衡与集群:当多实例部署时,需要使用共享存储(如Redis)管理会话信息,确保同一客户端的请求能路由到正确的实例。
- 监控与健康检查:添加Spring Boot Actuator,提供
/actuator/health端点,方便监控服务状态。
总结
通过Spring Boot和WebFlux,我们轻松构建了一个基于SSE的MCP天气查询服务器,让Claude Desktop可以通过HTTP远程调用Java服务,实现了真正的云端集成。相比stdio模式,SSE版本具有更好的扩展性和灵活性,适合微服务架构和分布式部署。你可以在此基础上继续封装更多的业务工具,让AI模型成为你系统的“超级入口”。赶快试试吧,让你的Claude不仅能聊天,还能实时感知天气变化!
参考资源
- MCP协议官方文档:https://modelcontextprotocol.io
- Spring WebFlux文档:https://docs.spring.io/spring-framework/reference/web/webflux.html
- 项目源码(示例):https://github.com/your-repo/weather-mcp-sse
内容来源:csdn.net
作者昵称:Tom·Ge
原文链接:https://blog.csdn.net/gedonshen/article/details/158317835
作者主页:https://blog.csdn.net/gedonshen
延伸阅读 · 我的付费专栏
觉得这篇文章对你有帮助?我把同类主题的系统化内容沉淀成了付费专栏,欢迎订阅支持持续输出:
| 专栏 | 定价 | 内容 |
|---|---|---|
| 大模型工程师修炼手记 | 19.9 元 | AI 编程 / Agent 深度实战 |
| AI时代程序员的自我提升 | 49.9 元 | AI 时代成长方法论 |
本文配套代码 / 资料包:欢迎在评论区留言「求代码」,我会私信发送完整资源!