Spring Boot 3.x WebFlux响应式编程高并发实战与性能调优

WebFlux解决什么场景的问题

Spring WebFlux基于Reactor和Netty,用少量线程处理大量并发连接。核心场景:IO密集型业务(网关聚合、消息推送、实时数据流),单机支撑10万+连接。而传统Spring MVC的Thread-Per-Request模型在IO等待时线程阻塞,线程池用满后请求排队超时。

选WebFlux还是MVC的判断标准:如果业务是CPU密集型(计算、加解密),MVC足够;如果业务大量等待外部IO(数据库、HTTP调用、MQ消费),WebFlux的资源利用率远高于MVC。混合场景可以在同一应用中同时使用(MVC处理计算类接口,WebFlux处理IO类接口)。

响应式编程核心模型

Reactor的两个核心类型:Mono(0或1个元素的异步序列)和Flux(0到N个元素的异步序列)。

操作链式调用 vs 命令式代码,是响应式编程的思维转换:

// 命令式:获取用户 → 获取订单 → 获取商品
public UserDetail getUserDetail(Long userId) {
    User user = userClient.getUser(userId);       // 阻塞等待
    List<Order> orders = orderClient.getOrders(user.getId()); // 再阻塞
    List<Product> products = productClient.getProducts(orders); // 再阻塞
    return new UserDetail(user, orders, products);
}

// 响应式:三步串联,全程非阻塞
public Mono<UserDetail> getUserDetail(Long userId) {
    return userClient.getUser(userId)
        .flatMap(user ->
            orderClient.getOrders(user.getId())
                .flatMap(orders ->
                    productClient.getProducts(orders)
                        .map(products -> new UserDetail(user, orders, products))
                )
        );
}

并行聚合——多个独立IO调用并行执行:

public Mono<DashboardData> getDashboard(Long userId) {
    Mono<User> userMono = userClient.getUser(userId).cache();
    Mono<List<Order>> ordersMono = orderClient.getRecentOrders(userId);
    Mono<AccountInfo> accountMono = accountClient.getAccount(userId);
    Mono<List<Notice>> noticeMono = noticeClient.getNotices(userId);

    return Mono.zip(userMono, ordersMono, accountMono, noticeMono)
        .map(tuple -> new DashboardData(
            tuple.getT1(), tuple.getT2(), tuple.getT3(), tuple.getT4()
        ));
}

Mono.zip并行发起4个请求,总耗时等于最慢的那个,而非4个请求耗时之和。这是WebFlux高并发的关键——减少等待时间,提高线程利用率。

WebFlux接口开发实战

Router + Handler模式(函数式端点,推荐):

@Configuration
public class OrderRouter {

    @Bean
    public RouterFunction<ServerResponse> orderRoutes(OrderHandler handler) {
        return RouterFunctions.route()
            .GET("/api/orders", handler::listOrders)
            .GET("/api/orders/{id}", handler::getOrder)
            .POST("/api/orders", handler::createOrder)
            .PUT("/api/orders/{id}", handler::updateOrder)
            .DELETE("/api/orders/{id}", handler::deleteOrder)
            .build();
    }
}

@Component
@RequiredArgsConstructor
public class OrderHandler {

    private final OrderService orderService;

    public Mono<ServerResponse> listOrders(ServerRequest request) {
        int page = Integer.parseInt(request.queryParam("page").orElse("1"));
        int size = Integer.parseInt(request.queryParam("size").orElse("20"));

        return orderService.listOrders(page, size)
            .collectList()
            .flatMap(data -> ServerResponse.ok()
                .contentType(MediaType.APPLICATION_JSON)
                .bodyValue(data));
    }

    public Mono<ServerResponse> createOrder(ServerRequest request) {
        return request.bodyToMono(CreateOrderRequest.class)
            .flatMap(req -> orderService.createOrder(req))
            .flatMap(order -> ServerResponse
                .created(URI.create("/api/orders/" + order.getId()))
                .bodyValue(order))
            .onErrorResume(ValidationException.class, e ->
                ServerResponse.badRequest().bodyValue(Map.of("error", e.getMessage())));
    }
}

SSE实时推送

// Server-Sent Events端点
public Mono<ServerResponse> streamNotifications(ServerRequest request) {
    String userId = request.pathVariable("userId");

    Flux<ServerSentEvent<Notification>> events = notificationService
        .streamNotifications(userId)
        .map(n -> ServerSentEvent.<Notification>builder()
            .id(String.valueOf(n.getId()))
            .event("notification")
            .data(n)
            .build());

    return ServerResponse.ok()
        .contentType(MediaType.TEXT_EVENT_STREAM)
        .body(events, ServerSentEvent.class);
}

响应式数据库访问

Spring Data R2DBC是WebFlux下数据库访问的标准方案,基于非阻塞IO:

// Repository
public interface OrderRepository extends ReactiveCrudRepository<Order, Long> {
    Flux<Order> findByUserIdAndStatus(Long userId, String status);
    Mono<Long> countByStatus(String status);

    @Query("SELECT * FROM orders WHERE created_at > :since ORDER BY created_at DESC LIMIT :limit")
    Flux<Order> findRecentOrders(@Param("since") LocalDateTime since, @Param("limit") int limit);
}

// Service
@Service
@RequiredArgsConstructor
public class OrderService {

    private final OrderRepository orderRepository;

    public Mono<Order> createOrder(CreateOrderRequest req) {
        return Mono.just(req)
            .map(r -> Order.builder()
                .userId(r.getUserId())
                .productId(r.getProductId())
                .quantity(r.getQuantity())
                .status("CREATED")
                .createdAt(LocalDateTime.now())
                .build())
            .flatMap(orderRepository::save);
    }

    public Flux<Order> listOrders(int page, int size) {
        return orderRepository.findAll()
            .skip((long) (page - 1) * size)
            .take(size);
    }
}

R2DBC连接池配置:

# application.yml
spring:
  r2dbc:
    url: r2dbc:mysql://db-host:3306/order_db
    username: app_user
    password: ${DB_PASSWORD}
    pool:
      max-size: 50
      min-idle: 10
      max-idle-time: 30m
      validation-depth: LOCAL
      initial-size: 10

背压与流控

Flux的生产速度可能远超消费速度,背压机制控制数据流速:

// 限速:每秒处理100条
dataFlux
    .limitRate(100)  // 每次预取100条
    .delayElements(Duration.ofMillis(10))  // 每条间隔10ms

// 缓冲:积攒批量处理
dataFlux
    .buffer(500)  // 每500条一批
    .flatMap(batch -> batchSaveService.saveBatch(batch), 4)

// 采样:只取最新
dataFlux
    .sample(Duration.ofSeconds(1))

错误处理与熔断

响应式链路的错误处理比命令式复杂——异常可能在任意操作符中抛出,需要沿链路向上传播:

// 全局错误处理
@Component
@Order(-2)
public class GlobalErrorHandler implements WebExceptionHandler {

    @Override
    public Mono<Void> handle(ServerWebExchange exchange, Throwable ex) {
        if (ex instanceof ResponseStatusException rse) {
            exchange.getResponse().setStatusCode(rse.getStatus());
        } else if (ex instanceof ValidationException) {
            exchange.getResponse().setStatusCode(HttpStatus.BAD_REQUEST);
        } else {
            exchange.getResponse().setStatusCode(HttpStatus.INTERNAL_SERVER_ERROR);
        }

        byte[] bytes = Map.of("error", ex.getMessage()).toString().getBytes();
        DataBuffer buffer = exchange.getResponse().bufferFactory().wrap(bytes);
        exchange.getResponse().getHeaders().setContentType(MediaType.APPLICATION_JSON);
        return exchange.getResponse().writeWith(Mono.just(buffer));
    }
}

Resilience4j熔断器集成:

@Bean
public Customizer<ReactiveCircuitBreakerFactory<?, ?>> circuitBreakerCustomizer() {
    return factory -> factory.configureDefault(id ->
        new ReactiveCircuitBreakerConfig()
            .setSlidingWindowSize(20)
            .setFailureRateThreshold(50)
            .setWaitDurationInOpenState(Duration.ofSeconds(30))
            .setPermittedNumberOfCallsInHalfOpenState(5)
    );
}

// 使用熔断器
@Service
public class ResilientOrderService {

    @Autowired
    private ReactiveCircuitBreakerFactory<?, ?> cbFactory;

    public Mono<Order> getOrder(Long orderId) {
        return cbFactory.create("order-service")
            .run(
                orderClient.getOrder(orderId),
                throwable -> fallbackOrder(orderId)
            );
    }
}

性能调优关键参数

Netty线程模型调优:

# application.yml
server:
  netty:
    connection-timeout: 5000
    idle-timeout: 30000

# JVM参数(8核16G服务器示例)
# -Xms4g -Xmx4g
# -XX:+UseG1GC
# -XX:MaxGCPauseMillis=100
# -Dreactor.netty.ioWorkerCount=8
# -Dreactor.netty.ioSelectCount=2
# -Dreactor.netty.pool.maxConnections=500

监控指标:

WebFlux关键指标与MVC不同——不关注线程池使用率,而是关注:

Reactor延迟分布:onNext延迟P50/P95/P99
背压信号:request调用频率与数量
取消率:订阅后取消的比例(客户端断连)
错误率:onError信号占比

WebFlux不是银弹。如果团队对响应式编程不熟悉,强行引入会增加调试难度和上线风险。从高IO场景的接口开始试点,积累经验后再扩大范围,是更稳妥的路径。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot3xwebflux-xiang-ying-shi-bian-cheng-gao-bing-fa/

(0)
小编小编
上一篇 14小时前
下一篇 14小时前

相关推荐