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/