在Project Reactor中调用同步方法的问题及Reactor实现方式问询
问题分析与解决方案
一、当前实现的核心问题
- 阻塞Netty NIO线程:Netty的事件循环(EventLoop)线程负责处理所有非阻塞IO操作,这类线程设计上严禁被长时间阻塞。你的
synchronized refreshToken()方法会导致持有锁的NIO线程陷入阻塞,其他竞争锁的NIO线程也会进入等待状态,直接拉低应用的IO处理吞吐量,严重时会引发服务响应缓慢甚至无响应。 - 违背Reactor异步设计原则:Reactor基于异步非阻塞模型构建,直接在IO线程中调用同步阻塞方法,会破坏整个异步链的非阻塞特性,引发线程调度紊乱。
二、Reactor中调用同步/阻塞方法的正确姿势
所有同步阻塞操作必须切换到专门的阻塞操作线程池执行,避免占用IO线程。常用实现方式有两种:
- 使用
publishOn:切换链中后续操作的执行线程(适合在异步链中间插入阻塞操作的场景) - 使用
subscribeOn:切换整个Mono/Flux的订阅线程(适合源头本身是阻塞操作的场景)
推荐使用Reactor内置的Schedulers.boundedElastic(),它是专门为阻塞操作设计的线程池,会根据负载动态扩容缩容,避免资源耗尽。
三、实现"仅刷新一次令牌"的最优方案
不需要依赖synchronized这类同步锁,利用Reactor的异步特性即可实现并发安全的单例令牌刷新。以下是两种可落地的方案:
方案1:利用Mono.cache()实现令牌缓存与单例刷新
private AtomicReference<Mono<String>> currentTokenMono = new AtomicReference<>(); private final HttpClient reactorHttpClient; // 获取有效令牌,确保同一时刻只有一个刷新请求发出 private Mono<String> getValidToken() { return currentTokenMono.updateAndGet(mono -> mono == null ? refreshTokenMono().doFinally(s -> currentTokenMono.compareAndSet(mono, null)) : mono ); } // 异步执行令牌刷新(同步操作切换到阻塞线程池) private Mono<String> refreshTokenMono() { return Mono.fromCallable(() -> { // 原refreshToken()中的同步逻辑:请求服务器获取新令牌 Response response = syncTokenClient.getNewToken(); return response.getToken(); }).subscribeOn(Schedulers.boundedElastic()) .cache(); // 缓存令牌,直到过期或出错 } // 修正后的业务处理方法 Mono<Result> process(Request request) { return getValidToken() .map(this::buildAuthHeader) .flatMap(header -> reactorHttpClient.get() .uri(request.getUri()) .header(header) .responseSingle((resp, body) -> { if (resp.status().code() == 401) { return Mono.error(new InvalidAuthException("令牌已过期")); } return body.map(this::convertToResult); }) ) .onErrorResume(InvalidAuthException.class, e -> { // 令牌过期时清空缓存,触发重新刷新 currentTokenMono.set(null); return process(request); }); } // 辅助方法:用令牌生成认证Header private Header buildAuthHeader(String token) { return new Header("Authorization", "Bearer " + token); }
方案2:结合Semaphore与retryWhen实现重试+单例刷新
这种方案更适合需要重试逻辑的场景,同时确保刷新操作仅执行一次:
private volatile String authString; private final Semaphore refreshSemaphore = new Semaphore(1); private final HttpClient reactorHttpClient; Mono<Result> process(Request request) { return Mono.defer(() -> { Header header = buildAuthHeader(); return reactorHttpClient.get() .uri(request.getUri()) .header(header) .responseSingle((resp, body) -> { if (resp.status().code() == 401) { return Mono.error(new InvalidAuthException("令牌已过期")); } return body.map(this::convertToResult); }); }) .retryWhen(Retry.backoff(3, Duration.ofSeconds(1)) .filter(InvalidAuthException.class::isInstance) .doBeforeRetry(retrySignal -> { // 用Semaphore确保同一时刻只有一个线程执行刷新 if (refreshSemaphore.tryAcquire()) { try { // 同步刷新操作切换到阻塞线程池 Mono.fromCallable(() -> { Response response = syncTokenClient.getNewToken(); authString = response.getToken(); return authString; }).subscribeOn(Schedulers.boundedElastic()).block(); } finally { refreshSemaphore.release(); } } else { // 等待其他线程完成刷新 try { refreshSemaphore.acquire(); refreshSemaphore.release(); } catch (InterruptedException ignored) { Thread.currentThread().interrupt(); } } }) ); } private Header buildAuthHeader() { return new Header("Authorization", "Bearer " + authString); }
四、关键注意事项
- 严禁在Netty EventLoop线程执行阻塞操作:包括同步方法、线程睡眠、锁等待等,必须切换到
Schedulers.boundedElastic()这类专门的阻塞线程池。 - 避免用传统同步锁控制并发:在Reactor异步链中,
synchronized或ReentrantLock容易导致IO线程阻塞,优先用Reactor的异步特性(如cache()、原子引用、信号量配合线程池)实现并发安全。 - 令牌缓存的主动失效:可以在
cache()中设置过期时间,或者在捕获到令牌过期异常时主动清空缓存,确保下次请求触发新的刷新。
内容的提问来源于stack exchange,提问作者akshay
相关产品推荐
相关产品推荐

