一、先搞懂:WebFlux里的“红绿灯”和“快递队”

WebFlux是Spring提供的响应式编程框架,很多人刚接触的时候会觉得它和传统的Servlet不一样,其实可以用两个生活化的类比来理解:背压就是“红绿灯”,线程模型就是“快递队”。

1.1 背压:流量的“红绿灯”

传统的Servlet服务就像快递员一直猛跑,消费者(接口调用方)能接多少就接多少,结果就是快递堆成山,消费者垮掉。而WebFlux的背压就是红绿灯:上游(比如上游服务、消息队列)发数据的时候,会问下游“你能处理多少?”,下游说“我现在能处理5个”,上游就只发5个,不会塞太多,避免下游被压垮,相当于“按需供给”。

1.2 线程模型:处理任务的“快递队”

WebFlux的线程模型就像不同的快递队,有两种常用的:一种是“临时快递队(elastic)”,适合处理IO密集型任务(比如调用接口、读数据库),需要经常等,所以队员工资高(线程资源按需分配);另一种是“固定人数的快递队(parallel)”,适合CPU密集型任务(比如计算、排序),人数固定,不会随便加人,避免浪费资源。

二、真实踩坑:线上卡成龟速的那个午夜故障

我们团队曾经遇到过一次,凌晨2点运维告警,说核心接口的P99延迟从100ms涨到了10s,QPS直接掉了80%,但CPU使用率才20%(正常是50%左右),内存也没爆,排查了好久才找到问题,就是背压和线程模型的锅。

2.1 故障现场还原

当时的场景是我们做了一个消息订阅服务,用WebFlux消费Kafka的消息,然后转成API返回给前端。上线前做过小压测,一切正常,但一到凌晨流量峰值(有个活动的消息通知),就直接卡壳了。现场的日志里,全是“处理超时”“线程等待”,没有明显的错误栈,因为是异步模型,排查起来比传统服务麻烦。

2.2 找坑第一步:背压“红绿灯”坏了

当时我们写的代码里,为了方便,在Flux处理链里直接用了同步阻塞的操作,比如Thread.sleep()来模拟第三方接口调用,这里的问题就是:背压机制本来是控制上游不发太多数据,但下游的处理线程被阻塞了,相当于“绿灯一直亮,但行人走不动”,上游不知道下游的处理速度,还是猛发数据,下游的队列(背压缓冲)被撑爆,所有请求堆积,服务直接卡住。 我们当时写的错误代码是这样的:

// 技术栈:Spring Boot 2.7.12 + WebFlux + Reactor 3.5.12
// 错误示例:未正确处理背压,使用同步阻塞操作导致处理速度跟不上
import reactor.core.publisher.Flux;
import org.springframework.stereotype.Component;

@Component
public class ErrorMessageHandler {

    // 消费Kafka消息的Flux,每秒模拟收到1000条消息
    public Flux<String> consumeKafkaMessages() {
        return Flux.range(1, 10000) // 模拟10000条消息
                .map(messageId -> {
                    try {
                        // 模拟每条消息的处理需要100ms,这里是同步阻塞
                        Thread.sleep(100); 
                        return "处理完成:" + messageId;
                    } catch (InterruptedException e) {
                        Thread.currentThread().interrupt();
                        return "处理中断:" + messageId;
                    }
                });
    }
}

为什么这个代码会有问题?因为Flux的处理线程是reactor的默认线程,这个线程本来要处理请求,但这里被阻塞了,导致后面的请求无法被及时处理,背压的信号传不回去,上游一直发,下游的线程都被占满,新的请求进不来,服务就卡了。 后来我们改成了正确的背压处理方式,用异步的线程池来处理,代码变成这样:

// 技术栈:Spring Boot 2.7.12 + WebFlux + Reactor 3.5.12
// 正确示例:正确处理背压,用异步线程池避免阻塞主处理线程
import reactor.core.publisher.Flux;
import reactor.core.scheduler.Schedulers;
import org.springframework.stereotype.Component;

@Component
public class CorrectMessageHandler {

    public Flux<String> consumeKafkaMessages() {
        return Flux.range(1, 10000)
                .flatMap(messageId -> 
                    // 把消息处理放到parallel线程池里,不阻塞主处理线程
                    Mono.fromCallable(() -> {
                        // 模拟业务处理逻辑,这里的100ms阻塞是在线程池里,不会影响主流程
                        Thread.sleep(100);
                        return "处理完成:" + messageId;
                    }).subscribeOn(Schedulers.parallel()),
                    5 // 并发数限制:最多同时处理5个消息,这就是背压的核心,保证不会过载
                );
    }
}

这里的关键是,把处理逻辑放到了单独的线程池里,主处理线程不会被阻塞,同时用flatMap的第二个参数限制了并发数,也就是告诉上游“我每次最多处理5个”,背压生效,不会有太多请求堆积。

2.3 再查坑:线程模型的“快递队”超载

除了背压的问题,还有线程模型的配置不对。当时我们用的是默认的线程池,reactor的elastic线程池,这个线程池的最大线程数是动态扩容的,不过因为我们错误的阻塞操作,导致线程被占满,新的请求无法创建线程,就排队了。后来我们调整了线程池的配置,在application.yml里加了:

# Spring Boot WebFlux线程池配置,修正后的配置
spring:
  reactor:
    netty:
      # Netty主工作线程配置,处理网络请求
      worker-group:
        select-count: 4 # 选择器线程数,和CPU核心数一致,比如4核就4个
        worker-count: 8 # 工作线程数,CPU密集型是2*CPU,IO密集型可以调大到20左右
  tasks:
    # Spring异步任务线程池,处理业务逻辑
    core-pool-size: 10 # 核心线程数
    max-pool-size: 50 # 最大线程数,避免线程无限增长
    queue-capacity: 100 # 任务队列容量,超过后会拒绝新任务

调整之后,线程不会被轻易占满,服务的恢复速度快了很多。

三、避坑指南:线上运行要注意的几个关键点

遇到这个故障之后,我们总结了几个核心的避坑点,不管是新手还是老司机,都要注意。

3.1 背压的正确姿势

背压不是“设置了就完事”,要根据业务场景选对方式:

  • 如果业务数据丢了没关系(比如日志、非核心通知),就用onBackpressureDrop(),直接丢多余的,不会影响服务;
  • 如果数据不能丢,但允许少量延迟,就用onBackpressureBuffer(1000),设置缓冲大小,超过后拒绝新数据;
  • 如果只要最新的那条数据(比如实时监控),就用onBackpressureLatest(),只保留最新的,旧的直接丢。 举个例子,非核心的通知消息,用onBackpressureDrop的代码:
// 示例:非核心业务用onBackpressureDrop处理背压
return Flux.range(1, 10000)
        .onBackpressureDrop(dropMessage -> {
            // 把丢的消息记录下来,后面可以补偿,或者打日志
            log.warn("丢弃了消息:{}", dropMessage);
        })
        .flatMap(this::processMessage);

3.2 线程模型的配置技巧

线程池的配置要和业务匹配:

  • CPU密集型业务(比如计算、加密):用parallel线程池,核心数设为CPU核心数的1-2倍,避免线程切换;
  • IO密集型业务(比如调用第三方接口、读数据库):用elastic线程池,核心数可以设为20-50,因为大部分时间线程在等待,不需要太多CPU;
  • 不要在主处理线程(reactor的主线程)里做阻塞操作,所有阻塞的逻辑都要放到单独的线程池里。

3.3 压测和监控要跟上

一定要做压测,用工具比如JMeter、k6模拟高并发,重点看两个指标:

  • 背压相关的指标:reactor.backpressure.buffer.size(缓冲队列大小),如果这个值持续上涨,说明背压有问题;
  • 线程相关的指标:reactor.active.threads(活跃线程数),如果持续达到最大线程数,说明线程不够,需要调整。 监控这些指标可以用Spring Boot Actuator,比如访问/actuator/metrics/reactor.backpressure.buffer.size就能看到实时的缓冲大小,提前发现问题。

四、总结:踩坑后的几个核心感悟

WebFlux的优势很明显,适合高并发的IO密集型服务,比如API网关、消息订阅、实时通知,能扛住比传统Servlet多几倍的流量,资源利用率更高。但它的坑也不少,尤其是背压和线程模型,没弄好的话,线上故障比传统服务更难排查,因为是异步模型,堆栈不清晰,请求的链路也长。 核心的注意点就是:第一,不要随便用同步阻塞操作,把所有阻塞逻辑放到单独的线程池;第二,背压要根据业务选对应的策略,不能一概而论;第三,线程池的配置要和业务匹配,不能用默认的就不管了;第四,压测和监控是必须的,线上故障大部分是没提前压测出来的。 总的来说,WebFlux不是银弹,适合的场景才用,用的时候一定要把背压和线程模型搞懂,不然很容易踩坑,就像我们那次午夜的故障,排查了3个小时才找到问题,后来调整之后,服务的稳定性提升了很多,峰值流量的时候也没再卡过。