如何实现Apollo Server订阅中间件以拦截请求并转发至远程服务器?
Apollo Server订阅请求转发至远程服务器实现方案
步骤1:安装依赖
首先需要安装ws库来处理WebSocket连接:
npm install ws # 或使用yarn yarn add ws
步骤2:实现自定义订阅转发服务器
创建自定义WebSocket服务器,处理客户端连接并转发请求到远程服务:
import { ApolloServerExpressConfig, ApolloServer } from 'apollo-server-express'; import { WebSocketServer, WebSocket } from 'ws'; import { createServer as createHttpServer } from 'http'; import express from 'express'; // 替换为你的远程GraphQL订阅服务地址 const REMOTE_SUBSCRIPTION_URL = 'ws://remote-server-domain.com/graphql/subscriptions'; // 创建负责转发的订阅服务器 function createForwardingSubscriptionServer(schema, httpServer) { const wsServer = new WebSocketServer({ server: httpServer, path: '/core/graphql/subscriptions', // 与原Apollo配置的订阅路径保持一致 }); wsServer.on('connection', (clientWs) => { let remoteWs: WebSocket; // 监听客户端发来的消息 clientWs.on('message', async (data) => { const message = JSON.parse(data.toString()); // 处理客户端的初始化连接请求 if (message.type === 'connection_init') { const connectionParams = message.payload; // 根据connectionParams创建到远程服务器的WebSocket连接 remoteWs = new WebSocket(REMOTE_SUBSCRIPTION_URL, { headers: { // 示例:将客户端传入的认证令牌转发到远程服务 ...(connectionParams?.authToken && { Authorization: `Bearer ${connectionParams.authToken}` }), // 可根据远程服务要求添加其他请求头 }, }); // 转发远程服务器的响应到客户端 remoteWs.on('message', (remoteData) => { clientWs.send(remoteData.toString()); }); // 转发客户端的后续消息到远程服务器(跳过已处理的connection_init) clientWs.on('message', (clientData) => { const clientMsg = JSON.parse(clientData.toString()); if (clientMsg.type !== 'connection_init') { remoteWs.send(clientData.toString()); } }); // 同步关闭两端连接 clientWs.on('close', () => remoteWs?.close()); remoteWs.on('close', () => clientWs.close()); // 处理远程连接错误 remoteWs.on('error', (err) => { console.error('远程WebSocket连接错误:', err); clientWs.close(1011, '远程服务连接异常'); }); } }); // 处理客户端连接错误 clientWs.on('error', (err) => { console.error('客户端WebSocket连接错误:', err); remoteWs?.close(); }); }); return wsServer; }
步骤3:修改Apollo Server启动逻辑
放弃默认的订阅配置,改用自定义转发服务器:
// 假设你已定义了mergedSchema、mergedSchemaWithPermissions、TEST_MODE、GraphQLContext等变量 async function startApolloServer() { const app = express(); const httpServer = createHttpServer(app); // 初始化转发订阅服务器 createForwardingSubscriptionServer( TEST_MODE ? mergedSchema : mergedSchemaWithPermissions, httpServer ); const graphqlServerConfig: ApolloServerExpressConfig = { schema: TEST_MODE ? mergedSchema : mergedSchemaWithPermissions, context: ({ req, res, connection }): GraphQLContext => { return { req, res, csctx, }; }, // 移除原有的subscriptions配置,使用自定义实现 playground: { endpoint: '/core/graphql/', subscriptionEndpoint: '/core/graphql/subscriptions', settings: { 'request.credentials': 'include', }, }, }; const server = new ApolloServer(graphqlServerConfig); await server.start(); server.applyMiddleware({ app, path: '/core/graphql/' }); // 启动HTTP服务器 httpServer.listen(4000, () => { console.log(`服务启动完成:http://localhost:4000/core/graphql/`); console.log(`订阅服务地址:ws://localhost:4000/core/graphql/subscriptions`); }); } startApolloServer();
核心逻辑说明
- 连接初始化:当客户端发送
connection_init消息时,提取其中的connectionParams,以此配置到远程服务器的WebSocket连接(比如认证信息)。 - 双向消息转发:客户端的后续订阅请求直接转发到远程服务器,远程返回的订阅响应原样回传给客户端。
- 连接同步:客户端或远程服务任意一端关闭连接时,同步关闭另一端的连接,避免资源泄漏。
- 错误处理:捕获两端的连接错误,及时关闭连接并反馈给客户端。
内容的提问来源于stack exchange,提问作者Sagar Saud
相关产品推荐
相关产品推荐

