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

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方法被多次调用。结合代码场景,核心原因是:

  1. Reactor的Mono是冷流,默认每次订阅都会重新执行整个请求链路;
  2. 高并发场景下,缓存逻辑或Akka Stream的重复处理可能导致同一个Mono实例被意外重复订阅;
  3. ConcurrentHashMap.computeIfAbsent的原子性虽能保证同一key只执行一次fetch,但如果fetch返回的Future对应的Mono未被正确缓存,仍可能触发重复订阅。

调试步骤

  1. 追踪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;
    }
    
  2. 验证并发场景
    用压测工具模拟高并发请求同一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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 14:04:54