Spring Cloud Stream中WebClient报错致Dispatcher无订阅者的解决方法咨询
问题描述
我们使用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

