1. 为什么 Java 接入远程 MCP 服务时,SSE 传输总让人卡壳
如果你之前只用过 Stdio 方式跑 MCP,客户端启动服务端进程,通过管道一问一答,逻辑非常直白。可一旦服务端跑在远端机器上,你没法启动它的进程,只能走 HTTP 通信,问题就来了:HTTP 是请求-响应模型,服务端没法主动往客户端推消息,而 MCP 的工具调用、通知、进度上报又都需要服务端主动推送。
MCP 官方给出的解法是 SSE 传输,核心是双通道设计。下行通道是一条GET /sse的长连接,服务端通过它把 JSON-RPC 响应推给客户端;上行通道是普通的POST /messages/?session_id=xxx,客户端通过它把 JSON-RPC 请求发给服务端。两条通道靠session_id绑定到同一个会话。
这个设计比 WebSocket 简单,完全基于 HTTP,比轮询实时,服务端能主动推。但对 Java 开发者来说,真正的难点不在协议本身,而在于:请求和响应走了两条不同的通道,而且是异步推送的,你怎么知道收到的这条 SSE 消息对应的是哪个请求?答案是用请求 ID 做异步匹配。
这篇文章聚焦 Java 场景,把双通道初始化、SSE 事件流解析、请求 ID 匹配这条完整链路拆开,给你可复制的配置骨架和验证动作。适合需要对接远程 MCP 服务、或者想理解 SSE 传输机制的 Java 开发者。读完你能在本地复现并验证传输层行为,而不是停留在“连上就能用”的模糊认知。
2. TaoToken 前置准备:拿到 Base URL、API Key 和 Model ID
在写 Java 客户端之前,先把服务端侧的接入信息准备好。我用 TaoToken 作为远程 MCP 服务的接入入口,它提供统一的 API 网关,省去自己搭服务端的麻烦。你需要准备三样东西:Base URL、API Key、Model ID。
Base URL 是https://taotoken.net/api,这是所有请求的根地址。API Key 在控制台的 API Keys 页面生成,格式类似sk-开头的一串字符。Model ID 取决于你要调用的模型,在模型列表里能看到具体标识。
拿到这三样之后,先做一次连通性验证,确认 Key 有效、网络可达。用 curl 发一个最简单的请求:
curl -X POST https://taotoken.net/api/v1/chat/completions \ -H "Authorization: Bearer sk-你的Key" \ -H "Content-Type: application/json" \ -d '{ "model": "你的ModelID", "messages": [{"role": "user", "content": "ping"}], "max_tokens": 10 }'如果返回正常的 JSON 响应,说明 Base URL 和 Key 都没问题。如果返回 401,检查 Key 是否复制完整、有没有多余空格。如果连接超时,检查网络是否能访问taotoken.net。
这一步很关键,因为后面 Java 客户端的所有请求都会复用这套凭证。我建议把这三个值写进配置文件,而不是硬编码在代码里。比如用一个mcp.properties:
mcp.base.url=https://taotoken.net/api mcp.api.key=sk-你的Key mcp.model.id=你的ModelID mcp.sse.url=https://taotoken.net/api/sse注意mcp.sse.url是 SSE 长连接的地址,通常是在 Base URL 后面加/sse。不同服务端的路径可能不同,以实际文档为准。TaoToken 的接入文档里有完整的端点说明,配置前先对一遍。
如果你还没生成 Key,去控制台的 API Keys 页面创建一个。创建时注意权限范围,MCP 场景通常需要读写权限。Key 只在创建时显示一次,记得保存好。
3. 可复制的 Java 配置骨架:双通道初始化与异步匹配
这一节给你完整的 Java 代码骨架,基于 OkHttp 的 EventSource 和 CompletableFuture。先看依赖,Maven 里加这几项:
<dependency> <groupId>com.squareup.okhttp3</groupId> <artifactId>okhttp</artifactId> <version>4.12.0</version> </dependency> <dependency> <groupId>com.squareup.okhttp3</groupId> <artifactId>okhttp-sse</artifactId> <version>4.12.0</version> </dependency> <dependency> <groupId>com.fasterxml.jackson.core</groupId> <artifactId>jackson-databind</artifactId> <version>2.17.0</version> </dependency>核心类叫McpSseConnection,负责管理双通道和异步匹配。先定义请求和响应的数据结构:
public class McpRequest { public String jsonrpc = "2.0"; public Integer id; public String method; public Object params; } public class McpResponse { public String jsonrpc; public Object id; // 兼容数字和字符串 public Object result; public McpError error; } public class McpError { public int code; public String message; public Object data; }注意McpResponse.id用Object类型,因为 JSON-RPC 规范允许 id 是数字或字符串,不同服务端实现可能不一样。后面匹配时会统一转成 Integer。
接下来是连接类的核心字段:
public class McpSseConnection { private final String serverName; private final String sseUrl; private final OkHttpClient httpClient; private final ObjectMapper mapper = new ObjectMapper(); private EventSource eventSource; private String messageEndpoint; private volatile boolean connected = false; private final CountDownLatch endpointLatch = new CountDownLatch(1); private final Map<Integer, CompletableFuture<McpResponse>> pendingResponses = new ConcurrentHashMap<>(); private final AtomicInteger requestIdSeq = new AtomicInteger(1); }pendingResponses是异步匹配的核心,key 是请求 ID,value 是等待响应的 Future。endpointLatch用来等待服务端推送 endpoint 事件。
建立连接的方法:
public void connect() throws IOException { try { startSseConnection(); boolean received = endpointLatch.await(10, TimeUnit.SECONDS); if (!received || messageEndpoint == null) { throw new IOException("等待 endpoint 事件超时,请检查服务端是否正常运行"); } connected = true; performHandshake(); log.info("[{}] SSE 连接建立完成,endpoint: {}", serverName, messageEndpoint); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new IOException("连接被中断", e); } }startSseConnection启动 SSE 监听,注册事件回调:
private void startSseConnection() { Request request = new Request.Builder() .url(sseUrl) .header("Accept", "text/event-stream") .header("Cache-Control", "no-cache") .build(); EventSource.Factory factory = EventSources.createFactory(httpClient); eventSource = factory.newEventSource(request, new EventSourceListener() { @Override public void onEvent(EventSource es, String id, String type, String data) { if ("endpoint".equals(type)) { handleEndpointEvent(data); } else if ("message".equals(type) || type == null) { processResponse(data); } } @Override public void onFailure(EventSource es, Throwable t, Response resp) { log.error("[{}] SSE 连接异常", serverName, t); connected = false; } @Override public void onClosed(EventSource es) { connected = false; } }); }handleEndpointEvent解析服务端推送的请求端点,拼上 Base URL:
private void handleEndpointEvent(String data) { String baseUrl = sseUrl.substring(0, sseUrl.lastIndexOf("/sse")); if (data.startsWith("http")) { messageEndpoint = data; } else if (data.startsWith("/")) { messageEndpoint = baseUrl + data; } else { messageEndpoint = baseUrl + "/" + data; } endpointLatch.countDown(); }发送请求和异步匹配是重点。先注册 Future,再发 POST,然后阻塞等待:
public synchronized McpResponse sendRequest(String method, Object params) throws Exception { int currentId = requestIdSeq.getAndIncrement(); McpRequest request = new McpRequest(); request.id = currentId; request.method = method; request.params = params; CompletableFuture<McpResponse> future = new CompletableFuture<>(); pendingResponses.put(currentId, future); try { sendHttpPost(mapper.writeValueAsString(request)); return future.get(30, TimeUnit.SECONDS); } catch (TimeoutException e) { pendingResponses.remove(currentId); throw new IOException("请求超时:" + method); } }收到 SSE 消息后,按 id 取出 Future 并 complete:
private void processResponse(String data) { try { McpResponse response = mapper.readValue(data, McpResponse.class); Integer id = extractId(response); if (id == null) { return; } CompletableFuture<McpResponse> future = pendingResponses.remove(id); if (future == null) { log.warn("[{}] 收到未知请求的响应,id={}", serverName, id); return; } if (response.error != null) { future.completeExceptionally( new McpException(response.error.code, response.error.message, response.error.data)); } else { future.complete(response); } } catch (Exception e) { log.error("[{}] 响应解析失败:{}", serverName, data, e); } }extractId处理 id 类型兼容:
private Integer extractId(McpResponse response) { if (response.id instanceof Integer) { return (Integer) response.id; } else if (response.id instanceof String) { try { return Integer.parseInt((String) response.id); } catch (NumberFormatException e) { log.warn("无效的响应 id:{}", response.id); return null; } } return null; }发送 HTTP POST 的方法:
private static final MediaType JSON = MediaType.get("application/json"); private void sendHttpPost(String requestJson) throws IOException { RequestBody body = RequestBody.create(requestJson, JSON); Request request = new Request.Builder() .url(messageEndpoint) .post(body) .header("Content-Type", "application/json") .build(); try (Response response = httpClient.newCall(request).execute()) { if (!response.isSuccessful()) { throw new IOException("HTTP POST 失败,状态码:" + response.code()); } } }关闭连接时,主动 complete 所有未完成的 Future,避免调用方永久阻塞:
public void close() { connected = false; if (eventSource != null) { eventSource.cancel(); } pendingResponses.forEach((id, future) -> future.completeExceptionally(new IOException("连接已关闭"))); pendingResponses.clear(); }这套骨架的核心就一句话:请求线程先把 Future 放进 Map,然后阻塞等待;SSE 线程收到响应后按 id 从 Map 取出 Future,complete 它,请求线程随即被唤醒。理解了这个模式,其余代码都是工程细节。
4. 验证请求与成功结果:本地联调步骤
代码写完了,怎么确认它真的能跑通?我按顺序给你验证动作。
第一步,启动连接。写一个 main 方法:
public static void main(String[] args) throws Exception { McpSseConnection conn = new McpSseConnection( "taotoken", "https://taotoken.net/api/sse", new OkHttpClient.Builder() .connectTimeout(10, TimeUnit.SECONDS) .readTimeout(0, TimeUnit.SECONDS) // SSE 长连接不设读超时 .build() ); conn.connect(); System.out.println("连接成功,endpoint: " + conn.getMessageEndpoint()); }注意readTimeout要设为 0,否则 OkHttp 会在读超时后断开 SSE 长连接。这是很多人踩过的坑。
第二步,调用initialize完成握手。MCP 协议要求先初始化:
McpResponse initResp = conn.sendRequest("initialize", Map.of( "protocolVersion", "2024-11-05", "capabilities", Map.of(), "clientInfo", Map.of("name", "java-client", "version", "1.0") )); System.out.println("initialize 响应: " + initResp.result);如果成功,你会看到服务端返回的能力列表,包含protocolVersion、serverInfo、capabilities等字段。
第三步,发送initialized通知。通知没有 id,不需要注册 Future:
conn.sendNotification("notifications/initialized", Map.of());第四步,调用tools/list列出可用工具:
McpResponse toolsResp = conn.sendRequest("tools/list", Map.of()); System.out.println("工具列表: " + toolsResp.result);成功的话,控制台会打印出工具数组,每个工具包含name、description、inputSchema。
第五步,调用一个具体工具,比如callTool:
McpResponse callResp = conn.sendRequest("tools/call", Map.of( "name", "你的工具名", "arguments", Map.of("param1", "value1") )); System.out.println("工具调用结果: " + callResp.result);整个流程跑通后,你会看到类似这样的日志:
[taotoken] SSE 连接建立完成,endpoint: https://taotoken.net/api/messages/?session_id=abc123 initialize 响应: {protocolVersion=2024-11-05, serverInfo={...}, capabilities={...}} 工具列表: {tools=[{name=..., description=..., inputSchema=...}]} 工具调用结果: {content=[{type=text, text=...}]}如果某一步卡住,先看日志里有没有等待 endpoint 事件超时,再看pendingResponses里是不是有未完成的 Future。用 jstack 抓线程栈,能看到请求线程是不是阻塞在future.get()。
5. 本篇常见错误排查:401、local proxy failed、reading choices、OAuth
实际联调时,报错往往集中在几个地方。我按真实遇到的顺序列出来。
401 Unauthorized。这是最常见的,通常是 API Key 没带对。检查三处:请求头是不是Authorization: Bearer sk-xxx,Key 有没有多余空格,Key 是不是已经过期或被删除。如果用的是 TaoToken,去控制台确认 Key 状态。还有一种情况是 SSE 长连接和 POST 请求用了不同的 Key,导致一边通一边不通。
local proxy failed。这个报错通常出现在客户端配置了本地代理,但代理进程没启动或端口不对。检查你的 HTTP 客户端有没有走系统代理,OkHttp 默认会用ProxySelector.getDefault()。如果不需要代理,显式设置proxy(Proxy.NO_PROXY)。另外检查环境变量HTTP_PROXY、HTTPS_PROXY有没有设置成无效值。
reading choices 相关报错。这个一般出现在响应解析阶段,说明返回的 JSON 结构和你的McpResponse对不上。常见原因是服务端返回的是 OpenAI 格式的choices数组,而你的代码按 MCP 的result字段解析。确认你调用的端点是不是 MCP 端点,而不是普通的 chat completions 端点。MCP 的 SSE 端点是/sse,POST 端点是/messages/,别搞混。
OAuth 相关报错。如果服务端要求 OAuth 认证,而你的请求只带了 API Key,会返回 401 或 403 并附带 OAuth 挑战头。检查响应头里有没有WWW-Authenticate。MCP 的 OAuth 流程需要先获取 access token,再拿 token 去请求。如果你用的是 API Key 模式,确认服务端支持这种认证方式。
Future 永久阻塞。请求发出去了,但future.get()一直不返回,直到超时。原因通常是响应到了但 id 没匹配上。检查extractId有没有正确处理字符串类型的 id,检查pendingResponses.put是不是在sendHttpPost之前执行。顺序反了的话,响应可能在 put 之前就到达,导致 Future 永远不会被 complete。
SSE 连接频繁断开。如果日志里反复出现onFailure和重连,检查readTimeout是不是设成了非 0 值。另外有些中间层会在 60 秒无数据时断开空闲连接,服务端会定期发: ping注释行保活,客户端要能忽略注释行而不是当成错误。
排查时我习惯先看 HTTP 状态码,再看响应体,最后看线程栈。大部分问题在前两步就能定位。
6. 从传输层到工程落地:把 SSE 接入你的 Java 项目
把上面的骨架接进真实项目时,还有几个工程细节值得处理。
连接池和重连。生产环境不能只建一条连接就完事,要处理断线重连。我的做法是给McpSseConnection加一个状态监听器,在onFailure和onClosed里触发重连逻辑,用指数退避避免雪崩。重连后要重新走connect()和performHandshake(),因为session_id会变。
超时分层。连接超时、请求超时、SSE 读超时是三个不同的概念。连接超时控制 TCP 握手,请求超时控制future.get()的等待,SSE 读超时控制长连接的存活。我一般设连接超时 10 秒,请求超时 30 秒,SSE 读超时 0(不超时)。
内存泄漏防护。如果请求超时后 Future 没被清理,pendingResponses会持续增长。加一个定时清理任务:
private final ScheduledExecutorService cleaner = Executors.newSingleThreadScheduledExecutor(r -> { Thread t = new Thread(r, "mcp-cleaner-" + serverName); t.setDaemon(true); return t; }); public void startCleaner() { cleaner.scheduleAtFixedRate(() -> { pendingResponses.entrySet().removeIf(entry -> { if (!entry.getValue().isDone()) { entry.getValue().completeExceptionally( new TimeoutException("请求已过期")); return true; } return false; }); }, 60, 60, TimeUnit.SECONDS); }如果你用 j-langchain 这类框架,这些细节已经被封装在McpSseConnection内部。McpConnectionFactory.createConnection("name", sseConfig)创建连接实例,connect()完成 SSE 建立和握手,后续调用listTools()和callTool()与 Stdio 方式完全一致,使用方感知不到传输层的差异。这也是双通道设计的好处:协议层统一,传输层可替换。
最后给一个实用建议:把 Base URL、API Key、Model ID 三件套统一放在配置中心或环境变量里,代码里只读不写。切换环境时改配置就行,不用重新编译。验证模型连通性可以用模型对话页面快速测一下,长期跑编码或 Agent 任务的话,Coding Plan 的额度更划算。接入文档里有完整的端点说明和示例,配置前对一遍能省不少排查时间。