Spring WebFlux & Reactive: Mono, Flux & WebClient
Learn Spring WebFlux reactive programming: Master Mono and Flux from Project Reactor, build non-blocking controllers, and use WebClient for async HTTP calls.
Learn Spring WebFlux reactive programming: Master Mono and Flux from Project Reactor, build non-blocking controllers, and use WebClient for async HTTP calls. The guide uses practical examples to explain spring webflux & reactive: mono, flux & webclient, project reactor core types and shows how to apply the ideas in a Spring Boot project.
Introduction
Java code has traditionally treated thread blocking as normal. Call a database, the thread sits there waiting. Hit an external API, it idles. Under load this wastes resources.
Reactive programming inverts this. Instead of blocking while you wait, you register a callback and move on. The result arrives whenever it arrives, and something else handles the notification. Spring WebFlux shipped this model as a mainstream option.
WebFlux landed in Spring Framework 5.0 as a fully non-blocking web framework built on reactive streams. It runs on servers like Netty without servlet container overhead. WebFlux implements the Reactive Streams specification through Project Reactor, so every operator and backpressure mechanism follows a shared contract.
The benefit is throughput. A blocking server exhausts its thread pool under load. A WebFlux server handles more concurrent requests on the same hardware, because threads never sit idle waiting for network responses.
Spring WebFlux & Reactive: Mono, Flux & WebClient
Project Reactor Core Types
Project Reactor is the reactive library that powers WebFlux. It has two primary types, and getting comfortable with both is prerequisite to anything else.
Mono
Mono represents a pipeline that emits zero or one element. Think of it as an async Optional, or a CompletionStage that can also be empty. You reach for Mono when an operation yields a single result or nothing at all.
Mono<User> findById(Long id);
Mono<Void> save(User user);
Mono<Boolean> existsByUsername(String username);
The empty case deserves attention. A Mono can complete without emitting any value at all, which is meaningfully different from “no result found” in blocking code. Your code has to handle both the value-present and the empty-complete paths.
Flux
Flux represents a pipeline that emits zero to n elements. It is the reactive equivalent of a Stream, except it can also signal completion or an error. You use Flux when an operation might return multiple results.
Flux<User> findAll();
Flux<Order> findByCustomerId(Long customerId);
Both Mono and Flux are lazy. Nothing runs until you subscribe. This is a fundamental shift from imperative code where calling a method executes it immediately. With reactive types, calling findAll() does not hit the database — it builds a description of a pipeline that will hit the database when something subscribes.
Key Operators
Reactor ships with an enormous operator library. A handful show up constantly in actual applications.
map transforms each element in the pipeline. It runs synchronously on whatever thread the previous step was using.
Flux<User> activeUsers = userRepository.findAll()
.filter(user -> user.isActive())
.map(user -> user.toSummaryDto());
flatMap transforms each element into a Mono or Flux and then merges all those nested pipelines into one. This is how you chain asynchronous operations.
Flux<OrderDto> orderDtos = orderRepository.findByCustomerId(id)
.flatMap(order -> productService.fetchProduct(order.getProductId())
.map(product -> new OrderDto(order, product)));
zip combines multiple reactive pipelines pairwise. Both inputs need to emit the same number of elements, and zip fires whenever all sources have each emitted a new item.
Mono<CombinedReport> report = Mono.zip(
userService.fetchUser(id).single(),
orderService.fetchRecentOrders(id).collectList(),
(user, orders) -> new CombinedReport(user, orders)
);
switchIfEmpty substitutes an alternative Mono when the source completes without emitting. It handles the empty case without throwing.
Mono<Cart> cart = cartRepository.findByUserId(id)
.switchIfEmpty(Mono.just(Cart.empty(userId)));
When to Use Reactive — And When Not To
Reactive programming solves real problems. It also introduces complexity that is not always justified. Here is how to think it through.
Choose Reactive When
You are building a system that fans out to multiple external services. If service A takes 100ms, B takes 150ms, and C takes 80ms, calling them sequentially in blocking code costs 330ms. Calling them reactively with flatMap or zip costs roughly 150ms — the slowest one. That math adds up fast at scale.
You are building a data streaming endpoint. Server-sent events, WebSocket streams, anything where data arrives incrementally maps naturally to Flux.
You are using a database with a reactive driver. R2DBC works for relational databases, and reactive MongoDB and Cassandra drivers are solid. If your database has no reactive driver, you are blocking anyway, and the complexity is pure overhead.
Avoid Reactive When
Your team is new to reactive programming. The debugging experience is genuinely different. Stack traces span multiple threads. NullPointerException inside an operator behaves oddly. Tracking down a blocking call buried in a reactive chain can eat half a day.
Your dependencies do not support reactive. Wrapping blocking calls in Mono.fromCallable() or Schedulers.boundedElastic() technically works, but you have just hidden the blocking. The thread pool managing those calls becomes your bottleneck and your blast radius.
Simple CRUD with no meaningful concurrency needs? Just use Spring MVC. It is easier to write, easier to debug, and easier to hand off to the next developer.
Trade-off Table
| Factor | Blocking (Spring MVC) | Reactive (WebFlux) |
|---|---|---|
| Thread utilization | Poor under I/O wait | High — threads never block |
| Throughput under load | Limited by thread pool size | High — event-loop based |
| Learning curve | Low | Steep — new mental model |
| Debugging | Familiar stack traces | Multi-thread, async stack traces |
| Ecosystem support | Universal JDBC drivers | Reactive drivers only |
| Backpressure | Not built in | Native via Reactive Streams spec |
| Memory per connection | One thread ~1MB stack | A few kilobytes |
| Code readability | Straightforward | Requires operator fluency |
| Right fit | Simple CRUD, CPU-bound | I/O-heavy, high concurrency |
Reactive Flow Architecture
Understanding the path a request takes through a reactive system is worth your time. Here is how a typical call moves through the stack.
graph TD
A["Client Request"] --> B["Netty Event Loop"]
B --> C["WebFlux DispatcherHandler"]
C --> D["Handler Function / Controller"]
D --> E["Service Layer<br/>Mono / Flux"]
E --> F["Reactive Repository<br/>R2DBC / Reactive Mongo"]
F --> E
E --> G["Data Transform Operators"]
G --> H["Response Encoding"]
H --> I["Client Response"]
The Netty event loop never blocks waiting for the database. When the database responds, the event loop picks up the continuation and resumes the pipeline. This is the core mechanism that lets a handful of threads handle massive concurrency.
Implementation Snippets
Annotation-Based Controller
The annotation-based style looks familiar if you have used Spring MVC. The key difference is the return types.
@RestController
@RequestMapping("/api/users")
public class UserController {
private final UserRepository userRepository;
public UserController(UserRepository userRepository) {
this.userRepository = userRepository;
}
@GetMapping("/{id}")
public Mono<ResponseEntity<UserDto>> getUser(@PathVariable Long id) {
return userRepository.findById(id)
.map(user -> ResponseEntity.ok(user.toDto()))
.defaultIfEmpty(ResponseEntity.notFound().build());
}
@GetMapping
public Flux<UserDto> getAllUsers(
@RequestParam(defaultValue = "0") int page,
@RequestParam(defaultValue = "20") int size) {
return userRepository.findAll()
.skip((long) page * size)
.take(size)
.map(User::toDto);
}
@PostMapping
public Mono<ResponseEntity<UserDto>> createUser(
@RequestBody @Valid UserCreateRequest request) {
return userRepository.save(request.toEntity())
.map(user -> ResponseEntity
.created(URI.create("/api/users/" + user.getId()))
.body(user.toDto()));
}
@DeleteMapping("/{id}")
public Mono<ResponseEntity<Void>> deleteUser(@PathVariable Long id) {
return userRepository.deleteById(id)
.then(Mono.just(ResponseEntity.noContent().<Void>build()));
}
}
Functional Routing Configuration
Functional routing separates routing configuration from the actual handler logic. Routes are defined as a bean returning a RouterFunction.
@Configuration
public class UserRoutingConfig {
@Bean
public RouterFunction<ServerResponse> userRoutes(UserHandler userHandler) {
return route()
.GET("/api/users/{id}", userHandler::getUser)
.GET("/api/users", userHandler::getAllUsers)
.POST("/api/users", userHandler::createUser)
.DELETE("/api/users/{id}", userHandler::deleteUser)
.build();
}
}
The handler receives ServerRequest objects instead of method parameters parsed by annotations.
@Component
public class UserHandler {
private final UserRepository userRepository;
public UserHandler(UserRepository userRepository) {
this.userRepository = userRepository;
}
public Mono<ServerResponse> getUser(ServerRequest request) {
Long id = Long.parseLong(request.pathVariable("id"));
return userRepository.findById(id)
.map(user -> ServerResponse.ok().bodyValue(user.toDto()))
.defaultIfEmpty(ServerResponse.notFound().build())
.flatMap(sr -> sr);
}
public Mono<ServerResponse> getAllUsers(ServerRequest request) {
return ServerResponse.ok()
.body(userRepository.findAll().map(User::toDto), UserDto.class);
}
public Mono<ServerResponse> createUser(ServerRequest request) {
return request.bodyToMono(UserCreateRequest.class)
.flatMap(userRepository::save)
.map(user -> ServerResponse
.created(URI.create("/api/users/" + user.getId()))
.bodyValue(user.toDto()))
.flatMap(sr -> sr);
}
}
WebClient with Retry and Timeout
WebClient is the reactive HTTP client that comes with WebFlux. It replaces RestTemplate in reactive applications.
@Configuration
public class ExternalApiConfig {
@Bean
public WebClient webClient(WebClient.Builder builder) {
return builder
.baseUrl("https://api.example.com")
.defaultHeader(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE)
.build();
}
}
@Service
public class ProductService {
private final WebClient webClient;
public ProductService(WebClient webClient) {
this.webClient = webClient;
}
public Mono<ProductSummary> fetchProductSummary(String productId) {
return webClient.get()
.uri("/products/{id}/summary", productId)
.retrieve()
.bodyToMono(ProductSummary.class)
.retryWhen(Retry.backoff(3, Duration.ofMillis(200))
.filter(ex -> ex instanceof WebClientResponseException.ServiceUnavailable)
.doBeforeRetry(signal -> {
System.err.println("Retrying after: " + signal.failure().getMessage());
}))
.timeout(Duration.ofMillis(1500))
.onErrorResume(WebClientResponseException.class,
ex -> Mono.just(ProductSummary.unavailable(productId)));
}
public Flux<Review> fetchReviews(String productId) {
return webClient.get()
.uri("/products/{id}/reviews", productId)
.retrieve()
.bodyToFlux(Review.class)
.take(10)
.timeout(Duration.ofSeconds(5))
.onErrorResume(Exception.class, ex -> Flux.empty());
}
}
The retry policy uses exponential backoff and filters to retry only on service-unavailable errors, not 404s or 400s. The timeout operator prevents calls from hanging forever, and onErrorResume provides a graceful fallback.
Failure Scenarios
Backpressure Mishandling
Reactive streams carry backpressure — the downstream consumer can signal the upstream to slow down. In theory this prevents memory exhaustion. In practice, if you have a Flux emitting thousands of items per second and your database write cannot keep up, buffering kicks in. With no limit, you run out of memory.
Use .limitRate() or .onBackpressureBuffer() intentionally.
// Buffers everything — dangerous with large datasets
userRepository.findAll()
.map(this::processUser)
.subscribe();
// Backpressure applied — processes 100 items at a time
userRepository.findAll()
.limitRate(100)
.flatMap(this::processUser)
.subscribe();
Blocking Inside Reactive Pipelines
Sometimes a blocking call slips into a map() or flatMap(). The code compiles. It might pass unit tests. In production it deadlocks the event loop thread and your server stops handling requests entirely.
// This compiles but will deadlock your server under load
userRepository.findAll()
.map(user -> someBlockingLegacyService.getInfo(user.getId()))
.subscribe();
The fix is moving the blocking call to a separate scheduler.
userRepository.findAll()
.flatMap(user -> Mono.fromCallable(() ->
someBlockingLegacyService.getInfo(user.getId()))
.subscribeOn(Schedulers.boundedElastic()))
.subscribe();
boundedElastic is the scheduler for blocking calls. It manages a bounded thread pool for exactly this situation. Do not use it on hot paths, but it is the right escape hatch when you need it.
Null in Reactive Streams
Reactive streams do not allow null values. Emitting null is treated the same as completing the stream with an error. If your method might return null, wrap it.
// Throws NullPointerException at runtime
userRepository.findById(999L)
.map(user -> user.getNickname())
.subscribe();
// Wrap nullable values in Optional
userRepository.findById(999L)
.map(user -> Optional.ofNullable(user.getNickname()))
.subscribe();
Memory Leaks from Unsubscribed Pipelines
If you create a Flux from a hot source — a message queue, a Kafka consumer — and never subscribe, the source keeps buffering messages. Always track your subscriptions and cancel them when the owning component is destroyed.
private Flux<Message> messageStream;
@PostConstruct
public void init() {
messageStream = messageBroker.subscribe("topic").share();
}
@Override
public void dispose() {
messageStream.subscribe().dispose();
}
Spring’s @PreDestroy and DisposableBean patterns work here. Make sure reactive components clean up after themselves.
Observability Checklist
Reactive streams can hide latency problems that are obvious in blocking code. A slow call in a blocking system shows up immediately in the thread stack. With reactive, you need explicit instrumentation.
Logging
Project Reactor provides a debugging helper that instruments operators with checkpoint descriptions.
// Enable at startup to instrument all operators
Hooks.enableRecycledChecks();
userRepository.findAll()
.map(user -> processUser(user))
.checkpoint("UserProcessing.fetchAndProcess")
.subscribe();
For production, use ContextualLogging via the log() operator, which attaches context across async boundaries.
userRepository.findAll()
.flatMap(user -> processUser(user))
.contextWrite(ctx -> ctx.put("requestId", requestId))
.log("com.example.user-service")
.subscribe();
Metrics
Micrometer integrates with Project Reactor through ReactorMetrics. Register the meter registry and reactive pipelines automatically report subscription timings.
@Bean
public ReactiveWebServerFactoryAccountant reactorMetrics(MeterRegistry registry) {
return new ReactorNettyMeterBinder(registry);
}
Metrics worth watching:
reactor.subscribed— active subscriptionsreactor.flowed— items that passed through the pipelinereactor cancelled— subscriptions that cancelled early- HTTP client metrics: connect time, response time, retry count
Security Notes
Security in WebFlux uses a different model than Spring Security’s servlet filter chain. WebFlux uses WebFilter chains composed of ServerWebFilter instances.
Basic Authentication with WebFilter
@Configuration
@EnableWebFluxSecurity
public class SecurityConfig {
@Bean
public SecurityWebFilterChain securityWebFilterChain(ServerHttpSecurity http) {
return http
.csrf(ServerHttpSecurity.CsrfSpec::disable)
.authorizeExchange(exchanges -> exchanges
.pathMatchers("/public/**").permitAll()
.pathMatchers("/api/admin/**").hasRole("ADMIN")
.anyExchange().authenticated())
.httpBasic(Customizer.withDefaults())
.build();
}
}
JWT in WebFlux
JWT authentication in reactive context requires reactive token validation.
@Component
public class JwtAuthenticationWebFilter implements WebFilter {
private final JwtTokenProvider jwtTokenProvider;
public JwtAuthenticationWebFilter(JwtTokenProvider jwtTokenProvider) {
this.jwtTokenProvider = jwtTokenProvider;
}
@Override
public Mono<Void> filter(ServerWebExchange exchange, WebFilterChain chain) {
String path = exchange.getRequest().getPath().value();
if (path.startsWith("/public")) {
return chain.filter(exchange);
}
return jwtTokenProvider.validateTokenAsync(
exchange.getRequest().getHeaders().getFirst("Authorization"))
.flatMap(claims -> {
if (claims != null) {
String username = claims.getSubject();
List<SimpleGrantedAuthority> authorities = extractRoles(claims);
return chain.filter(
exchange.mutate()
.principal(new UsernamePasswordAuthenticationToken(
username, null, authorities))
.build());
}
return chain.filter(exchange);
})
.switchIfEmpty(chain.filter(exchange));
}
}
Rate Limiting
Rate limiting works as a WebFilter that tracks request counts per IP using a concurrent map.
@Component
public class RateLimitingWebFilter implements WebFilter {
private final Map<String, AtomicInteger> requestCounts = new ConcurrentHashMap<>();
private static final int MAX_REQUESTS = 100;
@Override
public Mono<Void> filter(ServerWebExchange exchange, WebFilterChain chain) {
String clientIp = exchange.getRequest().getRemoteAddress()
.getAddress().getHostAddress();
AtomicInteger count = requestCounts.computeIfAbsent(clientIp, k -> new AtomicInteger(0));
if (count.incrementAndGet() > MAX_REQUESTS) {
exchange.getResponse().setStatusCode(HttpStatus.TOO_MANY_REQUESTS);
return exchange.getResponse().setComplete();
}
return chain.filter(exchange)
.doFinally(signal -> count.decrementAndGet());
}
}
Common Pitfalls / Anti-Patterns
Block-Hunting
Never call block() or blockFirst() on a reactive pipeline from an event loop thread. If you need the result synchronously at some boundary, use toFuture() at the edge of your system, not in the middle of a reactive chain.
// WRONG — blocks the event loop thread
@GetMapping("/sync")
public UserDto getUser(Long id) {
return userRepository.findById(id).block();
}
// CORRECT — return reactive type from controller
@GetMapping("/async")
public Mono<ResponseEntity<UserDto>> getUser(Long id) {
return userRepository.findById(id)
.map(user -> ResponseEntity.ok(user.toDto()))
.defaultIfEmpty(ResponseEntity.notFound().build());
}
subscribe() Variants
The subscribe() method has several overloads. The no-arg variant swallows errors silently.
// Silent failures — errors go nowhere
userRepository.findAll().subscribe();
// Handle each signal explicitly
userRepository.findAll()
.subscribe(
item -> processItem(item),
error -> logger.error("Failed", error),
() -> logger.info("Completed"));
In tests, always use StepVerifier.
@Test
public void testGetUser() {
StepVerifier.create(userRepository.findById(1L))
.expectNextMatches(user -> user.getName().equals("Alice"))
.verifyComplete();
}
@Test
public void testNotFound() {
StepVerifier.create(userRepository.findById(999L))
.verifyComplete(); // completes without emitting = not found
}
concatMap vs flatMap
Both flatten nested reactive pipelines, but they behave differently.
flatMap subscribes to all inner pipelines immediately and interleaves their outputs as they arrive. Order is not preserved. Use it for independent operations that can run concurrently.
// A, B, C start simultaneously, results interleave as they complete
userRepository.findAll()
.flatMap(user -> externalService.enrich(user))
.subscribe();
concatMap subscribes to the next inner pipeline only after the previous one completes. Order is preserved, but operations run sequentially.
// Processes A completely, then B, then C
userRepository.findAll()
.concatMap(user -> writeToAuditLog(user))
.subscribe();
Mixing these up causes subtle correctness bugs. flatMap for independent concurrent work, concatMap when order matters or operations must not overlap.
Production Failure Scenarios
Reactive systems fail in ways that look nothing like blocking code failures.
Blocking the Event Loop Thread
A blocking call buried inside a reactive pipeline deadlocks the event loop thread. The code compiles, unit tests pass in isolation, and production grinds to a halt. The symptom is sessions timing out across the entire service, not just on the endpoint making the bad call. Common culprits: jdbcTemplate.queryForObject() inside a map(), a synchronous cache lookup inside a flatMap(), any call to a blocking library that has no reactive alternative. Profile blocking calls with Schedulers.assertThatOperatorsAreNotBlocking() in tests. Move blocking operations to Mono.fromCallable() with subscribeOn(Schedulers.boundedElastic()).
null Emission Triggering NPE in Pipelines
Reactive streams do not allow null values. Emitting null is treated the same as completing with an error. If your database query returns no results and you map that to null instead of Mono.empty(), the pipeline throws a NullPointerException at subscription time. This is particularly insidious because the error happens far from where the null originated. Always wrap nullable values in Optional. Use switchIfEmpty() to handle the empty case explicitly rather than relying on null.
Hot Source Memory Leaks from Unsubscribed Streams
A Flux created from a hot source like a message queue or Kafka consumer keeps buffering messages if no subscription exists. If your application receives messages but the consumer is never subscribed, those messages accumulate in the source’s internal buffer until memory is exhausted. Always track subscriptions and cancel them in @PreDestroy or DisposableBean.dispose(). Use share() or publish() with explicit subscription management rather than assuming the pipeline will clean itself up.
Backpressure Mishandling Causing Memory Exhaustion
A Flux emitting thousands of items per second that meets a slow consumer buffers everything by default. Without .limitRate() or explicit backpressure configuration, memory grows until the application crashes or the consumer catches up. The fix is applying .limitRate() or .onBackpressureBuffer() with an explicit overflow strategy. Choose the strategy based on your failure mode: drop new events, drop old events, or signal backpressure to the producer.
Quick Recap Checklist
- Mono emits 0-1 elements; Flux emits 0-n. Pick based on how many results your operation produces.
- Nothing runs until you subscribe. Reactive pipelines are lazy by design.
- Blocking inside a reactive chain deadlocks the event loop. Move blocking calls to
Schedulers.boundedElastic(). - Reactive streams do not allow null. Wrap nullable values in Optional.
- Use
.limitRate()or.onBackpressureBuffer()to prevent memory exhaustion on large streams. - WebClient replaces RestTemplate in reactive applications. Use it for all HTTP calls.
- Attach retry and timeout to every external HTTP call.
.switchIfEmpty()handles the empty Mono case without throwing.flatMapinterleaves concurrent work;concatMappreserves order.- Use StepVerifier in tests. Never use no-arg
subscribe()in production. - WebFlux security uses WebFilter, not servlet filters.
- Cancel hot streams on component destroy to avoid resource leaks.
- Use
checkpoint()orlog()for observability across async boundaries. - Functional routing is worth considering for new WebFlux applications — it is more explicit and easier to test than annotation-based controllers.
Interview Questions
Mono represents a pipeline that emits zero or one element — like an async Optional that can also be empty. Flux represents a pipeline that emits zero to n elements, similar to a reactive Stream. Use Mono for single-result operations like fetching one user. Use Flux when an operation returns multiple results, like querying all users or streaming data from a WebSocket.
WebFlux runs on a small fixed number of event-loop threads, typically one per CPU core. Calling block() on one of those threads freezes it completely — it cannot handle any other request. Under any meaningful load, a handful of blocking calls can saturate all event-loop threads and bring the server to a standstill. The right approach is to stay fully asynchronous: return Mono or Flux from your controller, and only block synchronously at the absolute system boundary if you absolutely must.
flatMap subscribes to all inner pipelines immediately and interleaves their results as they arrive. Order is not preserved, but throughput is maximized because everything runs concurrently. concatMap subscribes to the next inner pipeline only after the current one finishes, so order is preserved but operations run sequentially. Use flatMap for independent concurrent work where order does not matter. Use concatMap when order matters or when operations must not overlap, like audit logging where you need each write to complete before the next starts.
WebClient uses the retryWhen() operator with a Retry policy — typically exponential backoff with a max attempt count and error filtering so you do not retry on client errors like 404s. The timeout() operator sets a deadline for the entire pipeline. Both matter because network calls fail in various ways: transiently under load, permanently when a service goes down. Without retry, a single blip fails the whole request. Without timeout, a hanging call ties up resources indefinitely. Together they give you explicit, configurable client-side resilience.
Backpressure is how a downstream consumer tells an upstream producer to slow down, preventing the producer from overwhelming the consumer with more data than it can process. Without it, a fast producer feeding a slow consumer causes unbounded buffering and memory exhaustion. Project Reactor implements backpressure through the Reactive Streams specification. Operators like limitRate() request items in batches, onBackpressureBuffer() buffers with configurable overflow behavior, and onBackpressureDrop() discards items when downstream cannot keep up. Choosing the right strategy for your data source and processing speed is one of the most important production readiness decisions in a reactive system.
Spring MVC uses a one-thread-per-request model where each incoming request occupies a thread for its entire lifetime, including wait times during I/O operations. A thread consumes roughly 1MB of stack memory. WebFlux uses an event-loop model with a small fixed number of threads (typically one per CPU core) that never block — they handle thousands of concurrent connections by registering callbacks and resuming when data arrives. Memory per connection is a few kilobytes. This difference enables WebFlux to handle an order of magnitude more connections than Spring MVC on the same hardware, though it requires careful code to avoid blocking calls.
switchIfEmpty substitutes an alternative Mono when the source completes without emitting any value. It is the reactive equivalent of providing a default value for an absent result, rather than throwing an exception or returning null. Common use cases include returning a default object when a database lookup finds nothing, or falling back to a cached response when the primary source is empty. It differs from defaultIfEmpty in that defaultIfEmpty always emits a value (the default) regardless of the source, while switchIfEmpty only activates when the source completes empty.
Annotation-based controllers in WebFlux look identical to Spring MVC but return reactive types. They use familiar annotations like @RestController, @GetMapping, and rely on Spring's parameter resolution. Functional routing defines routes as a RouterFunction bean with Route.route() DSL, where each route points to a handler method. The handler receives a ServerRequest and returns a ServerResponse, with all parsing done explicitly in the handler. Functional routing is more explicit, easier to unit test in isolation, and scales better for large route sets, but has a steeper initial learning curve. Annotation-based is easier to migrate from MVC but handler logic can become mixed with routing concerns.
Error handling in reactive pipelines uses operators rather than try-catch blocks. onErrorResume() catches an error and replaces it with an alternative Mono or Flux — use this for fallback values or alternative data sources. onErrorReturn() catches and returns a static default value. doOnError() lets you log or monitor errors without altering the stream. retry() resubscribes to the source when an error occurs, useful for transient failures. For WebClient responses, onStatus() lets you map HTTP error codes to custom exceptions, which onErrorResume then catches. Always handle errors explicitly rather than letting them propagate to subscribers by default.
R2DBC (Reactive Relational Database Connectivity) is a specification for reactive database access that provides non-blocking database drivers for relational databases like PostgreSQL, MySQL, and H2. Traditional JDBC is inherently blocking — each operation occupies a thread waiting for the database. R2DBC drivers return Mono or Flux from database operations, keeping the thread free to handle other requests. This matters because if you use a blocking JDBC driver inside a WebFlux application, you defeat the purpose of the event-loop model. Spring Data R2DBC and DatabaseClient are the main Spring abstractions over R2DBC drivers. Without a reactive database driver, your Spring WebFlux application will still block at the persistence layer.
Spring MVC security runs as a servlet filter chain that processes HttpServletRequest and HttpServletResponse. WebFlux security runs as a WebFilter chain processing ServerHttpRequest and ServerHttpResponse — these are reactive, non-blocking types. Authentication in WebFlux returns Mono<Authentication>, and the security chain passes ServerWebExchange through each filter. The configuration DSL uses ServerHttpSecurity instead of HttpSecurity. JWT validation in WebFlux uses ReactiveJwtDecoder which returns a Mono<Jwt>, fitting naturally into the reactive pipeline. The mental shift is that everything is async and returns Mono or Flux rather than blocking on security operations.
Schedulers.boundedElastic() is a dedicated thread pool in Project Reactor designed for wrapping blocking calls that you cannot avoid. It creates a bounded pool of threads (by default one per 250ms of blocking time, capped at a reasonable maximum) that execute blocking work without interfering with the event-loop threads. Use subscribeOn(Schedulers.boundedElastic()) or Mono.fromCallable() with subscribeOn() when you must call a blocking legacy library, JDBC code, or synchronous HTTP client from within a reactive pipeline. The caveat is that boundedElastic threads still block, so avoid using them on hot paths or in high-throughput scenarios — they are an escape hatch, not a default solution.
StepVerifier is Project Reactor's testing tool that subscribes to a Mono or Flux and verifies the signals it receives. Use StepVerifier.create(pipeline) then chain expectations: expectNext() for expected items, expectNextMatches() with a predicate for conditional matching, expectError() or expectErrorMatches() for expected errors, verifyComplete() to assert the stream ended successfully, and verifyTimeout() to set a max duration. For example, StepVerifier.create(userRepository.findAll()) with .expectNextCount(5) and .verifyComplete() asserts a flux emits 5 items then completes. This approach is far superior to using subscribe() in tests because it gives precise control over what signals you expect and when.
A cold source produces data from the beginning for each subscriber — each subscription gets the full sequence from start to end. HTTP requests, database queries, and file reads are typical cold sources. A hot source does not reset for new subscribers — it broadcasts data to all current subscribers from the point they subscribed forward. Kafka consumers, message queues, and UI event streams are typical hot sources. Hot sources can cause memory leaks if they buffer data for subscribers that never arrive, because the source does not know when it is safe to drop old events. Use share() or publish().refCount() to convert a hot source into one that auto-cleans when all subscribers leave, and always clean up subscriptions in @PreDestroy.
Mono.zip() and Flux.zip() combine multiple reactive pipelines pairwise, emitting a combined result only when all sources have each produced a new item. For Mono, zip emits when all Monos have resolved — useful for fan-in scenarios where you need multiple data sources before proceeding, like fetching a user and their permissions concurrently and combining them. For Flux, zip emits tuples of corresponding elements, so it pairs item 1 from each source, then item 2, and so on. This is useful for joining streams where you need correlated data from multiple sources at the same logical position. The downstream waits for all upstreams to be ready before emitting, so slow sources throttle the entire zip.
WebClient is fully asynchronous and non-blocking, whereas RestTemplate is synchronous and blocking by design. WebClient uses an event-loop model (Netty by default) that can handle many concurrent HTTP calls with a small number of threads, while RestTemplate occupies a thread for the full duration of each HTTP call. WebClient has a fluent API with better support for reactive pipeline composition, integrates naturally with retry and timeout operators, and supports HTTP/2. RestTemplate is legacy and Spring deprecated it in favor of WebClient. For Spring WebFlux, WebClient is the only sensible choice. For Spring MVC, WebClient is still preferable for high-concurrency scenarios, though RestTemplate remains acceptable for simple low-concurrency use cases where asynchronicity is not a concern.
Reactor Context is a reactive equivalent of ThreadLocal — it carries state across async boundaries in a reactive pipeline. Unlike ThreadLocal which is thread-bound, Context propagates through flatMap, Mono.fromRunnable(), and other async boundary operators. You write to the Context with contextWrite() and read from it with contextRead(). Common uses include request-scoped values like a request ID or tenant ID that need to travel through the entire pipeline. Context is immutable and each operator can add to it, but existing keys cannot be overwritten. It is essential for proper logging correlation across async boundaries because log() can read from Context to include request IDs in log entries.
Project Reactor uses Scheduler instances to determine which thread executes which part of a pipeline. publishOn() changes the threading context for subsequent operators — everything after it runs on the specified scheduler. subscribeOn() changes the thread where subscription starts, affecting the entire pipeline upstream. Default behavior runs everything on the subscriber's thread. WebFlux uses Schedulers.parallel() for event-loop work (do not block these) and Schedulers.boundedElastic() for blocking work. Schedulers.single() provides a single reusable thread useful for sequential operations. Understanding scheduler semantics is critical because the most common production bug in WebFlux is blocking a parallel/event-loop scheduler, which deadlocks the server.
onErrorReturn() catches a specific error type and replaces it with a static fallback value — the simplest error handling. onErrorResume() catches an error and switches to an alternative Mono or Flux as the fallback — more flexible because the fallback can itself be async. onErrorMap() catches an error and transforms it into a different exception before passing it downstream — useful for converting low-level exceptions into domain-specific ones. The choice depends on what you need: return a constant default (onErrorReturn), compute a fallback reactively (onErrorResume), or translate the error type (onErrorMap). All three terminate the error signal and replace it with a success path.
Choose WebFlux when you need high concurrency with limited hardware — microservices that fan out to many downstream services, chat servers, live dashboards, event streaming endpoints, or anything handling thousands of concurrent long-lived connections. WebFlux excels at I/O-bound workloads where requests spend most of their time waiting. Choose Spring MVC when your team is new to reactive programming, when your dependencies lack reactive drivers (especially relational databases without R2DBC), when you need blocking libraries that cannot be easily replaced, or when your problem is CPU-bound rather than I/O-bound. Spring MVC is simpler to debug, has a larger ecosystem, and benefits from universal JDBC driver support. A microservices architecture that mixes MVC and WebFlux services is common and often the right call.
Further Reading
R2DBC Deep Dive
R2DBC (Reactive Relational Database Connectivity) is the cornerstone of fully reactive Spring applications with relational databases. Unlike JDBC’s synchronous model, R2DBC exposes database operations as Mono and Flux types, enabling true end-to-end non-blocking data access. Spring Data R2DBC provides repository abstractions that return reactive types, and the DatabaseClient gives you fine-grained control over queries and transactions. R2DBC also introduces backpressure at the database layer, which means slow consumers can signal the database to slow down, preventing unbounded memory growth. Pair R2DBC with reactive transaction management via @Transactional for consistent data access across multiple operations.
WebSocket & Server-Sent Events (SSE) in WebFlux
WebFlux has first-class support for WebSocket connections and Server-Sent Events, both of which map naturally to reactive streams. A WebSocket endpoint returns a Flux<WebSocketMessage> in each direction, letting you process incoming and outgoing messages reactively. SSE endpoints return a Flux<ServerSentEvent> which WebFlux encodes automatically. Both work seamlessly with backpressure — if a client reads slowly, the server buffers within configured limits rather than accumulating unbounded data. For real-time dashboards or notification systems, WebSocket and SSE in WebFlux provide better resource efficiency than polling or long-polling alternatives.
Testing WebFlux Applications
Testing reactive code requires different tools than testing blocking code. StepVerifier is the primary tool — it subscribes to a reactive pipeline and asserts each expected signal. For unit testing controllers, WebTestClient binds to your router or controller and lets you assert the response with StepVerifier. For integration testing, MockServerHttpRequest and MockServerHttpResponse simulate HTTP interactions without a real server. Always test error paths: use verifyError() or verifyErrorMatches() in StepVerifier. For performance testing, Project Reactor provides FluxSink.next() backpressure tests. Profile blocking calls in tests with Schedulers.assertThatOperatorsAreNotBlocking() to catch accidental blocking calls before they reach production.
Conclusion
Spring WebFlux represents a fundamental shift in how Java applications handle concurrency. Rather than dedicating a thread to each blocking I/O operation, WebFlux uses an event-loop model that can handle thousands of concurrent connections with a small number of threads. The payoff is dramatic throughput improvements for I/O-bound workloads, but the cost is a steeper learning curve and debugging experience that differs significantly from traditional blocking code.
The key to successful WebFlux adoption is recognizing when the tradeoff makes sense. High-concurrency scenarios with multiple external service calls, streaming data endpoints, and applications that need to handle thousands of concurrent long-lived connections are WebFlux’s natural habitat. Simple CRUD operations with no meaningful concurrency needs are better served by Spring MVC. The ecosystem matters too: without reactive database drivers, you inherit blocking at the persistence layer and lose the core benefit.
If you do adopt WebFlux, invest heavily in the team’s operator fluency and testing practices. Reactive pipelines that compile and pass unit tests can still deadlock production servers when a blocking call slips into an operator chain. Schedulers.assertThatOperatorsAreNotBlocking() in tests catches these before they reach production. The observability story requires explicit instrumentation too — where blocking code shows up in thread stacks naturally, reactive code needs checkpoint() and Micrometer integration to be traceable.
Category
Related Posts
Spring Boot Build Tools: Maven & Gradle
Configure Maven and Gradle for Spring Boot projects—plugins, dependency management, packaging JARs and WARs, and build automation essentials.
Embedded Web Servers in Spring Boot: Tomcat, Jetty, Undertow
Configure embedded servers in Spring Boot: compare Tomcat, Jetty, and Undertow, tune thread pools, enable access logs, and switch implementations.
JUnit 5 & Jupiter: Lifecycle, Nested & Parameterized Tests
Explore JUnit 5 Jupiter features: master test lifecycle annotations, organize tests with @Nested, and parameterize tests with @CsvSource and @MethodSource.