GraphQL 订阅、SSE 与 WebSocket 实时推送实战

GraphQL Subscription 三种传输方案深度对比:WebSocket(graphql-ws)、SSE(Server-Sent Events)、HTTP 长轮询。含 Apollo/GraphQL Yoga 完整实现代码、认证方案与广播策略。

实时推送是现代 Web 应用的核心能力。从即时消息到股票行情,再到多人协作编辑,都离不开一条稳定、低延迟的数据通道。GraphQL 在 2016 年规范中引入 Subscription 操作类型,为服务端主动推送数据提供了统一语义。

但规范只定义了订阅语法与执行语义,并未绑定传输协议。落地中通常面临三种选择:WebSocket(双向全双工)、SSE(单向推送)、HTTP 长轮询(兼容性兜底)。本文围绕三种方案深度对比,覆盖协议原理、认证方案、Apollo/Yoga 完整代码,并延伸到 Redis PubSub 广播与多节点负载均衡。

实时推送的典型场景

  • 即时消息与聊天:用户 A 发消息需立即推送给对方或群组,要求低延迟、有序到达。
  • 系统通知与告警:监控告警、订单变更、社交提醒,需支持去重和批量聚合。
  • 实时数据仪表盘:金融行情、IoT 遥测、业务大盘,数据量大更新频繁。
  • 协作编辑:Notion、Figma 类工具,需操作序列化、光标同步,通常要双向高频通信。

聊天和协作编辑需要双向通道,仪表盘和通知仅需服务端单向推送。

GraphQL Subscription 规范解析

Subscription Operation 语法

subscription OnMessageAdded($channelId: ID!) {
  messageAdded(channelId: $channelId) {
    id
    content
    sender { id name }
    createdAt
  }
}

与 Query 的关键区别:root field resolver 必须返回一个 AsyncIterable,引擎每次迭代出新 payload 时序列化为 GraphQL 响应推送客户端。

Root Field Resolver 与 AsyncIterable

const resolvers = {
  Subscription: {
    messageAdded: {
      subscribe: async (_parent, args) =>
        pubsub.asyncIterator(`MESSAGE_ADDED_${args.channelId}`),
      resolve: (payload) => payload.messageAdded,
    },
  },
};

subscribe 返回 AsyncIterator,Apollo Server 通过 for await...of 消费,每次结果传入 resolve 解析。生命周期分四阶段:连接建立、初始化与认证、注册订阅、事件推送与取消。


WebSocket 方案:graphql-ws v5

WebSocket 是 GraphQL 订阅最主流的传输方案。graphql-ws 实现了现代化的 GraphQL over WebSocket 协议(相比已废弃的 subscriptions-transport-ws,v5 在安全、心跳和语义上显著改进)。

协议升级与心跳保活

WebSocket 始于 HTTP Upgrade 请求,服务端返回 101 Switching Protocols 后切换为帧协议,Query/Mutation/Subscription 均可在全双工通道复用。graphql-ws 内置双向 ping/pong 心跳,超时未收到帧即判定为死连接自动关闭:

import { useServer } from 'graphql-ws/lib/use/ws';
import { WebSocketServer } from 'ws';

const wsServer = new WebSocketServer({ server: httpServer, path: '/graphql' });

useServer({
  schema,
  keepAlive: 12000,
  onConnect: async (ctx) => {
    const token = ctx.connectionParams?.token;
    if (!token) return false;
    ctx.extra.user = await verifyToken(token);
    return true;
  },
}, wsServer);

connectionInit 认证

graphql-ws v5 用 connectionInit 解决连接未建立时的认证难题:

import { GraphQLWsLink } from '@apollo/client/link/subscriptions';
import { createClient } from 'graphql-ws';

const wsLink = new GraphQLWsLink(
  createClient({
    url: 'wss://api.example.com/graphql',
    connectionParams: async () => ({
      token: await getAccessToken(), clientName: 'web-dashboard',
    }),
    retryAttempts: 5,
    shouldRetry: () => true,
  })
);

connectionParams 支持返回 Promise 的函数,适合 Token 刷新后注入新凭证。

Apollo Client 订阅代码

使用 split 将订阅路由到 WebSocket Link:

import { ApolloClient, InMemoryCache, split, HttpLink } from '@apollo/client';
import { getMainDefinition } from '@apollo/client/utilities';
import { GraphQLWsLink } from '@apollo/client/link/subscriptions';
import { createClient } from 'graphql-ws';

const httpLink = new HttpLink({ uri: 'https://api.example.com/graphql' });
const wsLink = new GraphQLWsLink(createClient({
  url: 'wss://api.example.com/graphql',
  connectionParams: { token: getAccessToken() },
}));

const splitLink = split(
  ({ query }) => {
    const def = getMainDefinition(query);
    return def.kind === 'OperationDefinition' && def.operation === 'subscription';
  },
  wsLink, httpLink
);

const client = new ApolloClient({ link: splitLink, cache: new InMemoryCache() });

// 组件中使用
function MessageStream({ channelId }: { channelId: string }) {
  const { data, error } = useSubscription(MESSAGE_ADDED_SUBSCRIPTION, {
    variables: { channelId },
  });
  if (error) return <div>连接异常: {error.message}</div>;
  if (!data) return <div>等待消息...</div>;
  return <MessageCard message={data.messageAdded} />;
}

订阅权限控制(AbortOperation)

useServer({
  schema,
  onSubscribe: async (ctx, msg) => {
    const { channelId } = msg.payload.variables || {};
    if (!(await isMember(ctx.extra.user.id, channelId))) {
      return [new GraphQLError('无权订阅此频道', { extensions: { code: 'FORBIDDEN' } })];
    }
    return undefined;
  },
}, wsServer);

这种细粒度控制是 SSE 较难实现的——SSE 每个 HTTP 请求对应单一数据流,缺乏会话上下文。


SSE 方案:Server-Sent Events

如果业务只需服务端向客户端单向推送,SSE 是比 WebSocket 更轻量的选择,基于 HTTP/1.1 chunked transfer 或 HTTP/2 multiplexing,使用 text/event-stream

协议原理与自动重连

SSE 以 event:data:id: 字段组织,每块以两个换行符分隔:

HTTP/1.1 200 OK
Content-Type: text/event-stream
Cache-Control: no-cache

id: 42
event: message
data: {"data":{"messageAdded":{"id":"1","content":"Hello"}}}

原生 EventSource 内置自动重连,但不支持自定义请求头,这是认证方面的最大局限。

认证方案

  • Cookie/Session:同域或跨域配置正确时,SSE 自动携带 Cookie,最自然。
  • URL Query Token:将短期 Token 附加到查询参数,会暴露在浏览器历史和日志中,生产环境应使用一次性 Token。
  • Fetch + ReadableStream:现代浏览器支持 fetch ReadableStream,可替代 EventSource 并自由设置 Authorization:
async function createSseStream(url: string, token: string) {
  const response = await fetch(url, {
    headers: { Accept: 'text/event-stream', Authorization: `Bearer ${token}` },
  });
  const reader = response.body!.getReader();
  const decoder = new TextDecoder();
  while (true) {
    const { done, value } = await reader.read();
    if (done) break;
    const chunk = decoder.decode(value, { stream: true });
    // 解析 SSE 帧...
  }
}

SSE 单向语义决定它不适合聊天输入、光标同步等场景。但在实时通知、仪表盘、行情推送等"服务端广播"场景下有独特优势:协议简单无 Upgrade、原生自动重连、Last-Event-ID 恢复断线消息、防火墙友好、HTTP/2 下不受单域名连接数限制。


HTTP 长轮询:兼容性兜底

在旧浏览器(IE11)、极端严格防火墙/代理、旧版 Serverless 平台仍有价值。

客户端发起 HTTP 请求,服务端保持挂起直到有新数据或超时(通常 25-30 秒)。Apollo Client 可检测 WebSocket 支持情况做降级:

const client = new ApolloClient({
  link: split(
    ({ query }) => getMainDefinition(query).operation === 'subscription',
    typeof WebSocket !== 'undefined' ? wsLink : httpPollingLink,
    httpLink,
  ),
  cache: new InMemoryCache(),
});

GraphQL 规范未定义标准 HTTP 长轮询语义,仅在极端兼容需求下作为降级策略。


三大方案选型总表

维度WebSocket (graphql-ws)SSE (Server-Sent Events)HTTP 长轮询
协议复杂度高,需 Upgrade 握手和帧协议低,基于标准 HTTP最低,纯 HTTP
通信方向全双工(双向)单向(服务端→客户端)单向
实时性极低延迟低延迟中延迟
自动重连客户端库实现原生 EventSource 内置手动轮询
浏览器兼容性现代浏览器 + IE10+现代浏览器(IE 不支持原生 API)全浏览器兼容
认证便利性connectionInit / connectionParamsCookie 完美支持;Token 需 URL 或 Fetch 方案与普通 HTTP 一致
防火墙友好度一般,代理可能拦截 Upgrade极高,标准 HTTP极高
HTTP/2 多路复用独立连接天然共享连接共享连接
可扩展性需 Sticky Session 或集中 PubSub无状态 HTTP,负载均衡友好无状态 HTTP
断线消息恢复需业务层自定义原生 Last-Event-ID需业务层维护游标
适合场景聊天、协作编辑、在线游戏通知、行情、IoT 遥测、仪表盘旧浏览器兼容、极端网络

选型决策树

是否需要客户端高频上报?
├── 是(聊天输入、光标同步)→ WebSocket(graphql-ws)
└── 否(纯服务端推送)
    ├── 是否需要 IE 支持?→ HTTP 长轮询
    ├── 是否需要 Token 认证且不想 URL 泄露?→ WebSocket
    ├── 是否需要极致负载均衡友好性 + 自动重连?→ SSE
    └── 否则 → WebSocket 或 SSE 均可

广播与扩展:从单机到多节点

单机使用内存 PubSub。多节点部署(K8s 多 Pod)时各节点只持本地连接,事件无法跨节点广播,必须引入外部消息中间件作为事件总线。

Redis PubSub

最轻量方案,延迟亚毫秒,但不保证可靠性、不支持持久化,对允许丢失的场景很合适:

import { RedisPubSub } from 'graphql-redis-subscriptions';
import Redis from 'ioredis';

const pubsub = new RedisPubSub({
  publisher: new Redis(redisConfig),
  subscriber: new Redis(redisConfig),
});

const resolvers = {
  Subscription: {
    messageAdded: {
      subscribe: (_parent, { channelId }) =>
        pubsub.asyncIterator(`MESSAGE_ADDED_${channelId}`),
    },
  },
  Mutation: {
    addMessage: async (_parent, { channelId, content }, { user }) => {
      const message = await db.message.create({
        data: { channelId, content, senderId: user.id },
      });
      await pubsub.publish(`MESSAGE_ADDED_${channelId}`, { messageAdded: message });
      return message;
    },
  },
};

RabbitMQ / Kafka

RabbitMQ 支持消费者确认(Ack)、死信队列、路由绑定,可构建可靠广播。Kafka 天生高吞吐,支持消费者组,适合 IoT 遥测、金融行情等海量数据场景。Kafka 消息可通过 AsyncQueue 转发给 GraphQL 迭代器消费。

多节点 Subscription 负载均衡

GraphQL Subscription Connection 是有状态连接。Nginx/ALB 负载均衡时必须保证同一客户端 WebSocket 始终路由到同一后端节点:

upstream graphql_backend {
    ip_hash;
    server graphql-pod-1:4000;
    server graphql-pod-2:4000;
}

server {
    location /graphql {
        proxy_pass http://graphql_backend;
        proxy_http_version 1.1;
        proxy_set_header Upgrade $http_upgrade;
        proxy_set_header Connection "upgrade";
        proxy_read_timeout 86400;
    }
}

ALB 可在 Target Group 中启用 Stickiness,或使用独立 WebSocket 网关层将连接层与业务层分离。


Apollo GraphQL Subscriptions 完整代码示例

Apollo Server v4 + graphql-ws v5 的完整启动代码:

import { ApolloServer } from '@apollo/server';
import { expressMiddleware } from '@apollo/server/express4';
import { ApolloServerPluginDrainHttpServer } from '@apollo/server/plugin/drainHttpServer';
import { makeExecutableSchema } from '@graphql-tools/schema';
import { WebSocketServer } from 'ws';
import { useServer } from 'graphql-ws/lib/use/ws';
import express from 'express';
import { createServer } from 'http';
import { PubSub } from 'graphql-subscriptions';
import { json } from 'body-parser';
import cors from 'cors';

const typeDefs = `#graphql
  type Message { id: ID! content: String! sender: String! createdAt: String! }
  type Query { messages(channelId: ID!): [Message!]! }
  type Mutation { sendMessage(channelId: ID!, content: String!): Message! }
  type Subscription { messageAdded(channelId: ID!): Message! }
`;

const pubsub = new PubSub();

const resolvers = {
  Query: {
    messages: (_parent, { channelId }) =>
      db.messages.filter((m) => m.channelId === channelId),
  },
  Mutation: {
    sendMessage: async (_parent, { channelId, content }, { user }) => {
      const message = {
        id: crypto.randomUUID(), content,
        sender: user?.name || 'anonymous',
        createdAt: new Date().toISOString(),
      };
      await pubsub.publish(`MSG_${channelId}`, { messageAdded: message });
      return message;
    },
  },
  Subscription: {
    messageAdded: {
      subscribe: (_parent, { channelId }) => pubsub.asyncIterator(`MSG_${channelId}`),
    },
  },
};

const schema = makeExecutableSchema({ typeDefs, resolvers });
const app = express();
const httpServer = createServer(app);

const apolloServer = new ApolloServer({
  schema,
  plugins: [
    ApolloServerPluginDrainHttpServer({ httpServer }),
    {
      async serverWillStart() {
        return { async drainServer() { await serverCleanup.dispose(); } };
      },
    },
  ],
});

await apolloServer.start();

app.use(
  '/graphql',
  cors({ origin: ['https://app.example.com'], credentials: true }),
  json(),
  expressMiddleware(apolloServer, {
    context: async ({ req }) => {
      const token = req.headers.authorization?.replace('Bearer ', '');
      return { user: token ? await verifyToken(token) : null };
    },
  })
);

const wsServer = new WebSocketServer({ server: httpServer, path: '/graphql' });
const serverCleanup = useServer({
  schema,
  onConnect: async (ctx) => {
    const token = ctx.connectionParams?.token as string;
    if (!token) return false;
    ctx.extra.user = await verifyToken(token);
    return true;
  },
  onSubscribe: async (ctx, msg) => {
    const { channelId } = msg.payload.variables || {};
    if (!(await isMember(ctx.extra.user.id, channelId))) return [new GraphQLError('无权订阅')];
    return undefined;
  },
  context: async (ctx) => ({ user: ctx.extra.user }),
}, wsServer);

httpServer.listen(4000, () => {
  console.log('HTTP  http://localhost:4000/graphql');
  console.log('WS    ws://localhost:4000/graphql');
});

Yoga GraphQL SSE 完整代码示例

GraphQL Yoga 对 SSE 提供开箱即用支持:

import { createYoga, createSchema } from 'graphql-yoga';
import { createServer } from 'http';
import { EventEmitter, on } from 'events';

const eventEmitter = new EventEmitter();

const schema = createSchema({
  typeDefs: `#graphql
    type Notification { id: ID! title: String! body: String! createdAt: String! }
    type Query { health: String! }
    type Mutation { notify(userId: ID!, title: String!, body: String!): Notification! }
    type Subscription { notifications(userId: ID!): Notification! }
  `,
  resolvers: {
    Query: { health: () => 'ok' },
    Mutation: {
      notify: (_parent, { userId, title, body }) => {
        const notification = {
          id: crypto.randomUUID(), title, body,
          createdAt: new Date().toISOString(),
        };
        eventEmitter.emit(`NOTIFY_${userId}`, notification);
        return notification;
      },
    },
    Subscription: {
      notifications: {
        subscribe: async function* (_parent, { userId }) {
          const events = on(eventEmitter, `NOTIFY_${userId}`);
          for await (const [n] of events) yield { notifications: n };
        },
      },
    },
  },
});

const yoga = createYoga({ schema, graphiql: { subscriptionsProtocol: 'SSE' } });
createServer(yoga).listen(4000, () => {
  console.log('Yoga SSE Server ready at http://localhost:4000/graphql');
});

客户端连接 SSE 端点:

const source = new EventSource(
  'http://localhost:4000/graphql?query=subscription{notifications(userId:"user-1"){id,title,body}}',
  { withCredentials: true }
);

source.addEventListener('next', (e) => {
  const { data } = JSON.parse(e.data);
  console.log('收到通知:', data.notifications);
});

Yoga 的 SSE 实现默认支持 multipart/mixedtext/event-stream,通过 Accept 请求头自动协商。


一句话总结

需要双向高频通信选 WebSocket(graphql-ws),纯服务端推送优先 SSE,极端兼容兜底用 HTTP 长轮询;多节点务必引入 Redis/Kafka 做事件广播,并配置 Sticky Session 保证 WebSocket 连接粘性。


FAQ

Q1: 同一个浏览器可以并发多少个 WebSocket/SSE 连接?

Chrome 对同一域名默认限制约 255 个 WebSocket 连接,SSE 在 HTTP/2 下共享单一 TCP 连接。若需订阅大量主题,建议在应用层做订阅聚合(单连接多订阅)。

Q2: WebSocket 断线后如何恢复遗漏的消息?

graphql-ws 不内置恢复机制,需业务层实现:分配单调递增序列号,客户端重连时上报 Last-Received-Event-ID,服务端从 Redis Stream 等缓存补发。

Q3: SSE 是否支持多路订阅复用一条连接?

原生 SSE 一条连接对应一个订阅。如需多路订阅,可合并多个字段到同一 Subscription 文档,或使用 HTTP/2 multiplexing。

Q4: 移动端(React Native / Flutter)如何选择?

React Native 与 Flutter 都对 WebSocket 支持良好,graphql-wsgraphql_flutter 可直接使用。仅在 WebView 环境(如小程序)中才需评估 SSE 或长轮询。

Q5: 生产环境 WebSocket 连接数规划?

每个连接在 Node.js 中通常占用 1-3MB 内存。10 万并发单节点不现实,必须采用网关+多节点架构,或使用 Go/Rust 高性能网关做接入层。


相关阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「GraphQL」更多文章

  1. gRPC-Web 与 GraphQL 混合架构:微服务通信分层实战
  2. GraphQL 服务端实战:Apollo Server、GraphQL Yoga 与 Pothos 选型
  3. GraphQL 客户端状态管理:Apollo Client、Relay 与 urql 深度对比