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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 16:57:46