Spring GraphQL Subscription动态发布数据客户端无法接收问题求助
Spring GraphQL Subscription 客户端收不到推送问题修复
问题场景
- 基于Spring GraphQL Subscription + Reactor Sinks实现推送逻辑:NATS消息代理收到消息时,触发JPA查询
ToolOrder表全量数据,推送给所有在线订阅客户端 - 异常表现:NATS消息接收、数据库查询逻辑执行正常,Sinks发射日志正常打印,但客户端始终收不到服务端推送
根因分析
代码存在两个核心问题:
- Sinks实例不一致:在配置类、Resolver、NatsService中直接调用
valueSink()方法获取实例,如果配置类没有开启CGLIB代理(或直接手动实例化配置类),每次调用方法都会生成新的Sinks对象,导致NatsService发射数据的Sink和Resolver订阅Flux绑定的Sink完全不是同一个实例,数据自然无法流转到客户端。 - 缺少必要的Spring注解与配置:
- SubscriptionResolver方法没有加
@SubscriptionMapping注解,Spring GraphQL无法正确识别订阅映射 - 使用的
Sinks.many().multicast().onBackpressureBuffer()配置默认在无订阅者时直接丢弃发射的数据,如果先触发NATS消息、后发起客户端订阅,会直接发射失败返回FAIL_ZERO_SUBSCRIBER
- SubscriptionResolver方法没有加
修复后代码
1. 订阅配置类
@Configuration(proxyBeanMethods = true) // 开启配置类代理,保证@Bean方法返回单例 public class SubscribeConfiguration { @Bean public Sinks.Many<List<ToolOrder>> valueSink() { // 普通多播模式:仅推送给当前在线订阅者,无订阅者时丢弃数据 return Sinks.many().multicast().onBackpressureBuffer(); // 如果需要新订阅者连接后立刻拿到最新一次推送数据,替换为下面的重放模式 // return Sinks.many().replay().latest(); } @Bean public Flux<List<ToolOrder>> valueFlux(Sinks.Many<List<ToolOrder>> valueSink) { // 通过方法参数注入Sinks实例,避免直接调用方法产生多实例问题 return valueSink.asFlux(); } }
2. Subscription 解析器
@Controller // 必须加@Controller注解,让Spring扫描到这个GraphQL组件 public class ToolOrderSubscriptionResolver { private final Flux<List<ToolOrder>> toolOrderFlux; // 构造器注入提前初始化好的Flux单例,不要直接调用配置类方法 public ToolOrderSubscriptionResolver(Flux<List<ToolOrder>> toolOrderFlux) { this.toolOrderFlux = toolOrderFlux; } @SubscriptionMapping // 必须加这个注解,映射GraphQL订阅字段 public Publisher<List<ToolOrder>> getAllToolOrder() { System.out.println("客户端订阅连接建立成功"); return toolOrderFlux; // 移除无意义的map(e->e)操作 } }
3. NATS消息处理类
@Service public class NatsService { private static final Logger log = LoggerFactory.getLogger(NatsService.class); private final Sinks.Many<List<ToolOrder>> valueSink; private final ToolOrderRepository toolOrderRepository; // 构造器注入Sinks和Repository单例 public NatsService(Sinks.Many<List<ToolOrder>> valueSink, ToolOrderRepository toolOrderRepository) { this.valueSink = valueSink; this.toolOrderRepository = toolOrderRepository; } public void handleData(Message info) { String message = new String(info.getData()); log.info("received message from nats: {}, topic: {}", message, info.getSubject()); try { // JPA的findAll是阻塞方法,调度到专用阻塞线程池执行,避免占用NATS核心监听线程 List<ToolOrder> toolOrders = Mono.fromCallable(toolOrderRepository::findAll) .subscribeOn(Schedulers.boundedElastic) .block(); Sinks.EmitResult emitResult = valueSink.tryEmitNext(toolOrders); log.info("数据发射结果:{}", emitResult); if (emitResult.isFailure()) { log.warn("数据推送失败,失败原因:{}", emitResult); } } catch (JsonProcessingException e) { throw new RuntimeException("消息JSON解析失败", e); } } }
验证注意事项
- 确保GraphQL schema中正确定义了订阅字段,参考示例:
type Subscription { getAllToolOrder: [ToolOrder]! }
- 测试时先启动客户端发起订阅,看到控制台打印「客户端订阅连接建立成功」日志后,再触发NATS发送消息
- 如果日志打印发射结果为
FAIL_ZERO_SUBSCRIBER,说明订阅链路没有打通,优先检查注解、schema映射配置是否正确 - 如果需要新订阅者连上就拿到最新数据,把Sinks配置换成
replay().latest()模式即可
内容的提问来源于stack exchange,提问作者Shido
相关产品推荐
相关产品推荐

