You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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" };
  }
}

问题诊断与修复方案

核心问题点

  1. PubSub 实例不统一
    • 服务端构建Schema时注释了pubSub: new PubSub(),导致Resolver内创建的PubSub与Schema使用的实例完全独立,消息无法跨实例传递。
  2. WebSocket 配置冲突
    • 手动指定port:3001会与已启动的HTTP服务端口冲突,应该复用已创建的HTTP server实例,避免重复占用端口。
  3. 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"日志)
  • 调用sendMessage mutation后,检查服务端是否有消息发布的相关日志

内容的提问来源于stack exchange,提问作者Aqdas

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.20 06:45:36