Spring WebFlux WebClient重复订阅异常排查与修复咨询
Spring WebFlux WebClient 重复订阅异常调试与修复
问题背景
触发的异常栈:
java.lang.IllegalStateException: Spec. Rule 2.12 - Subscriber.onSubscribe MUST NOT be called more than once (based on object equality) at reactor.core.Exceptions.duplicateOnSubscribeException(Exceptions.java:181) at reactor.core.publisher.Operators.reportSubscriptionSet(Operators.java:1083) at reactor.core.publisher.Operators.setOnce(Operators.java:1188) at reactor.core.publisher.MonoFlatMap.onSubscribe(MonoFlatMap.java:102) at reactor.core.publisher.Operators.reportThrowInSubscribe(Operators.java:225) at reactor.core.publisher.Mono.subscribe(Mono.java:4255) at reactor.core.publisher.Mono.subscribeWith(Mono.java:4363) at reactor.core.publisher.Mono.toFuture(Mono.java:4697) at abc.bcd.cde.fgh.SettingsRepositoryImplP2.getApplicationSettings(SettingsRepositoryImplP2.java:32) at java.util.concurrent.ConcurrentHashMap.computeIfAbsent(ConcurrentHashMap.java:1705) at abc.bcd.cde.fgh.ApplicationSettingsCache.hasMaxCapture(ApplicationSettingsCache.java:38) at abc.bcd.cde.fgh.ApplicationDebugService.shouldBeCaptured(ApplicationDebugService.java:20) at abc.bcd.cde.fgh.ApplicationDebugService.lambdaf865c(ApplicationDebugService.java:26) at akka.stream.javadsl.Flow.(Flow.scala:739) at akka.stream.impl.fusing.MapAsync96775anon.onPush(Ops.scala:1307) at akka.stream.impl.fusing.GraphInterpreter.processPush(GraphInterpreter.scala:542) at akka.stream.impl.fusing.GraphInterpreter.execute(GraphInterpreter.scala:423) at akka.stream.impl.fusing.GraphInterpreterShell.runBatch(ActorGraphInterpreter.scala:650) at akka.stream.impl.fusing.GraphInterpreterShell.execute(ActorGraphInterpreter.scala:521) at akka.stream.impl.fusing.GraphInterpreterShell.processEvent(ActorGraphInterpreter.scala:625) at akka.stream.impl.fusing.ActorGraphInterpreter.akka96775processEvent(ActorGraphInterpreter.scala:800) at akka.stream.impl.fusing.ActorGraphInterpreter.akka96775shortCircuitBatch(ActorGraphInterpreter.scala:787) at akka.stream.impl.fusing.ActorGraphInterpreter96775anonfun.applyOrElse(ActorGraphInterpreter.scala:819) at akka.actor.Actor.aroundReceive(Actor.scala:537) at akka.actor.Actor.aroundReceive at akka.stream.impl.fusing.ActorGraphInterpreter.aroundReceive(ActorGraphInterpreter.scala:716) at akka.actor.ActorCell.receiveMessage96775(ActorCell.scala:580) at akka.actor.ActorCell.receiveMessage(ActorCell.scala) at akka.actor.ActorCell.invoke(ActorCell.scala:548) at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:270) at akka.dispatch.Mailbox.run(Mailbox.scala:231) at akka.dispatch.Mailbox.exec(Mailbox.scala:243) at java.util.concurrent.ForkJoinTask.doExec96775(ForkJoinTask.java:290) at java.util.concurrent.ForkJoinTask.doExec(ForkJoinTask.java) at java.util.concurrent.ForkJoinPool.topLevelExec(ForkJoinPool.java:1020) at java.util.concurrent.ForkJoinPool.scan(ForkJoinPool.java:1656) at java.util.concurrent.ForkJoinPool.runWorker(ForkJoinPool.java:1594) at java.util.concurrent.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:183)
相关代码:
WebClient请求逻辑:
public Mono<TestSettings> retrieveApplicationConfiguration(String applicationId) { var url = applicationConfigurationServiceUrl.concat("/conf").concat("/{id}"); return webClient.get() .uri(url, applicationId) .retrieve() .onStatus(HttpStatus::is4xxClientError, (clientResponse -> { log.info("Status code : {}", clientResponse.statusCode().value()); if (clientResponse.statusCode().equals(HttpStatus.NOT_FOUND)) { return Mono.error(new ApplicationConfigServiceException("Error fetching Applications for application " + applicationId, clientResponse.statusCode().value())); } return clientResponse.bodyToMono(String.class) .flatMap(response -> Mono.error(new ApplicationConfigServiceException(" Error fetching Applications for " + applicationId + response, clientResponse.statusCode().value()))); })) .onStatus(HttpStatus::is5xxServerError, (clientResponse -> { log.info("Status code : {}", clientResponse.statusCode().value()); return clientResponse.bodyToMono(String.class) .flatMap(response -> Mono.error(new ApplicationConfigServiceException(response, clientResponse.statusCode().value()))); })) .bodyToMono(ApplicationConfigurationResponse.class) .map(a -> { log.info("LOG : {}", a); return TestSettings.create(a); }) .log(); }
Mono转Future逻辑:
@Override public CompletableFuture<TestSettings> getApplicationSettings(String application) { return webClient.retrieveApplicationConfiguration(application) .toFuture(); }
上游缓存调用逻辑:
return testSettingsMap .computeIfAbsent(app, testSettings::fetch) .thenApply(TestSettings::debugSettings) .thenApply(s -> s.equals(MAXIMUM));
其中testSettingsMap定义为:private final ConcurrentHashMap<String, CompletableFuture<TestSettings>> testSettingsMap;
问题根源
异常提示违反Reactor规范:同一个Subscriber实例的onSubscribe方法被多次调用。结合代码场景,核心原因是:
- Reactor的
Mono是冷流,默认每次订阅都会重新执行整个请求链路; - 高并发场景下,缓存逻辑或Akka Stream的重复处理可能导致同一个
Mono实例被意外重复订阅; ConcurrentHashMap.computeIfAbsent的原子性虽能保证同一key只执行一次fetch,但如果fetch返回的Future对应的Mono未被正确缓存,仍可能触发重复订阅。
调试步骤
- 追踪Mono实例与订阅次数
在retrieveApplicationConfiguration中添加日志,确认同一applicationId是否复用同一Mono,以及订阅次数:public Mono<TestSettings> retrieveApplicationConfiguration(String applicationId) { var url = applicationConfigurationServiceUrl.concat("/conf").concat("/{id}"); AtomicInteger subscribeCount = new AtomicInteger(0); Mono<TestSettings> mono = webClient.get() // 原有逻辑不变 .doOnSubscribe(sub -> { int count = subscribeCount.incrementAndGet(); log.info("Mono[{}] for app {} subscribed {} times", mono.hashCode(), applicationId, count); if (count > 1) { log.warn("Mono[{}] for app {} subscribed multiple times!", mono.hashCode(), applicationId); } }) .log(); log.info("Generated Mono[{}] for app {}", mono.hashCode(), applicationId); return mono; } - 验证并发场景
用压测工具模拟高并发请求同一applicationId,观察日志中是否出现重复订阅的警告。
修复方案
方案1:将冷流转热流(推荐)
在retrieveApplicationConfiguration返回的Mono上添加.cache()操作符,确保多次订阅仅执行一次请求,同时避免重复订阅问题:
public Mono<TestSettings> retrieveApplicationConfiguration(String applicationId) { var url = applicationConfigurationServiceUrl.concat("/conf").concat("/{id}"); return webClient.get() .uri(url, applicationId) .retrieve() .onStatus(HttpStatus::is4xxClientError, (clientResponse -> { log.info("Status code : {}", clientResponse.statusCode().value()); if (clientResponse.statusCode().equals(HttpStatus.NOT_FOUND)) { return Mono.error(new ApplicationConfigServiceException("Error fetching Applications for application " + applicationId, clientResponse.statusCode().value())); } return clientResponse.bodyToMono(String.class) .flatMap(response -> Mono.error(new ApplicationConfigServiceException(" Error fetching Applications for " + applicationId + response, clientResponse.statusCode().value()))); })) .onStatus(HttpStatus::is5xxServerError, (clientResponse -> { log.info("Status code : {}", clientResponse.statusCode().value()); return clientResponse.bodyToMono(String.class) .flatMap(response -> Mono.error(new ApplicationConfigServiceException(response, clientResponse.statusCode().value()))); })) .bodyToMono(ApplicationConfigurationResponse.class) .map(a -> { log.info("LOG : {}", a); return TestSettings.create(a); }) .log() .cache(); // 转为热流,缓存请求结果 }
方案2:优化缓存逻辑(响应式最佳实践)
将缓存对象从CompletableFuture改为Mono,更贴合Reactor响应式设计,彻底避免重复订阅:
// 修改缓存定义 private final ConcurrentHashMap<String, Mono<TestSettings>> testSettingsMap; // 上游调用逻辑调整 return testSettingsMap .computeIfAbsent(app, this::fetchMono) .map(TestSettings::debugSettings) .map(s -> s.equals(MAXIMUM)) .toFuture(); // 新增fetchMono方法 private Mono<TestSettings> fetchMono(String application) { return webClient.retrieveApplicationConfiguration(application) .cache(); // 缓存请求结果 }
内容的提问来源于stack exchange,提问作者SRJ
相关产品推荐
相关产品推荐

