Programação Reativa com Spring WebFlux
WebFlux é o motor assíncrono do Spring — com um único Mono e o Flux de primeira classe, sua E/S não bloqueante suporta dezenas de milhares de conexões concorrentes.
1. O Que Você Vai Aprender
- Paradigma de Programação Reativa: Conceitos e Operadores Centrais de Mono / Flux
@RestController+Mono<T>/Flux<T>Escrevendo APIs Responsivas- WebFlux vs. WebMVC: Comparação de Casos de Uso e Guia de Seleção
- Repositórios Reativos (Spring Data R2DBC)
- Alice implementa uma API reativa capaz de processar milhares de consultas de pedidos por segundo
2. Uma História Real de um Arquiteto de Alta Concorrência
(1) Dor: Esgotamento do Pool de Threads
O OrderFlow de Alice enfrentou um gargalo durante um grande evento de vendas — com 500 usuários concorrentes consultando pedidos, o pool de threads do Tomcat (configurado com 200 threads por padrão) foi esgotado, as requisições foram enfileiradas e a latência P99 disparou para 3 segundos. Na maior parte do tempo, as threads estavam aguardando E/S (banco de dados, rede), e a utilização da CPU era de apenas 15%. Bob sugeriu aumentar o número de threads, mas adicionar 1.000 threads consumiria mais 1 GB de memória.
(2) A Solução WebFlux
O WebFlux gerencia alta concorrência com poucas threads — as threads não bloqueiam durante esperas de E/S e podem continuar atendendo outras requisições:
@GetMapping("/{id}")
public Mono<OrderResponse> getOrder(@PathVariable Long id) {
return orderRepository.findById(id)
.map(OrderResponse::from)
.switchIfEmpty(Mono.error(new ResourceNotFoundException("Order", id)));
}
(3) Resultado
Depois que Alice refez a API de consulta usando WebFlux, quatro threads foram capazes de lidar com 2.000 requisições concorrentes, o uso de memória foi reduzido em 70%, a latência P99 caiu de 3 segundos para 100 ms e a utilização da CPU aumentou para 60%.
3. Conceitos Centrais da Programação Reativa
(1) Mono e Flux
| Tipo | Significado | Analogia | Exemplo |
|---|---|---|---|
Mono<T> |
0 ou 1 elemento | Versão assíncrona do Optional | Consultar um único pedido |
Flux<T> |
0 a N elementos | Versão assíncrona do Stream | Consultar lista de pedidos |
Mono<Void> |
Sem retorno | Sinal de conclusão | Operação de exclusão |
graph LR
A["Publisher"] --> B["Mono<br/>0 ou 1 elemento"]
A --> C["Flux<br/>0 a N elementos"]
B --> D["onNext → onComplete"]
C --> E["onNext x N → onComplete"]
B --> F["onError"]
C --> F
(1) ▶ Exemplo: Operações Básicas com Mono e Flux
// Mono: valor único
Mono<String> mono = Mono.just("Hello OrderFlow");
Mono<String> empty = Mono.empty();
Mono<String> fromCallable = Mono.fromCallable(() -> fetchOrder(1L));
// Flux: múltiplos valores
Flux<Integer> flux = Flux.just(1, 2, 3, 4, 5);
Flux<Integer> range = Flux.range(1, 100);
Flux<Long> interval = Flux.interval(Duration.ofSeconds(1));
Saída:
// Execução bem-sucedida
(2) Operadores Comuns
| Operador | Função | Exemplo |
|---|---|---|
map |
Converter Elementos | .map(OrderResponse::from) |
flatMap |
Conversão Assíncrona | .flatMap(this::enrichOrder) |
filter |
Filtrar | .filter(order -> "PENDING".equals(order.getStatus())) |
switchIfEmpty |
Substituição para Nulo | .switchIfEmpty(Mono.error(...)) |
onErrorResume |
Degradar Erro | .onErrorResume(e -> fallbackOrder()) |
timeout |
Controle de Timeout | .timeout(Duration.ofSeconds(5)) |
retry |
Retentativa | .retry(3) |
subscribeOn |
Especificar Thread de Inscrição | .subscribeOn(Schedulers.boundedElastic()) |
(2) ▶ Exemplo: Encadeamento de Operadores
public Mono<OrderResponse> getOrderWithDetails(Long id) {
return orderRepository.findById(id)
.flatMap(order ->
productRepository.findById(order.getProductId())
.map(product -> OrderResponse.from(order, product)))
.switchIfEmpty(Mono.error(
new ResourceNotFoundException("Order", id)))
.timeout(Duration.ofSeconds(3))
.onErrorResume(TimeoutException.class,
e -> Mono.error(new RuntimeException("Timeout na consulta do pedido")))
.retry(2);
}
Saída:
// Execução bem-sucedida
4. WebFlux REST API
(1) ▶ Exemplo: Controller Reativo
@RestController
@RequestMapping("/api/v1/reactive/orders")
public class ReactiveOrderController {
private final ReactiveOrderService orderService;
public ReactiveOrderController(ReactiveOrderService orderService) {
this.orderService = orderService;
}
@GetMapping("/{id}")
public Mono<OrderResponse> getOrder(@PathVariable Long id) {
return orderService.findById(id);
}
@GetMapping
public Flux<OrderResponse> listOrders(
@RequestParam(defaultValue = "0") int page,
@RequestParam(defaultValue = "20") int size) {
return orderService.findAll(page, size);
}
@PostMapping
public Mono<ResponseEntity<OrderResponse>> createOrder(
@Valid @RequestBody CreateOrderRequest request) {
return orderService.create(request)
.map(order -> ResponseEntity
.status(HttpStatus.CREATED).body(order));
}
@DeleteMapping("/{id}")
public Mono<ResponseEntity<Void>> cancelOrder(@PathVariable Long id) {
return orderService.cancel(id)
.then(Mono.just(ResponseEntity.noContent().<Void>build()));
}
}
Saída:
// Execução bem-sucedida
| Dimensão | WebMVC | WebFlux |
|---|---|---|
| Tipo de Retorno | Objeto Síncrono | Mono<T> / Flux<T> |
| Modelo de Threads | Uma thread por requisição | Event loop (poucas threads) |
| Modelo de E/S | Bloqueante | Não bloqueante |
| Container | Tomcat / Jetty | Netty |
| Limite de Concorrência | Tamanho do Pool de Threads | Virtualmente Ilimitado |
5. Escolhendo Entre WebFlux e WebMVC
(1) Decisão de Seleção de Produto
| Cenário | Recomendação | Motivo |
|---|---|---|
| API CRUD, concorrência < 500 | WebMVC | Simples e intuitivo, ecossistema maduro |
| Intensivo em E/S, alta concorrência | WebFlux | Não bloqueante, suporta alta concorrência com poucas threads |
| Streaming em tempo real (SSE/WebSocket) | WebFlux | Suporte nativo |
| Grande número de operações bloqueantes (JDBC) | WebMVC | Operações bloqueantes são ainda piores no WebFlux |
| Cenários Mistos | WebMVC + Assincronia Parcial | Otimização Progressiva |
(1) ▶ Exemplo: Server-Sent Events (SSE)
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<OrderEvent> streamOrderEvents() {
return orderEventPublisher.eventStream()
.log("order-events");
}
Saída:
// Execução bem-sucedida
6. Spring Data R2DBC
(1) Acesso Reativo ao Banco de Dados
(1) ▶ Exemplo: Entidade e Repositório R2DBC
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-r2dbc</artifactId>
</dependency>
<dependency>
<groupId>io.r2dbc</groupId>
<artifactId>r2dbc-mysql</artifactId>
</dependency>
Saída:
// Execução bem-sucedida
@Table("products")
public class Product {
@Id
private Long id;
private String name;
private BigDecimal price;
private Integer stock;
public Product() {}
public Product(String name, BigDecimal price, Integer stock) {
this.name = name;
this.price = price;
this.stock = stock;
}
// getters
}
public interface ReactiveProductRepository
extends ReactiveCrudRepository<Product, Long> {
Flux<Product> findByNameContaining(String keyword);
Flux<Product> findByStockLessThan(Integer threshold);
}
| Dimensão | JPA + JDBC | R2DBC |
|---|---|---|
| Modelo do Driver | Bloqueante | Não bloqueante |
| Estilo da API | Síncrono | Reativo (Mono/Flux) |
| Mapeamento de Relacionamento | Suportado (@OneToMany) | Não Suportado |
| Linguagem de Consulta | JPQL | SQL Nativo / Derivação por Nome de Método |
| Transação | @Transactional |
@Transactional (Reativo) |
(2) ▶ Exemplo: Service Reativo
@Service
public class ReactiveOrderService {
private final ReactiveOrderRepository orderRepository;
private final ReactiveProductRepository productRepository;
public ReactiveOrderService(ReactiveOrderRepository orderRepository,
ReactiveProductRepository productRepository) {
this.orderRepository = orderRepository;
this.productRepository = productRepository;
}
public Mono<OrderResponse> create(CreateOrderRequest request) {
return productRepository.findById(request.productId())
.switchIfEmpty(Mono.error(
new ResourceNotFoundException("Product", request.productId())))
.flatMap(product -> {
if (product.getStock() < request.quantity()) {
return Mono.error(new InsufficientStockException(
product.getId(), product.getStock(), request.quantity()));
}
product.deductStock(request.quantity());
return productRepository.save(product)
.then(orderRepository.save(new Order(product.getId(), request.quantity())));
})
.map(OrderResponse::from);
}
public Flux<OrderResponse> findAll(int page, int size) {
return orderRepository.findAll()
.skip(page * size)
.take(size)
.map(OrderResponse::from);
}
public Mono<Void> cancel(Long orderId) {
return orderRepository.findById(orderId)
.flatMap(order -> {
order.setStatus("CANCELLED");
return orderRepository.save(order);
})
.then();
}
}
Saída:
// Execução bem-sucedida
7. Exemplo Completo: API de Consulta Reativa do OrderFlow
// ReactiveProductRepository.java
public interface ReactiveProductRepository
extends ReactiveCrudRepository<Product, Long> {
Flux<Product> findByNameContaining(String keyword);
}
// ReactiveOrderRepository.java
public interface ReactiveOrderRepository
extends ReactiveCrudRepository<Order, Long> {
Flux<Order> findByStatus(String status);
}
// ReactiveOrderController.java
@RestController
@RequestMapping("/api/v1/reactive/orders")
public class ReactiveOrderController {
private final ReactiveOrderService orderService;
public ReactiveOrderController(ReactiveOrderService orderService) {
this.orderService = orderService;
}
@GetMapping("/{id}")
public Mono<OrderResponse> get(@PathVariable Long id) {
return orderService.findById(id);
}
@GetMapping
public Flux<OrderResponse> list(
@RequestParam(defaultValue = "0") int page,
@RequestParam(defaultValue = "20") int size) {
return orderService.findAll(page, size);
}
@GetMapping("/status/{status}")
public Flux<OrderResponse> byStatus(@PathVariable String status) {
return orderService.findByStatus(status);
}
@PostMapping
public Mono<ResponseEntity<OrderResponse>> create(
@RequestBody CreateOrderRequest req) {
return orderService.create(req)
.map(o -> ResponseEntity.status(HttpStatus.CREATED).body(o));
}
@DeleteMapping("/{id}")
public Mono<ResponseEntity<Void>> cancel(@PathVariable Long id) {
return orderService.cancel(id)
.thenReturn(ResponseEntity.noContent().<Void>build());
}
}
// ReactiveOrderService.java
@Service
public class ReactiveOrderService {
private final ReactiveOrderRepository orderRepo;
private final ReactiveProductRepository productRepo;
public ReactiveOrderService(ReactiveOrderRepository orderRepo,
ReactiveProductRepository productRepo) {
this.orderRepo = orderRepo;
this.productRepo = productRepo;
}
public Mono<OrderResponse> findById(Long id) {
return orderRepo.findById(id)
.map(OrderResponse::from)
.switchIfEmpty(Mono.error(new ResourceNotFoundException("Order", id)));
}
public Flux<OrderResponse> findAll(int page, int size) {
return orderRepo.findAll().skip((long) page * size).take(size)
.map(OrderResponse::from);
}
public Flux<OrderResponse> findByStatus(String status) {
return orderRepo.findByStatus(status).map(OrderResponse::from);
}
public Mono<OrderResponse> create(CreateOrderRequest req) {
return productRepo.findById(req.productId())
.switchIfEmpty(Mono.error(new ResourceNotFoundException("Product", req.productId())))
.flatMap(p -> {
if (p.getStock() < req.quantity()) {
return Mono.error(new InsufficientStockException(p.getId(), p.getStock(), req.quantity()));
}
p.deductStock(req.quantity());
return productRepo.save(p)
.then(orderRepo.save(new Order(p.getId(), req.quantity())));
})
.map(OrderResponse::from);
}
public Mono<Void> cancel(Long id) {
return orderRepo.findById(id)
.flatMap(o -> { o.setStatus("CANCELLED"); return orderRepo.save(o); })
.then();
}
}
❓ Perguntas Frequentes
flatMap.Mono.fromCallable() + subscribeOn(Schedulers.boundedElastic()) para agendar operações bloqueantes no pool de threads elástico. No entanto, uso frequente indica que essa abordagem não é adequada para o WebFlux.onErrorResume / onErrorReturn; 2) @ExceptionHandler global (também suportado pelo WebFlux); 3) onErrorMap para converter tipos de exceção.WebTestClient em vez de MockMvc. @WebFluxTest carrega a camada WebFlux. StepVerifier é usado para testar sequências de elementos em Mono/Flux.📖 Resumo
- Mono (0/1 elementos) e Flux (0-N elementos) são os tipos centrais da programação reativa
- Chamadas encadeadas de operadores: map, flatMap, filter, switchIfEmpty, timeout, retry
- WebFlux usa Netty e event loop para suportar alta concorrência com um pequeno número de threads
- WebFlux é adequado para cenários intensivos em E/S e de alta concorrência, mas não é adequado para um grande número de operações bloqueantes
- R2DBC é um driver reativo de banco de dados; não suporta mapeamento relacional, portanto você deve combinar dados manualmente usando
flatMap - WebMVC e WebFlux não devem ser usados juntos; escolha um com base no cenário específico
📝 Exercícios
-
Problema Básico (Dificuldade ⭐): Implemente uma API reativa de consulta de produtos usando WebFlux (
GET /api/v1/reactive/products/{id}), conectando-se a um banco de dados em memória H2 via R2DBC. -
Exercício Avançado (Dificuldade: ⭐⭐): Implemente uma API reativa de realização de pedidos que cubra todo o processo — incluindo verificação de estoque, dedução de estoque e criação do pedido — usando
flatMappara combinar múltiplas operações reativas. -
Desafio (Dificuldade: ⭐⭐⭐): Use
WebTestClienteStepVerifierpara escrever um teste de API reativa, implemente uma interface SSE de push de eventos de pedidos em tempo real e compare as diferenças de desempenho entre WebMVC e WebFlux sob 1.000 conexões concorrentes.



