参考graphql-java订阅示例,如何在Spring Boot中实现WebSocket订阅?
我刚好折腾过这个场景,在Spring环境下复用graphql-java的WebSocket订阅能力其实没那么复杂,咱们一步步拆解来做:
1. 先搞定依赖配置
首先得把Spring WebSocket、graphql-java核心以及订阅相关的依赖加进项目:
Maven(pom.xml)
<dependencies> <!-- Spring WebSocket 核心 --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-websocket</artifactId> </dependency> <!-- graphql-java 核心库 --> <dependency> <groupId>com.graphql-java</groupId> <artifactId>graphql-java</artifactId> <version>19.0</version> <!-- 建议用最新稳定版 --> </dependency> <!-- graphql-java 订阅扩展 --> <dependency> <groupId>com.graphql-java</groupId> <artifactId>graphql-java-subscriptions</artifactId> <version>19.0</version> </dependency> <!-- 可选:如果用Spring管理GraphQL实例的starter --> <dependency> <groupId>com.graphql-java</groupId> <artifactId>graphql-spring-boot-starter</artifactId> <version>12.0</version> </dependency> </dependencies>
Gradle(build.gradle)
dependencies { implementation 'org.springframework.boot:spring-boot-starter-websocket' implementation 'com.graphql-java:graphql-java:19.0' implementation 'com.graphql-java:graphql-java-subscriptions:19.0' implementation 'com.graphql-java:graphql-spring-boot-starter:12.0' }
2. 配置WebSocket端点与GraphQL订阅处理器
Spring里需要注册graphql-java提供的WebSocketSubscriptionHandler,让它处理指定WebSocket端点的订阅请求:
import org.springframework.context.annotation.Configuration; import org.springframework.web.socket.config.annotation.EnableWebSocket; import org.springframework.web.socket.config.annotation.WebSocketConfigurer; import org.springframework.web.socket.config.annotation.WebSocketHandlerRegistry; import com.graphql.java.GraphQL; import com.graphql.java.subscriptions.WebSocketSubscriptionHandler; @Configuration @EnableWebSocket public class GraphQLWebSocketConfig implements WebSocketConfigurer { private final WebSocketSubscriptionHandler subscriptionHandler; // 注入Spring管理的GraphQL实例 public GraphQLWebSocketConfig(GraphQL graphQL) { this.subscriptionHandler = new WebSocketSubscriptionHandler(graphQL); } @Override public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) { // 注册订阅端点,比如 /graphql/subscriptions registry.addHandler(subscriptionHandler, "/graphql/subscriptions") .setAllowedOrigins("*"); // 生产环境一定要替换成具体允许的域名,别用* } }
3. 编写你的订阅Resolver(返回Publisher)
这部分和你之前的非Spring版本逻辑一致,Resolver方法返回一个符合Reactive Streams规范的Publisher就行,比如用Spring生态友好的Reactor Flux:
import org.reactivestreams.Publisher; import reactor.core.publisher.Flux; import java.time.Duration; import org.springframework.stereotype.Component; import com.graphql.java.annotations.GraphQLSubscription; @Component public class MessageSubscriptionResolver { // 订阅字段:实时推送新消息 @GraphQLSubscription public Publisher<Message> newMessage() { // 这里模拟每1秒生成一条新消息,实际场景可以替换成Kafka、RabbitMQ的消费流 return Flux.interval(Duration.ofSeconds(1)) .map(tick -> new Message("实时消息 #" + tick)); } // 消息实体类,对应GraphQL Schema中的Message类型 public static class Message { private String content; public Message(String content) { this.content = content; } public String getContent() { return content; } } }
4. 构建GraphQL实例并关联订阅Resolver
需要把订阅Resolver绑定到GraphQL Schema中,让GraphQL能识别你的订阅字段:
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import com.graphql.java.GraphQL; import com.graphql.java.GraphQLSchema; import com.graphql.java.SchemaParser; import com.graphql.java.TypeDefinitionRegistry; import com.graphql.java.schema.GraphQLObjectType; @Configuration public class GraphQLConfig { private final MessageSubscriptionResolver subscriptionResolver; public GraphQLConfig(MessageSubscriptionResolver subscriptionResolver) { this.subscriptionResolver = subscriptionResolver; } @Bean public GraphQL graphQL() { // 定义GraphQL SDL Schema String schemaSDL = """ type Query { hello: String # 随便加个查询字段,保证Schema合法 } type Subscription { newMessage: Message # 订阅字段,对应Resolver的方法 } type Message { content: String } """; // 解析SDL生成类型注册表 TypeDefinitionRegistry typeRegistry = new SchemaParser().parse(schemaSDL); // 构建订阅类型,绑定Resolver的dataFetcher GraphQLObjectType subscriptionType = GraphQLObjectType.newObject() .name("Subscription") .field(field -> field.name("newMessage") .type(new com.graphql.java.GraphQLTypeReference("Message")) .dataFetcher(env -> subscriptionResolver.newMessage())) .build(); // 组装完整Schema并创建GraphQL实例 GraphQLSchema schema = GraphQLSchema.newSchema() .subscription(subscriptionType) .build(); return GraphQL.newGraphQL(schema).build(); } }
5. 测试订阅功能
现在启动Spring Boot应用,用WebSocket客户端(比如wscat、Postman的WebSocket功能)连接ws://localhost:8080/graphql/subscriptions,然后发送以下订阅请求:
{ "type": "start", "id": "sub-1", "payload": { "query": "subscription { newMessage { content } }" } }
之后你就能每秒收到一条推送的消息了!
关键逻辑说明
你之前疑惑的「为什么返回Publisher就能和WebSocket结合」——其实graphql-java-subscriptions里的WebSocketSubscriptionHandler已经帮你做了所有中间工作:它会解析客户端的订阅请求,订阅你Resolver返回的Publisher,每当Publisher有新数据产生时,就自动把数据封装成GraphQL订阅响应,通过WebSocket推送给客户端。
内容的提问来源于stack exchange,提问作者user101010101

