实战指南:用 Java Netty 构建轻量级 MCP 中间件,打通大模型与企业级应用 2026-10-07 程序之旅,记录 暂无评论 4 次阅读 # 实战指南:用 Java Netty 构建轻量级 MCP 中间件,打通大模型与企业级应用 随着大模型(LLM)向深度业务场景渗透,如何让 AI 安全、高效地调用企业内部的存量系统(如 ERP、CRM、核心数据库)成为了一个痛点。Anthropic 推出的 **Model Context Protocol (MCP)** 提供了一套完美的标准协议。 如果你的企业系统基于庞大的微服务架构(Spring Cloud / Dubbo),直接在老业务代码中强行揉入 MCP 协议不仅极具侵入性,还可能引发稳定性风险。此时,开发一个轻量级的“MCP 中间件 / 网关模块”是最优雅的解法:它对外面向大模型提供标准的 MCP-SSE 接口,对内通过 RPC 或 HTTP 桥接企业旧系统。 在 Java 生态中,**Netty** 无疑是开发这类高并发、低内存占用中间件的最佳选择。今天,我们就来实战如何用 Netty 从零编写一个高性能的 MCP 服务端。 ## 一、 为什么选择 Netty? 相比于 Spring Boot,在这个特定场景下 Netty 具有压倒性优势: 1. **极致轻量**:没有臃肿的依赖注入和 Web 容器启动过程,打出来的 Jar 包极小,内存占用可控制在几十 MB,非常适合作为 Sidecar(边车)或代理网关部署。 2. **高并发长连接**:MCP 的云端通信依赖 SSE(Server-Sent Events)长连接。Netty 的异步事件驱动架构天然契合海量长连接场景。 3. **优雅的保活机制**:无需像 Spring 那样启动全局定时任务,Netty 原生的 `IdleStateHandler` 可以零开销实现完美的链路心跳检测。 ## 二、 核心架构:解构 MCP 的 SSE 通信 在 Netty 中,我们要用最底层的 HTTP 报文来模拟 MCP 的全双工通信。你需要实现两个核心端点: 1. **`GET /sse` (建立连接)**:客户端发起 GET 请求。服务端不要关闭 Channel,而是返回 `Transfer-Encoding: chunked`,并在连接成功时下发一个 `endpoint` 事件,告诉客户端后续的指令发往何处。 2. **`POST /message` (接收指令)**:客户端将 JSON-RPC 请求发到这里。服务端解析指令、调用企业内部应用后,将结果通过刚才 `GET` 阶段保存下来的 Channel 推送回去。 ## 三、 代码实战:Netty MCP 中间件实现 ### 1. Channel Pipeline 初始化与心跳保活 这是 Netty 的精髓。我们通过 `IdleStateHandler` 设定:如果 15 秒没有向客户端推送数据,就触发一个写空闲事件,下发 `ping` 报文保活,防止被 Nginx 等中间网关掐断连接。 ```java import io.netty.channel.ChannelInitializer; import io.netty.channel.ChannelPipeline; import io.netty.channel.socket.SocketChannel; import io.netty.handler.codec.http.HttpObjectAggregator; import io.netty.handler.codec.http.HttpServerCodec; import io.netty.handler.timeout.IdleStateHandler; public class McpServerInitializer extends ChannelInitializer { @Override protected void initChannel(SocketChannel ch) { ChannelPipeline pipeline = ch.pipeline(); // 1. HTTP 编解码器与报文聚合 (限制最大请求体为 1MB) pipeline.addLast(new HttpServerCodec()); pipeline.addLast(new HttpObjectAggregator(1048576)); // 2. 核心保活:设置写空闲时间为 15 秒 pipeline.addLast(new IdleStateHandler(0, 15, 0)); // 3. 自定义心跳处理器(捕获空闲事件并发送 ping) pipeline.addLast(new McpKeepAliveHandler()); // 4. MCP 业务路由与处理器 pipeline.addLast(new McpSseHandler()); } } ``` **心跳处理器实现:** ```java import io.netty.buffer.Unpooled; import io.netty.channel.ChannelDuplexHandler; import io.netty.channel.ChannelHandlerContext; import io.netty.handler.codec.http.DefaultHttpContent; import io.netty.handler.timeout.IdleState; import io.netty.handler.timeout.IdleStateEvent; import io.netty.util.CharsetUtil; public class McpKeepAliveHandler extends ChannelDuplexHandler { @Override public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception { if (evt instanceof IdleStateEvent) { IdleStateEvent e = (IdleStateEvent) evt; if (e.state() == IdleState.WRITER_IDLE) { // 触发写空闲,下发 SSE Ping 事件 String pingEvent = "event: ping\ndata: keep-alive\n\n"; ctx.writeAndFlush(new DefaultHttpContent(Unpooled.copiedBuffer(pingEvent, CharsetUtil.UTF_8))); } } else { super.userEventTriggered(ctx, evt); } } } ``` ### 2. 核心业务 Handler:处理 GET 与 POST 请求 在这里,我们全局维护一个 `SessionID -> Channel` 的映射字典。 ```java import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; import io.netty.channel.*; import io.netty.handler.codec.http.*; import io.netty.util.CharsetUtil; import java.util.Map; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; public class McpSseHandler extends SimpleChannelInboundHandler { // 保存所有活跃的长连接 Session private static final Map activeSessions = new ConcurrentHashMap<>(); @Override protected void channelRead0(ChannelHandlerContext ctx, FullHttpRequest request) { String uri = request.uri(); HttpMethod method = request.method(); // ========================================== // 阶段 1:建立 SSE 长连接 // ========================================== if (HttpMethod.GET.equals(method) && uri.startsWith("/sse")) { String sessionId = UUID.randomUUID().toString(); // 构造 SSE 协议规范的 HTTP 响应头 HttpResponse response = new DefaultHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK); response.headers().set(HttpHeaderNames.CONTENT_TYPE, "text/event-stream; charset=UTF-8"); response.headers().set(HttpHeaderNames.CACHE_CONTROL, "no-cache"); response.headers().set(HttpHeaderNames.CONNECTION, "keep-alive"); response.headers().set(HttpHeaderNames.TRANSFER_ENCODING, "chunked"); // 开启分块传输 // 发送响应头(注意:只 flush Header,绝对不关闭 Channel) ctx.writeAndFlush(response); // 注册到会话管理器,并监听连接断开事件以防内存泄漏 activeSessions.put(sessionId, ctx.channel()); ctx.channel().closeFuture().addListener(future -> activeSessions.remove(sessionId)); // MCP 核心规范:连接成功后必须立即告知客户端 POST 接口的路径 sendSseEvent(ctx.channel(), "endpoint", "/message?sessionId=" + sessionId); return; } // ========================================== // 阶段 2:接收客户端 JSON-RPC 请求 // ========================================== if (HttpMethod.POST.equals(method) && uri.startsWith("/message")) { // 解析参数获取目标 Session String sessionId = extractSessionId(uri); Channel targetChannel = activeSessions.get(sessionId); if (targetChannel == null || !targetChannel.isActive()) { sendErrorResponse(ctx, HttpResponseStatus.BAD_REQUEST, "Session timeout or not found"); return; } // 提取 JSON-RPC 报文 (通常是 Tools 调用请求) String jsonRpcRequest = request.content().toString(CharsetUtil.UTF_8); System.out.println("收到模型指令: " + jsonRpcRequest); // 立即给当前 POST 请求回复 202 Accepted,释放客户端的 HTTP 等待 FullHttpResponse postResponse = new DefaultFullHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.ACCEPTED); ctx.writeAndFlush(postResponse).addListener(ChannelFutureListener.CLOSE); // ⚠️ 在这里桥接企业应用: // 解析 jsonRpcRequest,根据 method 调用公司内部的 Dubbo 接口或 HTTP 服务... // 此处模拟企业应用返回的结果: String mockJsonResponse = "{\"jsonrpc\": \"2.0\", \"id\": 1, \"result\": {\"content\": [{\"type\": \"text\", \"text\": \"订单 #9527 查询成功,状态:已发货\"}]}}"; // 将企业应用的执行结果,通过刚才的长连接 Channel 推送给大模型 sendSseEvent(targetChannel, "message", mockJsonResponse); } } /** * 将字符串打包成 Chunked 格式并写入 Channel */ private void sendSseEvent(Channel channel, String eventName, String data) { if (channel != null && channel.isActive()) { String ssePayload = String.format("event: %s\ndata: %s\n\n", eventName, data); ByteBuf buffer = Unpooled.copiedBuffer(ssePayload, CharsetUtil.UTF_8); // 注意:必须包装为 DefaultHttpContent 才能被 Netty 作为 Chunk 发送 channel.writeAndFlush(new DefaultHttpContent(buffer)); } } private void sendErrorResponse(ChannelHandlerContext ctx, HttpResponseStatus status, String msg) { FullHttpResponse response = new DefaultFullHttpResponse( HttpVersion.HTTP_1_1, status, Unpooled.copiedBuffer(msg, CharsetUtil.UTF_8)); ctx.writeAndFlush(response).addListener(ChannelFutureListener.CLOSE); } private String extractSessionId(String uri) { try { return uri.split("sessionId=")[1]; } catch (Exception e) { return ""; } } } ``` ## 四、 协议精讲:JSON-RPC 2.0 交互规范 在 `POST /message` 接收到的数据,以及我们通过 `sendSseEvent` 返回的数据,都必须严格遵守 **JSON-RPC 2.0** 规范。作为中间件开发者,你只需要搞懂以下 4 种报文格式即可完成与内部业务逻辑的映射: ### 1. 客户端发来的请求报文 (Request) 大模型想要执行动作(比如调用工具)。必须带 `id` 字段。 ```json { "jsonrpc": "2.0", "method": "tools/call", "params": { "name": "query_order", "arguments": { "orderId": "9527" } }, "id": "req-001" } ``` ### 2. 我们推送回去的成功响应 (Success Response) 内部业务执行完毕后,回传给大模型的报文。**必须包含 `result` 字段,绝对不能有 `error`,且 `id` 必须原样奉还。** ```json { "jsonrpc": "2.0", "result": { "content": [ { "type": "text", "text": "订单已发货" } ] }, "id": "req-001" } ``` ### 3. 我们推送回去的错误响应 (Error Response) 如果内部 Dubbo 接口报错、或者大模型传参非法,返回此报文。**必须包含 `error` 字段,绝对不能有 `result`。** (标准错误码:`-32601` 方法不存在,`-32602` 参数无效,`-32603` 内部报错) ```json { "jsonrpc": "2.0", "error": { "code": -32603, "message": "Internal service error", "data": "下游订单中心 RPC 调用超时" }, "id": "req-001" } ``` ### 4. 通知报文 (Notification) 不需要结果的单向通信(通常是握手完成时)。**特征是没有 `id` 字段。** 收到这种报文,我们的 Netty 服务绝不能回复任何 Response。 ```json { "jsonrpc": "2.0", "method": "notifications/initialized" } ``` ## 结语 使用 Netty 编写 MCP 中间件,我们将复杂的大模型通信协议与企业级厚重的业务逻辑实现了解耦。这个 Netty 模块不仅占用极小,还能充当流量漏斗、权限校验层和审计日志节点。 通过这种“边车”模式,你可以在不修改老旧系统一行代码的前提下,将企业积累的 IT 资产全部暴露给 Claude 或其他支持 MCP 的 AI Agent,真正实现企业系统的 AI 智能化升级。 打赏: 微信, 支付宝 标签: Netty, ai, mcp 本作品采用 知识共享署名-相同方式共享 4.0 国际许可协议 进行许可。