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

第三方GraphQL客户端对接Quarkus SmallRye GraphQL订阅失效问题

问题背景

我是Quarkus + GraphQL技术栈初学者,出于学习目的搭建了一套GraphQL服务,代码仓库地址:quakus-graphql-demo。
服务核心Java代码如下:

import java.util.Collection;

import javax.inject.Inject;

import org.eclipse.microprofile.graphql.GraphQLApi;
import org.eclipse.microprofile.graphql.Mutation;
import org.eclipse.microprofile.graphql.Query;

import com.wangxiaohu.quarkus.graphql.demo.model.Person;
import com.wangxiaohu.quarkus.graphql.demo.service.PersonService;

import io.quarkus.logging.Log;
import io.smallrye.graphql.api.Subscription;
import io.smallrye.mutiny.Multi;
import io.smallrye.mutiny.operators.multi.processors.BroadcastProcessor;

@GraphQLApi
public class PersonResource {

    private final BroadcastProcessor<Person> _personBroadcastProcessor;

    public PersonResource() {
        _personBroadcastProcessor = BroadcastProcessor.create();
    }

    @Inject
    PersonService _personService;

    @Query("getAllPeople")
    public Collection<Person> getAllPeople() {
        return _personService.getAllPeople();
    }

    @Query("getPersonById")
    public Person getPerson(int id) {
        return _personService.getPerson(id);
    }

    @Mutation("createPerson")
    public Person createPerson(String firstName, String lastName) {
        Person person = _personService.createPerson(firstName, lastName);

        Log.info("signaling the person created...");
        _personBroadcastProcessor.onNext(person);
        Log.info("signaled the person created.");

        return person;
    }

    @Subscription("personCreated")
    public Multi<Person> subscribeToPersonCreation() {
        Log.info("subscribeToPersonCreation");
        return _personBroadcastProcessor;
    }
}

该服务设计支持三类操作:

  • 新增Person
  • 查询全量Person列表
  • 订阅Person创建事件

测试阶段首先编写Python客户端尝试订阅Person创建事件,代码路径:test/python,核心代码如下:

from gql import gql, Client
from gql.transport.websockets import WebsocketsTransport

if __name__ == '__main__':
    transport = WebsocketsTransport(
        url="ws://localhost:8080/graphql",
        subprotocols=[WebsocketsTransport.GRAPHQLWS_SUBPROTOCOL]
    )

    client = Client(transport=transport, fetch_schema_from_transport=True)

    query = gql(
        '''
        subscription subscribeToPersonCreation {
            personCreated{
                id
                firstName
                lastName
            }
        }
        '''
    )

    for result in client.subscribe(query):
        print(result)

测试时发现新增Person后订阅完全不触发,在_personBroadcastProcessor.onNext(person)位置加断点调试,看到执行到这行时subscribers列表为空。
调试时subscribers列表为空截图

后续又编写Node.js版本GraphQL客户端做订阅测试,新增Person时依然收不到任何推送,核心代码如下:

const ws = require('ws');
const Crypto = require('crypto');
const { createClient } = require('graphql-ws');

const client = createClient({
  url: "ws://localhost:8080/graphql",
  webSocketImpl: ws,
  generateID: () =>
    ([1e7] + -1e3 + -4e3 + -8e3 + -1e11).replace(/[018]/g, (c) =>
      (c ^ (Crypto.randomBytes(1)[0] & (15 >> (c / 4)))).toString(16),
    ),
});

(async () => {
  const onNext = (value) => {
    console.log(value);
  }

  await new Promise((resolve, reject) => {
    unsubscribe = client.subscribe(
      {
        query: `subscription subscribeToPersonCreation {
          personCreated {
            id
            firstName
            lastName
          }
        }`
      },
      {
        next: onNext,
        error: reject,
        complete: resolve,
      }
    );
  });
})();

对应测试代码路径:test/nodejs。


问题根因

核心问题是CDI Bean作用域配置错误,导致订阅和发布操作使用的BroadcastProcessor不是同一个实例:
@GraphQLApi标注的资源类默认CDI作用域为@RequestScoped,每次HTTP/WebSocket请求进入时都会创建全新的类实例:

  1. 客户端发起订阅连接时,创建第一个PersonResource实例,订阅关系绑定在这个实例持有的_personBroadcastProcessor上
  2. 后续发起createPerson的Mutation请求时,会创建第二个全新的PersonResource实例,调用onNext发布事件用的是这个新实例自己的_personBroadcastProcessor,这个处理器上没有任何订阅者,事件自然不会推送给客户端,断点看到的空subscribers列表就是这个新实例的处理器状态。

另外BroadcastProcessor本身的特性就是无订阅者时调用onNext会直接丢弃事件,不会缓存,进一步导致消息丢失。


修复方案

服务端修改

将事件广播处理器抽成应用全局单例的CDI Bean,保证所有订阅、发布操作共用同一个处理器实例:

  1. 新建全局事件发布类,标注@ApplicationScoped设为应用单例:
import io.smallrye.mutiny.Multi;
import io.smallrye.mutiny.operators.multi.processors.BroadcastProcessor;
import javax.enterprise.context.ApplicationScoped;

@ApplicationScoped
public class PersonEventPublisher {
    private final BroadcastProcessor<Person> broadcastProcessor = BroadcastProcessor.create();

    public void publish(Person person) {
        broadcastProcessor.onNext(person);
    }

    public Multi<Person> subscribe() {
        return broadcastProcessor;
    }
}
  1. 修改原PersonResource类,移除自己实例化的BroadcastProcessor,改为注入单例的发布器:
import java.util.Collection;

import javax.inject.Inject;

import org.eclipse.microprofile.graphql.GraphQLApi;
import org.eclipse.microprofile.graphql.Mutation;
import org.eclipse.microprofile.graphql.Query;

import com.wangxiaohu.quarkus.graphql.demo.model.Person;
import com.wangxiaohu.quarkus.graphql.demo.service.PersonService;
import com.wangxiaohu.quarkus.graphql.demo.event.PersonEventPublisher;

import io.quarkus.logging.Log;
import io.smallrye.graphql.api.Subscription;
import io.smallrye.mutiny.Multi;

@GraphQLApi
public class PersonResource {

    @Inject
    PersonService personService;

    @Inject
    PersonEventPublisher personEventPublisher;

    @Query("getAllPeople")
    public Collection<Person> getAllPeople() {
        return personService.getAllPeople();
    }

    @Query("getPersonById")
    public Person getPerson(int id) {
        return personService.getPerson(id);
    }

    @Mutation("createPerson")
    public Person createPerson(String firstName, String lastName) {
        Person person = personService.createPerson(firstName, lastName);
        Log.info("signaling the person created...");
        personEventPublisher.publish(person);
        Log.info("signaled the person created.");
        return person;
    }

    @Subscription("personCreated")
    public Multi<Person> subscribeToPersonCreation() {
        Log.info("new person creation subscription registered");
        return personEventPublisher.subscribe();
    }
}

客户端对接说明

Quarkus SmallRye GraphQL的WebSocket订阅默认端点为/graphql,同时兼容两种主流GraphQL WebSocket子协议,现有客户端代码不需要额外修改即可正常对接:

  • 老版Apollo协议(对应Python客户端使用的graphql-ws子协议):SmallRye原生支持,保持现有Python代码配置即可
  • 新版graphql-ws库协议(对应Node.js客户端使用的graphql-transport-ws子协议):SmallRye同样原生支持,现有Node.js代码不需要调整子协议配置

重启服务后,先启动任意客户端建立订阅连接,再调用createPerson接口即可正常收到推送事件。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.31 10:30:53