第三方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列表为空。
后续又编写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请求进入时都会创建全新的类实例:
- 客户端发起订阅连接时,创建第一个
PersonResource实例,订阅关系绑定在这个实例持有的_personBroadcastProcessor上 - 后续发起
createPerson的Mutation请求时,会创建第二个全新的PersonResource实例,调用onNext发布事件用的是这个新实例自己的_personBroadcastProcessor,这个处理器上没有任何订阅者,事件自然不会推送给客户端,断点看到的空subscribers列表就是这个新实例的处理器状态。
另外BroadcastProcessor本身的特性就是无订阅者时调用onNext会直接丢弃事件,不会缓存,进一步导致消息丢失。
服务端修改
将事件广播处理器抽成应用全局单例的CDI Bean,保证所有订阅、发布操作共用同一个处理器实例:
- 新建全局事件发布类,标注
@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; } }
- 修改原
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

