WebFlux中如何在MessageSource实现类中访问Reactor Context?
解决WebFlux中MessageSource访问Reactor Context的问题
你的核心问题在于AbstractMessageSource的resolveCode是同步阻塞方法,无法直接访问Reactor Context;而用block()会创建新订阅,丢失外层Context。下面给出两种解决方案,优先推荐反应式原生方案。
方案一:使用ReactiveMessageSource(推荐)
Spring 5.3+提供了ReactiveMessageSource接口,专门适配反应式场景,返回Mono<String>,可以直接通过Mono.deferContextual访问Reactor Context。同时替换阻塞的RestTemplate为反应式的WebClient,完全契合WebFlux模型。
实现ReactiveMessageSource
@RequiredArgsConstructor @Component public class ReactiveRemoteMessageSource implements ReactiveMessageSource { private final WebClient webClient; @Override public Mono<String> resolveMessage(String code, Locale locale, MessageSourceResolvable resolvable) { // 直接返回原始文本(如果code包含空格) if (StringUtils.containsWhitespace(code)) { return Mono.just(code); } // 从Reactor Context中获取HTTP头参数 return Mono.deferContextual(contextView -> { // 替换为你实际存入Context的头参数key String targetHeaderValue = contextView.get("YOUR_HEADER_KEY"); return webClient.get() .uri(uriBuilder -> uriBuilder .path("/public/v1/messages/i18n/{code}") .queryParam("yourParam", targetHeaderValue) // 传递头参数到远程接口 .build(code)) .header(HttpHeaders.ACCEPT_LANGUAGE, locale.toLanguageTag()) .retrieve() .bodyToMono(String.class) .onErrorMap(HttpStatusCodeException.class, CustomErrorException::of); }); } }
确保WebFilter正确存入Context
@Component public class RequestHeaderContextFilter implements WebFilter { @Override public Mono<Void> filter(ServerWebExchange exchange, WebFilterChain chain) { // 从请求头中获取目标参数 String headerValue = exchange.getRequest().getHeaders().getFirst("YOUR_HEADER_KEY"); // 将参数存入Reactor Context return chain.filter(exchange) .contextWrite(context -> context.put("YOUR_HEADER_KEY", headerValue)); } }
方案二:ThreadLocal桥接(不推荐,仅兼容旧场景)
如果因为依赖旧代码必须使用AbstractMessageSource,可以通过ThreadLocal将Reactor Context的值桥接到同步方法中。但要注意WebFlux的线程复用特性,必须严格清理ThreadLocal,避免内存泄漏或线程污染。
实现带ThreadLocal的WebFilter
@Component public class RequestHeaderThreadLocalFilter implements WebFilter { private static final ThreadLocal<String> HEADER_VALUE_THREAD_LOCAL = new ThreadLocal<>(); // 提供静态方法供MessageSource获取值 public static String getCurrentHeaderValue() { return HEADER_VALUE_THREAD_LOCAL.get(); } @Override public Mono<Void> filter(ServerWebExchange exchange, WebFilterChain chain) { String headerValue = exchange.getRequest().getHeaders().getFirst("YOUR_HEADER_KEY"); return chain.filter(exchange) .contextWrite(context -> context.put("YOUR_HEADER_KEY", headerValue)) .doOnSubscribe(sub -> HEADER_VALUE_THREAD_LOCAL.set(headerValue)) .doFinally(signal -> HEADER_VALUE_THREAD_LOCAL.remove()); // 必须清理ThreadLocal } }
修改MessageSource实现
@RequiredArgsConstructor @Component public class RemoteAbstractMessageSource extends AbstractMessageSource { private final RestTemplate restTemplate; @Override public String resolveCode(String code, Locale locale) { if (StringUtils.containsWhitespace(code)) { return code; } // 从ThreadLocal获取头参数 String headerValue = RequestHeaderThreadLocalFilter.getCurrentHeaderValue(); var headers = new HttpHeaders(); headers.setAcceptLanguageAsLocales(List.of(locale)); var entity = new HttpEntity<>(headers); try { return restTemplate.exchange(UriComponentsBuilder.fromPath("/public/v1/messages/i18n/{code}") .queryParam("yourParam", headerValue) .build(code).toString(), HttpMethod.GET, entity, String.class).getBody(); } catch (HttpStatusCodeException e) { throw CustomErrorException.of(e); } } }
总结
- 优先选择方案一:完全遵循WebFlux的反应式设计,无阻塞、无线程污染问题,是Spring官方推荐的适配方式。
- 方案二仅作为临时兼容方案:破坏了反应式模型的无状态特性,存在线程安全风险,不建议长期使用。
内容的提问来源于stack exchange,提问作者fernando1979
相关产品推荐
相关产品推荐

