使用Express-GraphQL与graphql-subscriptions时Subscription返回null问题排查
问题:GraphQL Subscription 返回始终为 null,无法正常工作
服务端代码
import express, { Express } from "express"; import { graphqlHTTP } from "express-graphql"; import { buildSchema } from "type-graphql"; import { TaskResolver } from "./resolvers/task.resolver"; import { pgDatasource } from "./configs/db.config"; import { SeatBandingResolver } from "./resolvers/seatBanding.resolver"; import { GuestChatResolver } from "./resolvers/guestChat.resolver"; import { RateResolver } from "./resolvers/rate.resolver"; import { YearResolver } from "./resolvers/year.resolver"; import { ImplementationRateResolver } from "./resolvers/implementationRate.resolver"; import { UserResolver } from "./resolvers/user.resolver"; import { ReportResolver } from "./resolvers/report.resolver"; // Subscriptions const ws = require("ws"); const { useServer } = require("graphql-ws/lib/use/ws"); const { execute, subscribe } = require("graphql"); const main = async () => { const app: Express = express(); try { //connect to db await pgDatasource.initialize(); } catch (err) { throw err; } //build gql schema let schema = await buildSchema({ resolvers: [ SeatBandingResolver, GuestChatResolver, RateResolver, YearResolver, ImplementationRateResolver, UserResolver, ], validate: false, // pubSub: new PubSub() }); let schemaDoc = await buildSchema({ resolvers: [ReportResolver], validate: false, }); //ql schema for report const docServer = graphqlHTTP((req, res) => { return { schema: schemaDoc, graphiql: true, context: { req: req, header: req.headers, }, }; }); //setting a graphql server instance const graphqServer = graphqlHTTP((req, res, graphQLParams) => { return { schema, context: { req: req, header: req.headers, }, graphiql: true, }; }); app.use(cors()); //graphql endpoint : change it to backend app.use("/graphql", graphqServer); //for report : change name to google api app.use("/doc", docServer); //test route app.get("/", (req, res) => { res.json({ message: "Hello world", }); }); let server = app.listen(3001, () => { console.log("server started"); const wsServer = new ws.WebSocketServer({ host: "localhost", // server, path: "/graphql", port: 3001, }); useServer( { schema, execute, subscribe, onConnect: (ctx) => { console.log("Connect"); }, onSubscribe: (ctx, msg) => { console.log("Subscribe"); }, onNext: (ctx, msg, args, result) => { console.debug("Next"); }, onError: (ctx, msg, errors) => { console.error("Error"); }, onComplete: (ctx, msg) => { console.log("Complete"); }, }, wsServer ); }); }; //starting a server main() .then(async (_) => { // await addColumn() }) .catch((err) => { console.log(err); });
客户端(Resolver)代码
import { Year } from "../entities/year.entity"; import { NewYear } from "../inputs/addYear.input"; import { Arg, Ctx, Field, Int, Mutation, ObjectType, Query, Resolver, Root, Subscription, UseMiddleware, } from "type-graphql"; import { Request } from "express"; import { Response } from "../helpers/response.helper"; import { Pagination } from "../inputs/pagination.input"; import { isAuth } from "../helpers/auth.helper"; import { PubSub, PubSubEngine } from "graphql-subscriptions"; const pubSub = new PubSub(); @ObjectType() class MessagePayload { @Field() message: string; } @Resolver(() => Year) export class YearResolver { @Mutation(() => String) async sendMessage(@Arg("message") message: string): Promise<string> { console.log("in send subscription"); pubSub.publish("MESSAGE_NOTIFICATION", { message }); return message; } //calling the subscription @Subscription(() => MessagePayload || null, { topics: "MESSAGE_NOTIFICATION", }) async receiveMessage( @Root() root: MessagePayload ): Promise<MessagePayload | null> { console.log("in publisher"); console.log(root, "in recieve message"); pubSub.asyncIterator("MESSAGE_NOTIFICATION"); return { message: "hello from the subscription" }; } }
问题诊断与修复方案
核心问题点
- PubSub 实例不统一
- 服务端构建Schema时注释了
pubSub: new PubSub(),导致Resolver内创建的PubSub与Schema使用的实例完全独立,消息无法跨实例传递。
- 服务端构建Schema时注释了
- WebSocket 配置冲突
- 手动指定
port:3001会与已启动的HTTP服务端口冲突,应该复用已创建的HTTP server实例,避免重复占用端口。
- 手动指定
- Subscription 方法逻辑错误
- 手动调用
pubSub.asyncIterator属于冗余操作,TypeGraphQL会通过topics配置自动处理订阅源;同时当前方法忽略了发布的实际消息,返回固定值,且类型定义不符合规范。
- 手动调用
修复步骤
服务端修复
- 全局创建并注入PubSub实例:
// 新增全局PubSub导入与实例化 import { PubSub } from "graphql-subscriptions"; export const pubSub = new PubSub(); // 构建Schema时注入pubSub let schema = await buildSchema({ resolvers: [ SeatBandingResolver, GuestChatResolver, RateResolver, YearResolver, ImplementationRateResolver, UserResolver, ], validate: false, pubSub: pubSub, // 取消注释并注入全局实例 });
- 修正WebSocket Server配置:
const wsServer = new ws.WebSocketServer({ server, // 复用已启动的HTTP server,删除host和port配置 path: "/graphql", });
客户端Resolver修复
- 导入服务端全局PubSub实例(避免重复创建),修改Subscription方法:
// 替换为服务端PubSub的实际路径 import { pubSub } from "../configs/pubSub"; // 修改Subscription装饰器与方法逻辑 @Subscription(() => MessagePayload, { topics: "MESSAGE_NOTIFICATION", }) receiveMessage(@Root() root: { message: string }): MessagePayload { console.log("收到订阅消息:", root.message); // 返回发布的实际消息,而非固定值 return { message: root.message }; }
- 移除冗余的
pubSub.asyncIterator调用,修正返回类型(使用MessagePayload | null替代MessagePayload || null,符合TypeScript语法)。
额外检查项
- 确保客户端GraphQL客户端(如Apollo Client)配置了正确的WebSocket端点:
ws://localhost:3001/graphql - 查看服务端日志,确认WebSocket连接、订阅事件正常触发(出现"Connect"、"Subscribe"日志)
- 调用
sendMessagemutation后,检查服务端是否有消息发布的相关日志
内容的提问来源于stack exchange,提问作者Aqdas
相关产品推荐
相关产品推荐

