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

Spring Cloud Stream中WebClient报错致Dispatcher无订阅者的解决方法咨询

Spring Cloud Stream响应式处理器因WebClient错误取消订阅问题解决

问题描述

我们使用Spring Cloud Stream响应式函数处理消息事件,通过Flux的flatMap调用WebClient发起外部HTTP请求。但只要WebClient请求出错,FluxMessageChannel的MessageHandler就会取消订阅通道,之后无法再发送任何事件,触发Dispatcher has no subscribers错误。尝试用onStatus处理但未成功,求避免处理器被取消订阅的方案。

原代码

...

private final WebClient profileServiceClient;

private Function<Message<MessageParticipant>, Mono<ProfileSchema>> getProfile = msg -> {
    final String profileId = msg.getPayload().getProfileId();
    return profileServiceClient.get()
            .uri(uriBuilder -> uriBuilder.path("/profiles/{id}").build(profileId))
            .retrieve()
            .bodyToMono(ProfileSchema.class);
};

@Bean
public Function<Flux<Message<MessageParticipant>>, Flux<Message<ProfileSchema>>> gatherProfile(ProfileService profileService) {
    return messages -> {
        return messages.log(log.getName(), Level.FINEST)
                .flatMap(getProfile)
                .map(p -> MessageBuilder.withPayload(p).build());
    };
}

错误信息

at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:76) ~[spring-integration-core-5.5.12.jar:5.5.12]
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:317) ~[spring-integration-core-5.5.12.jar:5.5.12]
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:272) ~[spring-integration-core-5.5.12.jar:5.5.12]
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:187) ~[spring-messaging-5.3.20.jar:5.3.20]
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:166) ~[spring-messaging-5.3.20.jar:5.3.20]
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:47) ~[spring-messaging-5.3.20.jar:5.3.20]
    at org.springframework.messaging.core.AbstractMessageSendingTemplate.send(AbstractMessageSendingTemplate.java:109) ~[spring-messaging-5.3.20.jar:5.3.20]
    at org.springframework.integration.endpoint.MessageProducerSupport.sendMessage(MessageProducerSupport.java:216) ~[spring-integration-core-5.5.12.jar:5.5.12]
    at de.idealo.spring.stream.binder.sqs.inbound.SqsInboundChannelAdapter.access$400(SqsInboundChannelAdapter.java:33) ~[spring-cloud-stream-binder-sqs-1.9.0.jar:na]
    at de.idealo.spring.stream.binder.sqs.inbound.SqsInboundChannelAdapter$IntegrationQueueMessageHandler.handleMessageInternal(SqsInboundChannelAdapter.java:162) ~[spring-cloud-stream-binder-sqs-1.9.0.jar:na]
    at org.springframework.messaging.handler.invocation.AbstractMethodMessageHandler.handleMessage(AbstractMethodMessageHandler.java:458) ~[spring-messaging-5.3.20.jar:5.3.20]
    at io.awspring.cloud.messaging.listener.SimpleMessageListenerContainer.executeMessage(SimpleMessageListenerContainer.java:222) ~[spring-cloud-aws-messaging-2.4.2.jar:2.4.2]
    at io.awspring.cloud.messaging.listener.SimpleMessageListenerContainer$MessageGroupExecutor.run(SimpleMessageListenerContainer.java:426) ~[spring-cloud-aws-messaging-2.4.2.jar:2.4.2]
    at io.awspring.cloud.messaging.listener.SimpleMessageListenerContainer$SignalExecutingRunnable.run(SimpleMessageListenerContainer.java:310) ~[spring-cloud-aws-messaging-2.4.2.jar:2.4.2]
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) ~[na:na]
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) ~[na:na]
    at java.base/java.lang.Thread.run(Thread.java:833) ~[na:na]
Caused by: org.springframework.integration.MessageDispatchingException: Dispatcher has no subscribers**
    at org.springframework.integration.dispatcher.UnicastingDispatcher.doDispatch(UnicastingDispatcher.java:139) ~[spring-integration-core-5.5.12.jar:5.5.12]
    at org.springframework.integration.dispatcher.UnicastingDispatcher.dispatch(UnicastingDispatcher.java:106) ~[spring-integration-core-5.5.12.jar:5.5.12]
    at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:72) ~[spring-integration-core-5.5.12.jar:5.5.12]
    ... 16 common frames omitted

尝试过的代码(未解决)

private Function<Message<MessageParticipant>, Mono<ProfileSchema>> getProfile = msg -> {
    final String profileId = msg.getPayload().getProfileId();
    return profileServiceClient.get()
            .uri(uriBuilder -> uriBuilder.path("/profiles/{id}").build(profileId))
            .retrieve()
            .onStatus(HttpStatus::isError, clientResponse -> {
                        log.error("Error while getting profile {} with status code {}", profileId, clientResponse.statusCode());
                        throw new ServiceException("Unable to retrieve profile for profile id : " + profileId);
                    })

            .bodyToMono(ProfileSchema.class);
};

问题根源

你在onStatus中直接抛出异常,这个异常会向上传递到整个Flux流中。根据响应式流的规则,一旦流中出现未处理的错误,整个流会立即终止并取消订阅,导致后续消息无法被处理,最终触发“无订阅者”错误。

解决办法

必须在单个Mono层面捕获并处理错误,不让错误扩散到整个Flux流。以下是几种可行方案:

1. 使用onErrorResume返回 fallback 值

如果业务允许,在WebClient请求出错时返回一个默认的ProfileSchema对象,或者根据错误类型返回不同的 fallback:

private Function<Message<MessageParticipant>, Mono<ProfileSchema>> getProfile = msg -> {
    final String profileId = msg.getPayload().getProfileId();
    return profileServiceClient.get()
            .uri(uriBuilder -> uriBuilder.path("/profiles/{id}").build(profileId))
            .retrieve()
            .bodyToMono(ProfileSchema.class)
            .onErrorResume(error -> {
                log.error("Failed to get profile {}: {}", profileId, error.getMessage());
                // 替换为你的默认对象或业务逻辑
                return Mono.just(new ProfileSchema()); 
            });
};

2. 使用onErrorContinue跳过错误消息

如果不需要处理错误消息,只想跳过它继续处理后续消息,可以在Flux链中添加onErrorContinue:

@Bean
public Function<Flux<Message<MessageParticipant>>, Flux<Message<ProfileSchema>>> gatherProfile(ProfileService profileService) {
    return messages -> messages.log(log.getName(), Level.FINEST)
            .flatMap(getProfile)
            .onErrorContinue((error, msg) -> {
                log.error("Error processing message {}: {}", msg, error.getMessage());
                // 记录错误后,流会继续处理下一条消息
            })
            .map(p -> MessageBuilder.withPayload(p).build());
}

注意:onErrorContinue会跳过出错的元素,但如果错误是系统性的(如外部服务持续不可用),会产生大量错误日志,建议结合重试机制使用。

3. 结合onStatus和onErrorResume优雅处理HTTP错误

先通过onStatus捕获HTTP错误状态码并转换为自定义异常,再用onErrorResume捕获并处理:

private Function<Message<MessageParticipant>, Mono<ProfileSchema>> getProfile = msg -> {
    final String profileId = msg.getPayload().getProfileId();
    return profileServiceClient.get()
            .uri(uriBuilder -> uriBuilder.path("/profiles/{id}").build(profileId))
            .retrieve()
            .onStatus(HttpStatus::isError, clientResponse -> {
                log.error("Profile service returned error {} for id {}", clientResponse.statusCode(), profileId);
                return Mono.error(new ServiceException("Profile service error: " + clientResponse.statusCode()));
            })
            .bodyToMono(ProfileSchema.class)
            .onErrorResume(ServiceException.class, error -> {
                log.error("Handling service exception: {}", error.getMessage());
                // 返回空Mono,flatMap会忽略该元素,流继续执行
                return Mono.empty(); 
            });
};

这种方式既能精确处理HTTP错误,又不会终止整个流。

4. 添加重试机制(可选)

如果外部服务的错误是临时的,可以添加重试逻辑减少错误发生概率:

private Function<Message<MessageParticipant>, Mono<ProfileSchema>> getProfile = msg -> {
    final String profileId = msg.getPayload().getProfileId();
    return profileServiceClient.get()
            .uri(uriBuilder -> uriBuilder.path("/profiles/{id}").build(profileId))
            .retrieve()
            .bodyToMono(ProfileSchema.class)
            .retryWhen(Retry.backoff(3, Duration.ofSeconds(1))
                    .filter(error -> error instanceof WebClientResponseException || error instanceof IOException)
                    .onRetryExhaustedThrow((retryBackoffSpec, retrySignal) -> {
                        log.error("Retry exhausted for profile {}", profileId);
                        return new ServiceException("Failed after retries");
                    }))
            .onErrorResume(error -> {
                log.error("Final error for profile {}: {}", profileId, error.getMessage());
                return Mono.empty();
            });
};

内容的提问来源于stack exchange,提问作者Christian Demers

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 06:14:56