- Agile 2
- Alta disponibilidad 1
- Alternativas cloud 1
- Aop 1
- Arquitectura 3
- Arquitectura distribuida 2
- Automatizacion 3
- Azure devops 1
- Base de datos 1
- Buenas practicas 19
- Cloud 1
- Colas 7
- Competing consumers 1
- Convenciones 11
- Copilot 1
- Diseno 6
- Docker 2
- Docker compose 1
- Documentacion 1
- Eda 11
- Equipos 1
- Escalabilidad 1
- Flujo de negocio 1
- Flujo de trabajo 3
- Flyway 1
- Git 4
- Gradle 3
- Herramientas digitales 1
- Ia 1
- Iam 1
- Infraestructura 2
- Java 14
- Jerarquia tecnica 1
- Jpa 1
- Jsonb 1
- Kafka 7
- Kubernetes 1
- Liderazgo en software 1
- Lineamientos 1
- Log 1
- Logging 3
- Microservicios 3
- Mongodb 1
- Monitoreo 1
- Nosql 3
- Observabilidad 4
- Open source 1
- Plugins 3
- Postgresql 1
- Privacidad 1
- Programacion funcional 1
- Programacion reactiva 4
- Rabbitmq 6
- Rotacion de talento 1
- Saga 2
- Scrum 2
- Security 1
- Seguridad 1
- Self hosting 1
- Sistemas legados 1
- Spring boot 3
- Spring mvc 2
- Sql 3
- Streams 1
- Threadlocal 1
- Trazabilidad 2
- Versionado 2
- Web 1
- Webflux 2
- Websockets 1
- Zero trust 1
Programacion reactiva
4 artículos
Spring WebFlux 4: Comunicación Avanzada, Pruebas y Producción
- Mauricio ECR
- Arquitectura
- 31 May, 2025
La serie Spring WebFlux nos ha llevado a través de un viaje fascinante por el mundo de la programación reactiva, desde sus fundamentos y el poder de Project Reactor hasta la construcción de arquit
Spring WebFlux 4: Comunicación Avanzada, Pruebas y Producción
- Mauricio ECR
- Arquitectura
- 31 May, 2025
La serie Spring WebFlux nos ha llevado a través de un viaje fascinante por el mundo de la programación reactiva, desde sus fundamentos y el poder de Project Reactor hasta la construcción de arquitecturas altamente concurrentes y la gestión de la comunicación con servicios externos y bases de datos. En esta cuarta parte, profundizaremos en aspectos más avanzados y críticos para el desarrollo y despliegue de aplicaciones WebFlux robustas y eficientes. Exploraremos desde la comunicación en tiempo real con Server-Sent Events y WebSockets, hasta la crucial gestión de la contrapresión, el contexto reactivo, las estrategias de testing y, por supuesto, la seguridad y las buenas prácticas en producción.
1. Server-Sent Events (SSE): Flujos de Eventos Unidireccionales
Los Server-Sent Events (SSE) son una tecnología web que permite a un servidor enviar actualizaciones automáticamente a un cliente a través de una conexión HTTP persistente y unidireccional. A diferencia de los WebSockets, que son bidireccionales y más complejos, los SSE están diseñados específicamente para escenarios donde el cliente solo necesita recibir datos del servidor. Piensa en ellos como un flujo continuo de noticias, actualizaciones de cotizaciones bursátiles o notificaciones en tiempo real.
¿Cómo funcionan los SSE?
El cliente establece una conexión HTTP normal con el servidor. Sin embargo, en lugar de cerrar la conexión después de enviar la respuesta inicial, el servidor la mantiene abierta y envía datos de forma continua. Cada "evento" se envía como un bloque de texto formateado de una manera específica, seguido de un salto de línea. El navegador o cliente (usando la API EventSource de JavaScript) interpreta estos bloques como eventos individuales.
SSE con Spring WebFlux
En Spring WebFlux, implementar SSE es sorprendentemente sencillo gracias a la naturaleza reactiva de Flux. Dado que un Flux puede emitir 0 a N elementos de forma asíncrona, es la elección natural para representar un flujo de eventos.
Para enviar eventos, simplemente necesitas devolver un Flux desde tu controlador. Spring WebFlux se encargará automáticamente de configurar los encabezados HTTP (Content-Type: text/event-stream) y formatear los datos para que el cliente los reciba como SSE.
import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import reactor.core.publisher.Flux;
import java.time.Duration;
import java.time.LocalDateTime;
@RestController
public class SseController {
@GetMapping(value = "/eventos", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> getEvents() {
return Flux.interval(Duration.ofSeconds(1)) // Emite un elemento cada segundo
.map(sequence -> "Evento #" + sequence + " a las " + LocalDateTime.now());
}
@GetMapping(value = "/data-stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<MyData> streamMyData() {
return Flux.interval(Duration.ofSeconds(2))
.map(sequence -> new MyData("Item " + sequence, Math.random() * 100))
.take(5); // Limita el número de elementos
}
}
En este ejemplo:
getEvents()envía una cadena de texto cada segundo.streamMyData()envía objetosMyData(que se serializarán a JSON automáticamente) cada dos segundos, limitando la emisión a 5 elementos.
Del lado del cliente (JavaScript):
const eventSource = new EventSource('/eventos');
eventSource.onmessage = function(event) {
console.log("Mensaje recibido:", event.data);
};
eventSource.onerror = function(error) {
console.error("Error en el flujo de eventos:", error);
eventSource.close();
};
// Si el servidor envía eventos con un 'event' type específico:
// eventSource.addEventListener('nombreDeEvento', function(event) {
// console.log("Evento con nombre específico:", event.data);
// });
Los SSE son ideales para dashboards en tiempo real, feeds de actividad o cualquier escenario donde se necesiten actualizaciones push del servidor sin la complejidad de una conexión bidireccional completa.
2. Backpressure: Gestionando el Flujo de Datos
El concepto de backpressure (contrapresión) es fundamental en la programación reactiva y, en particular, en Project Reactor y Spring WebFlux. Se refiere a la capacidad de un suscriptor (consumidor) de señalar a un publicador (productor) qué tan rápido o cuántos elementos puede procesar. En un flujo reactivo, si el productor es mucho más rápido que el consumidor, los datos se acumularán en el buffer del consumidor, lo que puede llevar a problemas de memoria o a la caída del sistema. La contrapresión resuelve esto permitiendo que el consumidor "tire" de los datos solo cuando está listo para manejarlos.
¿Por qué es crucial la contrapresión?
Imagina un río (el publicador) que fluye muy rápido hacia un balde (el suscriptor) que solo puede contener una pequeña cantidad de agua a la vez. Sin contrapresión, el balde se desbordaría rápidamente. Con contrapresión, el balde puede indicarle al río que disminuya el caudal o que le envíe agua solo cuando haya espacio.
En el contexto de Spring WebFlux, la contrapresión es vital para la estabilidad y eficiencia del sistema. Evita que un servicio backend sobrecargue a un cliente más lento (como un navegador o una API externa con límites de tasa) o que una base de datos reactiva inunde el servicio con resultados que no puede procesar a tiempo.
Implementación en Reactor
Project Reactor implementa la contrapresión según las especificaciones de Reactive Streams. Esto significa que los operadores de Mono y Flux manejan la contrapresión de forma nativa. Cuando un Subscriber se suscribe a un Publisher, lo primero que hace es solicitar un número inicial de elementos. Luego, a medida que procesa esos elementos, puede solicitar más (request(n)).
import reactor.core.publisher.Flux;
import org.reactivestreams.Subscription;
import org.reactivestreams.Subscriber;
public class BackpressureExample {
public static void main(String[] args) {
Flux.range(1, 100) // Publicador que emite 100 números
.subscribe(new Subscriber<Integer>() {
private Subscription s;
private int count = 0;
@Override
public void onSubscribe(Subscription s) {
this.s = s;
System.out.println("Suscrito. Solicitando 2 elementos.");
s.request(2); // Solicita inicialmente 2 elementos
}
@Override
public void onNext(Integer integer) {
System.out.println("Procesando: " + integer);
count++;
if (count % 2 == 0) { // Después de procesar 2 elementos, solicita 2 más
System.out.println("Procesados 2. Solicitando 2 más.");
s.request(2);
}
}
@Override
public void onError(Throwable t) {
System.err.println("Error: " + t);
}
@Override
public void onComplete() {
System.out.println("Completado.");
}
});
}
}
En este ejemplo simplificado, el Subscriber controla la velocidad de emisión al solicitar solo dos elementos a la vez. Este mecanismo es transparente en la mayoría de los casos cuando usas operadores de Reactor, pero es crucial entender que está ocurriendo "bajo el capó" para un comportamiento predecible y robusto.
3. Contexto Reactivo: Compartiendo Información
En la programación tradicional, ThreadLocal se utiliza comúnmente para compartir información a través de diferentes métodos en el mismo hilo de ejecución, como el contexto de seguridad o un ID de correlación para logging. Sin embargo, en un entorno reactivo y no bloqueante como Spring WebFlux, donde las operaciones pueden cambiar de hilo de forma asíncrona, ThreadLocal ya no es una opción viable porque la información se perdería entre los cambios de hilo.
Aquí es donde entra el Contexto Reactivo (Context) de Project Reactor. El Context es una característica que permite adjuntar datos a un flujo reactivo, haciéndolos disponibles para cualquier operador o suscriptor a lo largo de la cadena, independientemente de qué hilo esté ejecutando la operación.
¿Cómo funciona el Contexto Reactivo?
Cada flujo Mono o Flux tiene asociado un Context. Este Context es una estructura de datos inmutable (similar a un Map) que se propaga a lo largo de la cadena de operadores. Cuando un operador necesita acceder a información del contexto, puede hacerlo a través de métodos como contextWrite().
import reactor.core.publisher.Mono;
import reactor.core.publisher.Flux;
import reactor.util.context.Context;
public class ReactiveContextExample {
public static void main(String[] args) {
String correlationId = "corr-123";
Mono<String> dataMono = Mono.just("Hello")
.doOnNext(s -> {
// Acceder al contexto para obtener el correlationId
Mono.deferContextual(ctx -> {
String id = ctx.get("correlationId");
System.out.println("doOnNext: Data = " + s + ", Correlation ID from Context = " + id);
return Mono.empty();
}).subscribe(); // Suscribirse para activar el deferContextual
})
.contextWrite(Context.of("correlationId", correlationId)); // Escribir en el contexto
dataMono.subscribe(
data -> System.out.println("Subscriber: Data = " + data),
error -> System.err.println("Subscriber Error: " + error),
() -> System.out.println("Subscriber: Completed")
);
System.out.println("\n--- Otro ejemplo con Flux y múltiples valores ---");
Flux.just("Item A", "Item B")
.contextWrite(Context.of("traceId", "trace-xyz")) // Escribir en el contexto
.flatMap(item ->
Mono.deferContextual(ctx -> {
String traceId = ctx.get("traceId");
return Mono.just("Procesando " + item + " con Trace ID: " + traceId);
})
)
.subscribe(
result -> System.out.println("Subscriber: " + result),
error -> System.err.println("Subscriber Error: " + error)
);
}
}
Aplicaciones Comunes del Contexto Reactivo
- Propagación de IDs de Correlación/Traza: Esencial para el logging distribuido y la observabilidad. Puedes insertar un ID de correlación al inicio del flujo y que esté disponible en cada operador y en la capa de persistencia.
- Contexto de Seguridad: Información del usuario autenticado, roles, permisos.
- Parámetros de Configuración Dinámicos: Valores que pueden variar por solicitud pero que no son parte de la carga útil principal.
- Datos Transaccionales: Si bien Spring Data R2DBC maneja las transacciones reactivas, el contexto podría usarse para almacenar metadatos relacionados con la transacción.
El Context proporciona una forma segura y reactiva de pasar información a través de los límites de los hilos, manteniendo la integridad del flujo de datos.
4. Testing Reactivo: Garantizando la Robustez
Probar aplicaciones reactivas requiere un enfoque ligeramente diferente al de las aplicaciones síncronas debido a la naturaleza asíncrona y no bloqueante de los flujos. Spring WebFlux y Project Reactor ofrecen herramientas poderosas para facilitar este proceso, asegurando que tus flujos de datos se comporten como esperas.
TestUtils de Reactor: StepVerifier
La herramienta más importante para probar flujos Mono y Flux es StepVerifier de Project Reactor. Permite probar secuencias reactivas de manera determinista, verificando los valores emitidos, los errores y la finalización, e incluso simulando el tiempo para probar operadores basados en tiempo.
import org.junit.jupiter.api.Test;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
import java.time.Duration;
class ReactiveTestingExample {
// Prueba de un Mono simple
@Test
void testMono() {
Mono<String> mono = Mono.just("Hello Reactive!");
StepVerifier.create(mono)
.expectNext("Hello Reactive!") // Espera un valor específico
.expectComplete() // Espera que el flujo se complete
.verify(); // Inicia la verificación
}
// Prueba de un Flux con múltiples elementos
@Test
void testFlux() {
Flux<Integer> flux = Flux.just(1, 2, 3);
StepVerifier.create(flux)
.expectNext(1)
.expectNext(2)
.expectNext(3)
.expectComplete()
.verify();
}
// Prueba de un Flux con un error
@Test
void testFluxWithError() {
Flux<String> flux = Flux.just("data1", "data2")
.concatWith(Mono.error(new RuntimeException("Oops!")));
StepVerifier.create(flux)
.expectNext("data1", "data2")
.expectError(RuntimeException.class) // Espera un error de tipo RuntimeException
.verify();
}
// Prueba de un Flux con retardo (simulando tiempo)
@Test
void testFluxWithDelay() {
Flux<Long> flux = Flux.interval(Duration.ofSeconds(1)).take(3);
StepVerifier.withVirtualTime(() -> flux) // Usa tiempo virtual para acelerar la prueba
.expectSubscription()
.expectNoEvent(Duration.ofSeconds(1)) // No espera eventos por 1 segundo
.expectNext(0L)
.thenAwait(Duration.ofSeconds(1)) // Avanza el tiempo virtual 1 segundo
.expectNext(1L)
.thenAwait(Duration.ofSeconds(1))
.expectNext(2L)
.expectComplete()
.verify();
}
}
StepVerifier ofrece una API fluida y encadenable para definir las expectativas sobre el flujo. withVirtualTime() es particularmente útil para probar operadores basados en tiempo sin tener que esperar el tiempo real, acelerando significativamente las pruebas.
Testing de Controladores WebFlux
Para probar controladores WebFlux, puedes usar WebTestClient. Este cliente no bloqueante permite realizar solicitudes HTTP simuladas a tu aplicación WebFlux y verificar las respuestas reactivas. Es ideal para pruebas de integración o de slice.
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.autoconfigure.web.reactive.WebFluxTest;
import org.springframework.boot.test.mock.mockito.MockBean;
import org.springframework.test.web.reactive.server.WebTestClient;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import static org.mockito.Mockito.when;
@WebFluxTest(MyReactiveController.class) // Especifica el controlador a probar
class MyReactiveControllerTest {
@Autowired
private WebTestClient webTestClient; // Cliente para realizar solicitudes HTTP
@MockBean // Simula dependencias del controlador
private MyReactiveService myReactiveService;
@Test
void testGetHello() {
when(myReactiveService.getHelloMessage()).thenReturn(Mono.just("Hello from Service!"));
webTestClient.get().uri("/hello")
.exchange() // Realiza la solicitud
.expectStatus().isOk() // Verifica el código de estado HTTP
.expectBody(String.class).isEqualTo("Hello from Service!"); // Verifica el cuerpo de la respuesta
}
@Test
void testGetAllItems() {
when(myReactiveService.getAllItems()).thenReturn(Flux.just("Item1", "Item2"));
webTestClient.get().uri("/items")
.exchange()
.expectStatus().isOk()
.expectBodyList(String.class).containsExactly("Item1", "Item2"); // Verifica una lista de elementos
}
}
En este ejemplo:
@WebFluxTestconfigura un contexto de aplicación limitado para probar solo el controlador especificado.@MockBeanpermite simular las dependencias del controlador, lo que es crucial para aislar la lógica del controlador.WebTestClientsimula las solicitudes HTTP y permite verificar la respuesta de manera reactiva.
Combinando StepVerifier para la lógica reactiva de negocio y WebTestClient para las interacciones HTTP, puedes construir un conjunto de pruebas robusto para tus aplicaciones Spring WebFlux.
5. Seguridad en Aplicaciones WebFlux (Spring Security Reactivo)
La seguridad es un pilar fundamental en cualquier aplicación, y las aplicaciones reactivas no son la excepción. Spring Security Reactivo proporciona una integración fluida con Spring WebFlux, ofreciendo un modelo de seguridad no bloqueante que se adapta perfectamente al paradigma reactivo. A diferencia del Spring Security tradicional, que se basa en ThreadLocal y filtros de Servlet, la versión reactiva opera con Mono y Flux para mantener la reactividad de principio a fin.
Componentes Clave de Spring Security Reactivo
SecurityWebFilterChain: Reemplaza alFilterChainde Servlets y define la cadena de filtros de seguridad reactivos.ReactiveUserDetailsService: Para cargar detalles del usuario de forma reactiva.ReactiveAuthenticationManager: Para autenticar usuarios de forma reactiva.SecurityContextRepository: Para guardar y cargar el contexto de seguridad (ej. para sesiones o JWT).
Configuración Básica
Para habilitar Spring Security Reactivo, necesitas añadir la dependencia spring-boot-starter-security y configurar tu SecurityWebFilterChain.
// build.gradle
dependencies {
implementation 'org.springframework.boot:spring-boot-starter-webflux'
implementation 'org.springframework.boot:spring-boot-starter-security'
// ...
}
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.security.config.annotation.web.reactive.EnableWebFluxSecurity;
import org.springframework.security.config.web.server.ServerHttpSecurity;
import org.springframework.security.core.userdetails.MapReactiveUserDetailsService;
import org.springframework.security.core.userdetails.User;
import org.springframework.security.core.userdetails.UserDetails;
import org.springframework.security.crypto.bcrypt.BCryptPasswordEncoder;
import org.springframework.security.crypto.password.PasswordEncoder;
import org.springframework.security.web.server.SecurityWebFilterChain;
@Configuration
@EnableWebFluxSecurity
public class SecurityConfig {
@Bean
public SecurityWebFilterChain securityWebFilterChain(ServerHttpSecurity http) {
return http
.csrf(ServerHttpSecurity.CsrfSpec::disable) // Deshabilita CSRF para APIs sin estado
.authorizeExchange(exchanges -> exchanges
.pathMatchers("/public/**").permitAll() // Rutas públicas accesibles sin autenticación
.pathMatchers("/admin/**").hasRole("ADMIN") // Rutas solo para ADMIN
.anyExchange().authenticated() // Todas las demás rutas requieren autenticación
)
.httpBasic(httpBasic -> httpBasic.init(http)) // Habilita autenticación HTTP Basic
.formLogin(formLogin -> formLogin.disable()) // Deshabilita el formulario de login por defecto
.build();
}
@Bean
public MapReactiveUserDetailsService userDetailsService(PasswordEncoder passwordEncoder) {
UserDetails user = User.withUsername("user")
.password(passwordEncoder.encode("password"))
.roles("USER")
.build();
UserDetails admin = User.withUsername("admin")
.password(passwordEncoder.encode("adminpass"))
.roles("ADMIN")
.build();
return new MapReactiveUserDetailsService(user, admin);
}
@Bean
public PasswordEncoder passwordEncoder() {
return new BCryptPasswordEncoder();
}
}
En este ejemplo:
- Deshabilitamos CSRF (común para APIs RESTful sin estado).
- Definimos reglas de autorización para diferentes rutas (
/publices accesible por todos,/adminsolo por usuarios con rolADMIN). - Configuramos
HTTP Basicpara la autenticación simple. - Se define un
MapReactiveUserDetailsServicepara usuarios en memoria, aunque en un entorno real se usaría una base de datos reactiva.
Accediendo al Usuario Autenticado
En Spring WebFlux, puedes acceder al usuario autenticado usando Mono<Principal> o Mono<Authentication> en tus controladores o servicios.
import org.springframework.security.core.Authentication;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import reactor.core.publisher.Mono;
import java.security.Principal;
@RestController
public class SecuredController {
@GetMapping("/secure/user-info")
public Mono<String> getUserInfo(Mono<Principal> principalMono) {
return principalMono.map(principal -> "Hola, " + principal.getName() + "! Eres un usuario autenticado.");
}
@GetMapping("/admin/dashboard")
public Mono<String> getAdminDashboard(Mono<Authentication> authenticationMono) {
return authenticationMono.map(auth -> "Bienvenido al Dashboard de Admin, " + auth.getName() + "! Roles: " + auth.getAuthorities());
}
}
Uso de JWT (JSON Web Tokens)
Para aplicaciones sin estado, el uso de JWT es una práctica común. Spring Security Reactivo facilita la implementación de autenticación basada en JWT. Generalmente, esto implica:
- Un endpoint de login que recibe credenciales y devuelve un JWT.
- Un filtro de seguridad que intercepta las solicitudes, valida el JWT en el encabezado
Authorizationy construye unAuthenticationreactivo.
Puedes crear tu propio ServerWebExchangeMatcher y ServerAuthenticationConverter para procesar el token y autenticar al usuario sin necesidad de sesiones.
Spring Security Reactivo se integra perfectamente con el modelo de programación reactiva, asegurando que tus mecanismos de seguridad no introduzcan bloqueos o cuellos de botella en tus aplicaciones de alto rendimiento.
6. WebSockets con WebFlux: Comunicación Bidireccional en Tiempo Real
Mientras que Server-Sent Events (SSE) son excelentes para la comunicación unidireccional del servidor al cliente, las aplicaciones que requieren comunicación bidireccional en tiempo real, como chats, juegos en línea o herramientas de colaboración, necesitan WebSockets. WebSockets proporcionan un canal de comunicación dúplex completo a través de una única conexión TCP. Spring WebFlux ofrece un soporte robusto y reactivo para WebSockets.
¿Cómo funcionan los WebSockets?
A diferencia de HTTP, que es de corta duración y sin estado, los WebSockets comienzan con un handshake HTTP. Una vez que este handshake es exitoso, la conexión se "actualiza" a un protocolo WebSocket, permaneciendo abierta indefinidamente. Esto permite que tanto el cliente como el servidor envíen mensajes de forma asíncrona en cualquier momento.
WebSockets con Spring WebFlux
Spring WebFlux proporciona una API funcional para manejar WebSockets, aprovechando Flux y Mono para la gestión de mensajes reactivos.
WebSocketHandler: Es la interfaz principal que implementas para manejar la lógica de la conexión WebSocket. El métodohandlerecibe unWebSocketSessionque te permite enviar y recibir mensajes.WebSocketHandlerAdapterySimpleUrlHandlerMapping: Estos beans son necesarios para mapear las URLs a tusWebSocketHandlers específicos.
Configuración de WebSocket
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.web.reactive.handler.SimpleUrlHandlerMapping;
import org.springframework.web.reactive.socket.WebSocketHandler;
import org.springframework.web.reactive.socket.server.WebSocketService;
import org.springframework.web.reactive.socket.server.support.HandshakeWebSocketService;
import org.springframework.web.reactive.socket.server.support.WebSocketHandlerAdapter;
import org.springframework.web.reactive.socket.server.upgrade.ReactorNettyRequestUpgradeStrategy;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.time.Duration;
import java.util.HashMap;
import java.util.Map;
@Configuration
public class WebSocketConfig {
@Bean
public SimpleUrlHandlerMapping webSocketHandlerMapping(WebSocketHandler echoHandler) {
Map<String, WebSocketHandler> map = new HashMap<>();
map.put("/echo", echoHandler); // Mapea /echo a nuestro handler
map.put("/time-stream", new TimeStreamWebSocketHandler()); // Otro handler
return new SimpleUrlHandlerMapping(map);
}
@Bean
public WebSocketHandlerAdapter handlerAdapter(WebSocketService webSocketService) {
return new WebSocketHandlerAdapter(webSocketService);
}
@Bean
public WebSocketService webSocketService() {
// Usa Reactor Netty por defecto, que es el servidor webflux por defecto
return new HandshakeWebSocketService(new ReactorNettyRequestUpgradeStrategy());
}
@Bean
public WebSocketHandler echoHandler() {
return session -> session.send(
session.receive() // Recibe mensajes del cliente
.doOnNext(message -> System.out.println("Received: " + message.getPayloadAsText()))
.map(message -> session.textMessage("ECHO: " + message.getPayloadAsText())) // Eco de vuelta
).and(session.receive()
.doOnError(throwable -> System.err.println("Error en la conexión WebSocket: " + throwable.getMessage()))
.then()); // Mantener la conexión abierta hasta que se complete o haya un error
}
}
Creando un WebSocketHandler
Aquí tienes un ejemplo de un WebSocketHandler que envía la hora actual cada segundo:
// En un archivo separado o como inner class
import org.springframework.web.reactive.socket.WebSocketHandler;
import org.springframework.web.reactive.socket.WebSocketMessage;
import org.springframework.web.reactive.socket.WebSocketSession;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.time.Duration;
import java.time.LocalDateTime;
public class TimeStreamWebSocketHandler implements WebSocketHandler {
@Override
public Mono<Void> handle(WebSocketSession session) {
// Envía un mensaje cada segundo al cliente
Flux<WebSocketMessage> output = Flux.interval(Duration.ofSeconds(1))
.map(value -> session.textMessage("Current Time: " + LocalDateTime.now()));
// Mantén la conexión abierta para recibir mensajes (aunque este handler no los procese)
// La conexión se cierra cuando el Mono<Void> retornado se completa
return session.send(output)
.and(session.receive() // Esto es importante para mantener la conexión abierta
.doOnNext(message -> System.out.println("Received from client on time stream: " + message.getPayloadAsText()))
.then()); // No hacemos nada con los mensajes recibidos aquí, solo los logueamos
}
}
Cliente JavaScript para WebSockets
const ws = new WebSocket('ws://localhost:8080/echo'); // Para el handler de eco
ws.onopen = function(event) {
console.log("Conectado al WebSocket!");
ws.send("Hola desde el cliente!");
};
ws.onmessage = function(event) {
console.log("Mensaje recibido del servidor:", event.data);
};
ws.onclose = function(event) {
console.log("Conexión WebSocket cerrada:", event.code, event.reason);
};
ws.onerror = function(error) {
console.error("Error WebSocket:", error);
};
// Para enviar más mensajes:
// ws.send("Otro mensaje...");
// Para el handler de tiempo:
// const wsTime = new WebSocket('ws://localhost:8080/time-stream');
// wsTime.onmessage = function(event) {
// console.log("Tiempo recibido:", event.data);
// };
WebSockets con WebFlux te permiten construir aplicaciones de comunicación en tiempo real altamente eficientes, aprovechando la capacidad de Spring para manejar flujos de datos reactivos de forma nativa.
7. Buenas Prácticas en Producción para Aplicaciones WebFlux
Desarrollar una aplicación WebFlux es solo una parte del desafío; desplegarla y mantenerla en producción requiere atención a varias buenas prácticas para asegurar su rendimiento, estabilidad y observabilidad.
1. Monitoreo y Observabilidad
Las aplicaciones reactivas pueden ser más difíciles de depurar sin las herramientas adecuadas debido a la naturaleza asíncrona y la transición de hilos.
- Métricas (Micrometer/Prometheus): Spring Boot Actuator, combinado con Micrometer, facilita la exposición de métricas (JVM, WebFlux, Reactor, etc.) que pueden ser recolectadas por sistemas como Prometheus y visualizadas en Grafana. Monitorea la latencia, el rendimiento del Event Loop, el uso de memoria y el número de conexiones activas.
- Logging (Structured Logging): Utiliza un sistema de logging que soporte logging estructurado (ej. SLF4J con Logback configurado para JSON) para facilitar el análisis con herramientas como ELK Stack (Elasticsearch, Logstash, Kibana) o Grafana Loki.
- APM (Application Performance Monitoring): Herramientas como Dynatrace, New Relic o AppDynamics pueden proporcionar visibilidad profunda en el rendimiento de tu aplicación, incluyendo la trazabilidad de transacciones a través de hilos y servicios.
- Tracing (Brave/OpenTelemetry): Implementa Distributed Tracing (ej. con Spring Cloud Sleuth y Zipkin/Jaeger) para seguir el rastro de una solicitud a través de múltiples servicios, especialmente crucial en arquitecturas de microservicios reactivos.
2. Gestión de Recursos
- Connection Pooling: Asegúrate de que tus conexiones a bases de datos reactivas (R2DBC, MongoDB reactive drivers) o a otros servicios externos (WebClient) utilicen connection pooling para evitar la sobrecarga y el agotamiento de recursos.
- Timeouts: Configura timeouts apropiados en
WebClienty en tus servidores para evitar que las solicitudes de larga duración o los servicios externos lentos bloqueen los recursos del Event Loop. - Límites de Conexión: Establece límites de conexión adecuados en tus servidores (Netty, Undertow) para prevenir la sobrecarga.
3. Contrapresión Efectiva
Aunque Reactor maneja la contrapresión de forma nativa, es crucial entender cuándo y cómo se aplica, especialmente al integrar con sistemas que no son reactivos o que no la soportan. Asegúrate de que tus flujos de datos estén diseñados para manejar el backpressure correctamente para evitar la sobrecarga del consumidor.
4. Seguridad
- Principio de Mínimo Privilegio: Asegúrate de que tu aplicación solo tenga los permisos necesarios para realizar sus funciones.
- Secret Management: No guardes credenciales directamente en el código o en archivos de configuración. Utiliza soluciones de gestión de secretos como HashiCorp Vault, AWS Secrets Manager o Kubernetes Secrets.
- Actualizaciones y Parches: Mantén tus dependencias de Spring Boot, Spring Security y Reactor actualizadas para beneficiarte de las últimas correcciones de seguridad.
- HTTPS: Siempre utiliza HTTPS en producción para asegurar la comunicación cliente-servidor.
5. Configuración y Despliegue
- Externalización de la Configuración: Utiliza Spring Cloud Config Server, o simplemente
application.properties/application.ymlcon perfiles, y variables de entorno para gestionar la configuración de forma externa al artefacto de despliegue. - Contenedores (Docker/Kubernetes): Empaquetar tu aplicación en un contenedor Docker facilita el despliegue, la escalabilidad y la gestión de dependencias en entornos como Kubernetes.
- Liveness y Readiness Probes: En Kubernetes, configura Liveness y Readiness Probes para que el orquestador pueda saber cuándo tu aplicación está saludable y lista para recibir tráfico. Spring Boot Actuator proporciona endpoints
/actuator/healthque son perfectos para esto. - Escalabilidad: Las aplicaciones WebFlux son inherentemente escalables horizontalmente. Asegúrate de que tu infraestructura de despliegue (Kubernetes, balanceadores de carga) pueda escalar tu aplicación de manera eficiente.
6. Pruebas de Carga y Rendimiento
Realiza pruebas de carga exhaustivas para simular escenarios de alto tráfico y verificar cómo se comporta tu aplicación WebFlux bajo presión. Esto te ayudará a identificar cuellos de botella y a optimizar la configuración.
7. Manejo de Errores Robustos
- ErrorWebExceptionHandler: Asegúrate de tener un
ErrorWebExceptionHandlerglobal bien configurado para manejar excepciones no capturadas y proporcionar respuestas de error consistentes y amigables para el cliente, sin exponer detalles internos. - Circuit Breakers: Implementa patrones de Circuit Breaker (ej. con Resilience4j) al interactuar con servicios externos para evitar cascadas de fallos cuando un servicio dependiente no está disponible o es lento.
Al seguir estas buenas prácticas, puedes asegurar que tus aplicaciones Spring WebFlux no solo sean rápidas y eficientes en desarrollo, sino también robustas, seguras y fáciles de operar en producción.
Conclusión
En esta cuarta entrega de nuestra serie sobre Spring WebFlux, hemos explorado características avanzadas y cruciales que elevan el desarrollo de aplicaciones reactivas. Desde la implementación de Server-Sent Events (SSE) para flujos de datos unidireccionales hasta la robusta comunicación WebSocket para interacciones bidireccionales en tiempo real, hemos visto cómo Spring WebFlux simplifica la construcción de aplicaciones de tiempo real.
Hemos profundizado en la importancia de la contrapresión (backpressure), un mecanismo vital para garantizar la estabilidad del sistema al permitir que los consumidores controlen el flujo de datos. La gestión del contexto reactivo se ha revelado como una solución elegante para compartir información a través de los límites de los hilos en un entorno asíncrono, mientras que las herramientas de testing reactivo como StepVerifier y WebTestClient demuestran ser indispensables para asegurar la corrección de nuestros flujos. Finalmente, abordamos la integración de Spring Security Reactivo para asegurar nuestras aplicaciones de forma no bloqueante y delineamos un conjunto de buenas prácticas para la producción, fundamentales para el monitoreo, la estabilidad y la escalabilidad de nuestras aplicaciones WebFlux.
Esta serie ha cubierto los pilares esenciales para construir aplicaciones reactivas de alto rendimiento con Spring WebFlux. Con una base sólida en fundamentos, arquitectura, comunicación de datos, seguridad, pruebas y consideraciones de producción, estás bien preparado para enfrentar desafíos reales y llevar tus aplicaciones reactivas al siguiente nivel.
Como continuación natural de este camino, te recomendamos explorar algunas áreas complementarias que potenciarán aún más tus habilidades en entornos reactivos:
- R2DBC a profundidad: Explora la integración con bases de datos relacionales reactivas, optimización de consultas y rendimiento en entornos de alta demanda.
- Spring Cloud Gateway: Descubre cómo usar esta herramienta basada en WebFlux para implementar enrutamiento, seguridad y resiliencia en arquitecturas de microservicios.
- Programación Reactiva en el Frontend: Investiga cómo frameworks como React, Angular o Vue pueden conectarse eficientemente con backends WebFlux en escenarios de tiempo real.
- WebFlux y GraalVM Native Image: Evalúa las ventajas de empaquetar tus aplicaciones como imágenes nativas para mejorar el rendimiento y reducir el consumo de recursos.
- Patrones de resiliencia con WebFlux: Profundiza en técnicas como Circuit Breaker, Retry, Timeout y Rate Limiting mediante el uso de Resilience4j en entornos reactivos.
Explorar estos temas no solo ampliará tu dominio técnico, sino que también te permitirá diseñar soluciones más eficientes, resilientes y adaptadas a los retos actuales del desarrollo moderno. La programación reactiva, bien aplicada, abre la puerta a aplicaciones verdaderamente escalables y sensibles a la demanda del usuario.
Spring WebFlux 3: Comunicación, Datos y Errores Reactivos
- Mauricio ECR
- Arquitectura
- 24 May, 2025
¡Continuemos nuestro viaje por el fascinante mundo de Spring WebFlux! En la Parte 1, sentamos las bases de la programación reactiva y exploramos Project Reactor, el corazón de WebFlux. En la **Pa
Spring WebFlux 3: Comunicación, Datos y Errores Reactivos
- Mauricio ECR
- Arquitectura
- 24 May, 2025
¡Continuemos nuestro viaje por el fascinante mundo de Spring WebFlux!
En la Parte 1, sentamos las bases de la programación reactiva y exploramos Project Reactor, el corazón de WebFlux. En la Parte 2, nos adentramos en la arquitectura de WebFlux y aprendimos a construir endpoints utilizando tanto anotaciones como el enfoque funcional.
Ahora, en esta Parte 3, nos enfocaremos en cómo las aplicaciones WebFlux interactúan con el mundo exterior: cómo consumen otros servicios de manera reactiva, cómo persisten y recuperan datos en bases de datos reactivas, y, crucialmente, cómo gestionamos los errores que inevitablemente surgen en estos flujos asíncronos.
Comunicación con Servicios Externos (WebClient)
En el ecosistema de microservicios actual, es muy común que nuestras aplicaciones necesiten consumir APIs externas. Spring WebFlux nos proporciona una herramienta poderosa y reactiva para esto: WebClient. Es la contraparte no bloqueante de RestTemplate y la forma recomendada de hacer llamadas HTTP en un contexto reactivo.
WebClient
WebClient es un cliente HTTP no bloqueante que forma parte del módulo spring-webflux. Está diseñado para aprovechar la pila reactiva de principio a fin, lo que significa que no bloqueará hilos mientras espera respuestas de servicios externos, maximizando la eficiencia de tu aplicación WebFlux.
Su API es fluida y declarativa, similar a la forma en que construyes flujos con Mono y Flux.
Configuración Básica:
Puedes configurar WebClient de diversas maneras. La forma más común es inyectarlo como un bean en tu clase, o construir una instancia en línea. Puedes especificar una URL base, encabezados comunes, timeouts, filtros y más.
// Configuración como Bean (ejemplo en una clase @Configuration)
@Configuration
public class WebClientConfig {
@Bean
public WebClient externalApiClient(WebClient.Builder webClientBuilder) {
return webClientBuilder
.baseUrl("https://api.example.com") // URL base para todas las peticiones
.defaultHeader(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE) // Encabezado por defecto
.clientConnector(new ReactorClientHttpConnector(
HttpClient.create().responseTimeout(Duration.ofSeconds(5)) // Timeout de 5 segundos
))
.build();
}
}
Consumo de Respuestas Reactivas:
Después de definir la petición (GET, POST, PUT, DELETE, etc.), usas métodos como:
.retrieve(): Inicia la recuperación de la respuesta..bodyToMono(Class<T> type): Convierte el cuerpo de la respuesta en unMonode un objeto de tipoT. Útil cuando esperas una única respuesta (ej., un objeto JSON)..bodyToFlux(Class<T> type): Convierte el cuerpo de la respuesta en unFluxde objetos de tipoT. Útil para listas o streams de datos (ej., una lista de objetos JSON)..bodyToMono(ParameterizedTypeReference<T> typeRef)/.bodyToFlux(ParameterizedTypeReference<T> typeRef): Útil para tipos genéricos (ej.,List<MyObject>)..toEntity(Class<T> type)/.toEntityList(Class<T> type)/.toEntityFlux(Class<T> type): Devuelve unMono<ResponseEntity<T>>oMono<ResponseEntity<List<T>>>para acceder a la respuesta completa (estado HTTP, cabeceras, cuerpo).
Casos Típicos/Práctica
Llamada GET a un servicio externo y procesar la respuesta reactivamente:
Asumiendo que
externalApiClientes unWebClientbean inyectado.public Mono<MyObject> getObjectById(String id) { return externalApiClient.get() // Inicia una petición GET .uri("/objects/{id}", id) // Define la URI con PathVariable .retrieve() // Recupera la respuesta .bodyToMono(MyObject.class); // Convierte el cuerpo a Mono<MyObject> }Llamada POST enviando un
Mono<?>como body:public Mono<MyObject> createObject(Mono<MyObject> newObjectMono) { return externalApiClient.post() // Inicia una petición POST .uri("/objects") .body(newObjectMono, MyObject.class) // Envía el Mono<MyObject> como cuerpo .retrieve() .bodyToMono(MyObject.class); // Espera la respuesta como Mono<MyObject> }Manejar múltiples llamadas a servicios externos en paralelo (
Mono.zip,Flux.merge,flatMap):Mono.zip: Combina los resultados de múltiplesMonos (oFluxs que emiten un solo elemento) en un soloMonoque contiene una tupla de sus resultados. Las operaciones se ejecutan en paralelo. Ideal para combinar resultados de diferentes tipos que son necesarios simultáneamente.public Mono<CombinedData> getCombinedData(String id) { Mono<User> userMono = externalApiClient.get().uri("/users/{id}", id).retrieve().bodyToMono(User.class); Mono<Order> orderMono = externalApiClient.get().uri("/orders/{id}", id).retrieve().bodyToMono(Order.class); return Mono.zip(userMono, orderMono, (user, order) -> { // Aquí se combinan los resultados cuando ambos Monos han completado return new CombinedData(user, order); }); }Flux.merge: Combina múltiplesPublishers (Mono o Flux) en un únicoFlux, entrelazando sus elementos tan pronto como son emitidos. Las operaciones se ejecutan en paralelo, y el orden de los elementos resultantes no está garantizado.public Flux<Item> getItemsFromMultipleSources() { Flux<Item> source1 = externalApiClient.get().uri("/items/source1").retrieve().bodyToFlux(Item.class); Flux<Item> source2 = externalApiClient.get().uri("/items/source2").retrieve().bodyToFlux(Item.class); return Flux.merge(source1, source2); // Los ítems de source1 y source2 se entrelazan }flatMap: (Ya cubierto en Parte 1, pero clave aquí) Úsalo cuando la transformación de un elemento inicial te lleva a realizar otra operación asíncrona que devuelve unMonooFlux. Permite encadenar operaciones secuenciales asíncronas.public Mono<OrderDetail> getOrderDetails(String orderId) { return externalApiClient.get().uri("/orders/{id}", orderId).retrieve().bodyToMono(Order.class) // 1. Obtener la orden .flatMap(order -> externalApiClient .get() .uri("/products/{id}", order.getProductId()).retrieve().bodyToMono(Product.class) // 2. Obtener el producto de la orden .map(product -> new OrderDetail(order, product))); // 3. Combinar y devolver OrderDetail }
Manejar errores de un servicio externo llamado con WebClient:
WebClientlanzaWebClientResponseException(o subclases comoWebClientResponseException.NotFound) si la respuesta HTTP es un error (4xx, 5xx). Puedes usar operadores de manejo de errores de Reactor comoonErrorResumeoonErrorReturn.public Mono<MyObject> getObjectByIdHandlingError(String id) { return externalApiClient.get() .uri("/objects/{id}", id) .retrieve() .onStatus(HttpStatus.NOT_FOUND::equals, // Si el estado es 404 response -> Mono.error(new MyCustomNotFoundException("Object not found: " + id))) // Mapea a una excepción personalizada .onStatus(HttpStatus::is5xxServerError, // Si es un error 5xx response -> Mono.error(new RuntimeException("External service error"))) // Mapea a otra excepción .bodyToMono(MyObject.class) .onErrorResume(MyCustomNotFoundException.class, e -> { // Si es MyCustomNotFoundException, devuelve un Mono.empty() o un default System.err.println("Handling not found: " + e.getMessage()); return Mono.empty(); // O Mono.just(new MyObject("Default object")); }) .onErrorReturn(RuntimeException.class, new MyObject("Error occurred, returning default")); // Si es RuntimeException, devuelve un objeto por defecto }
Manejo de Datos Reactivos
Una aplicación reactiva es más eficiente si toda su pila es no bloqueante, y esto incluye la capa de persistencia de datos. Acceder a bases de datos de forma reactiva es crucial para evitar cuellos de botella por I/O bloqueante.
Integración de WebFlux con Bases de Datos Reactivas
Para bases de datos relacionales, la API estándar para acceso reactivo es R2DBC (Reactive Relational Database Connectivity). Es el equivalente reactivo de JDBC, pero diseñado desde cero para ser no bloqueante y asíncrono. Spring Data ha adoptado R2DBC, proporcionando integraciones para bases de datos como PostgreSQL, H2, MySQL (con driver de terceros) y SQL Server.
Para bases de datos NoSQL, muchos de los drivers ya están diseñados para ser reactivos. Por ejemplo, Spring Data tiene módulos reactivos para:
- MongoDB:
spring-data-mongodb-reactive - Cassandra:
spring-data-cassandra-reactive - Redis:
spring-data-redis-reactive
Repositorios Reactivos:
Spring Data extiende sus interfaces de repositorio para el contexto reactivo. En lugar de extender CrudRepository, extiendes interfaces como ReactiveCrudRepository, ReactiveMongoRepository, ReactiveCassandraRepository, etc. Los métodos de estas interfaces devuelven Mono<?> o Flux<?>.
Casos Típicos/Práctica
Asumiendo una entidad User y un repositorio UserRepository que extiende ReactiveCrudRepository<User, Long> (para R2DBC) o ReactiveMongoRepository<User, String> (para MongoDB).
Guardar (
save):// En un servicio @Autowired private UserRepository userRepository; public Mono<User> saveUser(User user) { return userRepository.save(user); // Devuelve Mono<User> }Encontrar por ID (
findById):public Mono<User> findUserById(Long id) { return userRepository.findById(id); // Devuelve Mono<User> }Encontrar todos (
findAll):public Flux<User> findAllUsers() { return userRepository.findAll(); // Devuelve Flux<User> }Manejo de Transacciones en un Contexto Reactivo: Este es un tema un poco más avanzado y complejo. En un contexto bloqueante, las transacciones se manejan con
@Transactional, que delega a unThreadLocal. Sin embargo, losThreadLocalno funcionan en un contexto reactivo porque los elementos pueden pasar por diferentes hilos en diferentes momentos.Para transacciones reactivas, Spring Data proporciona la anotación
@Transactionalen combinación con la infraestructura de transacciones reactivas de Spring (por ejemplo,ReactiveTransactionManagerpara R2DBC). Cuando usas@Transactionalen un método reactivo, Spring se asegura de que todas las operaciones reactivas dentro de ese método (que interactúan con la misma base de datos) se ejecuten dentro de la misma transacción.Es importante entender que una transacción se "adjunta" al
MonooFluxque se crea, no al hilo. Es decir, las operaciones dentro del flujo reactivo, si son parte de la misma transacción, se aseguran de comprometerse o revertirse juntas.@Service public class UserServiceImpl implements UserService { @Autowired private UserRepository userRepository; @Transactional // Esta anotación ahora trabaja con ReactiveTransactionManager public Mono<User> createUserAndAudit(User user) { return userRepository.save(user) // Guarda el usuario .flatMap(savedUser -> { // Simula una operación de auditoría que debe ser parte de la misma transacción // Si AuditRepository fuera reactivo y manejara transacciones. // return auditRepository.save(new AuditLog(savedUser.getId(), "User created")); System.out.println("User saved, attempting audit for: " + savedUser.getUsername()); return Mono.just(savedUser); // Devolver el usuario guardado }) .doOnError(e -> System.err.println("Transaction rolled back due to: " + e.getMessage())); // Manejo de error de transacción } }El desafío es que todas las operaciones dentro de la transacción deben ser reactivas y deben usar la misma conexión transaccional. Es un área donde la depuración puede ser más compleja que con las transacciones síncronas.
Manejo de Errores en Streams Reactivos
El manejo de errores es crucial en cualquier aplicación, y en los flujos reactivos tiene sus propias particularidades. Como ya mencionamos, cuando un error es emitido (onError), la secuencia se termina. Para evitar que toda la aplicación se caiga o para proporcionar una recuperación elegante, Reactor ofrece operadores específicos.
Operadores de Manejo de Errores
onErrorReturn(T fallbackValue): Cuando elPublisheremite un error, este operador intercepta el error, emite un valor de respaldo (fallbackValue), y luego completa la secuencia normalmente (onComplete). El error original es consumido.// Si ocurre un error, devuelve el valor por defecto "Default Message" Mono.error(new RuntimeException("Simulated error")) .onErrorReturn("Default Message") .subscribe(System.out::println, System.err::println); // Imprime "Default Message"onErrorResume(Function<Throwable, Mono<T>> fallbackMonoProvider): Si ocurre un error, este operador intercepta el error y cambia a unPublisheralternativo (fallbackMonoProvider). Es útil cuando necesitas ejecutar una lógica asíncrona para recuperarte del error.// Si ocurre un error, cambia a un Mono que simula una recuperación Mono.error(new RuntimeException("Simulated error")) .onErrorResume(e -> { System.err.println("Error caught, resuming with alternative: " + e.getMessage()); return Mono.just("Recovered from error!"); }) .subscribe(System.out::println, System.err::println); // Imprime "Recovered from error!"onErrorMap(Function<Throwable, Throwable> errorMapper): Transforma un tipo de excepción en otro. Esto es útil para encapsular excepciones internas en excepciones más significativas para tu dominio de negocio.// Transforma RuntimeException en CustomBusinessException Mono.error(new RuntimeException("Database error")) .onErrorMap(RuntimeException.class, e -> new MyCustomBusinessException("Failed to process data: " + e.getMessage())) .subscribe(System.out::println, System.err::println); // Lanza MyCustomBusinessExceptiondoOnError(Consumer<Throwable> errorConsumer): Ejecuta una acción de efecto secundario cuando un error ocurre, pero no consume el error. El error continúa propagándose por el stream. Útil para logging o métricas sin alterar el flujo de error.// Logea el error, pero el error sigue propagándose Mono.error(new RuntimeException("Another simulated error")) .doOnError(e -> System.err.println("Logging error before propagation: " + e.getMessage())) .subscribe(System.out::println, System.err::println); // Imprime el log y luego lanza RuntimeExceptionretry(long numRetries)/retryWhen(Function<Flux<Throwable>, Publisher<?>> retrySignal): Intenta re-suscribirse alPublisheroriginal un número de veces o bajo ciertas condiciones.
Manejo Global de Errores en WebFlux (ErrorWebExceptionHandler)
Para centralizar el manejo de errores y proporcionar respuestas HTTP consistentes (ej. JSON con un formato de error estándar), WebFlux proporciona la interfaz ErrorWebExceptionHandler. Puedes implementar esta interfaz y registrarla como un bean para manejar todas las excepciones no capturadas por los operadores en tus flujos.
@Component
@Order(-1) // Asegura que este handler sea el primero en la cadena
public class GlobalErrorWebExceptionHandler implements ErrorWebExceptionHandler {
@Override
public Mono<Void> handle(ServerWebExchange exchange, Throwable ex) {
HttpStatus status;
String errorMessage;
if (ex instanceof MyCustomNotFoundException) {
status = HttpStatus.NOT_FOUND;
errorMessage = ex.getMessage();
} else if (ex instanceof IllegalArgumentException) {
status = HttpStatus.BAD_REQUEST;
errorMessage = "Invalid input: " + ex.getMessage();
} else {
status = HttpStatus.INTERNAL_SERVER_ERROR;
errorMessage = "An unexpected error occurred: " + ex.getMessage();
// Considerar logear la excepción aquí
}
// Construir la respuesta de error JSON
ErrorResponse errorResponse = new ErrorResponse(status.value(), errorMessage);
DataBufferFactory bufferFactory = exchange.getResponse().bufferFactory();
DataBuffer buffer = bufferFactory.wrap(toJson(errorResponse).getBytes()); // Convierte el objeto a JSON
exchange.getResponse().setStatusCode(status);
exchange.getResponse().getHeaders().setContentType(MediaType.APPLICATION_JSON);
return exchange.getResponse().writeWith(Mono.just(buffer));
}
private String toJson(Object obj) {
// Implementa la lógica para convertir el objeto a JSON (ej. con ObjectMapper de Jackson)
try {
return new ObjectMapper().writeValueAsString(obj);
} catch (JsonProcessingException e) {
return "{\"status\":500, \"message\":\"Error converting error response to JSON\"}";
}
}
// Clase auxiliar para la respuesta de error
private static class ErrorResponse {
public int status;
public String message;
public ErrorResponse(int status, String message) { this.status = status; this.message = message; }
}
}
Casos Típicos/Práctica
Manejo de un error específico dentro de una cadena de operadores: Supongamos un servicio que busca un usuario, pero puede lanzar
UserNotFoundExceptionsi no lo encuentra.public Mono<User> getUserProfile(String userId) { return userRepository.findById(userId) // Simula buscar en DB .switchIfEmpty(Mono.error(new UserNotFoundException("User not found with ID: " + userId))) // Si Mono.empty(), lanza excepción .onErrorResume(UserNotFoundException.class, e -> { System.err.println("Handled specific UserNotFoundException: " + e.getMessage()); return Mono.just(new User("defaultUser", "Default User")); // Devuelve un usuario por defecto }); }Centralizar el manejo de errores para devolver respuestas HTTP consistentes: Como se mostró en el ejemplo de
GlobalErrorWebExceptionHandlerarriba.- 404 Not Found: Mapear
MyCustomNotFoundExceptionaHttpStatus.NOT_FOUND. - 500 Internal Server Error: Para excepciones inesperadas, mapear a
HttpStatus.INTERNAL_SERVER_ERROR. - 400 Bad Request: Para errores de validación o entrada incorrecta, mapear a
HttpStatus.BAD_REQUEST.
El
GlobalErrorWebExceptionHandleres el lugar ideal para definir el formato JSON estándar de tus mensajes de error y sus códigos de estado HTTP asociados, asegurando que todos los errores que atraviesan tu aplicación sean presentados de manera uniforme al cliente.- 404 Not Found: Mapear
Conclusión
En esta tercera entrega, hemos cubierto pilares fundamentales para construir aplicaciones WebFlux robustas: la comunicación reactiva con servicios externos utilizando WebClient, la persistencia de datos con bases de datos reactivas a través de Spring Data R2DBC o drivers NoSQL, y el vital manejo de errores en los flujos reactivos, tanto a nivel de operador como de forma global con ErrorWebExceptionHandler.
Estos conocimientos son esenciales para construir aplicaciones que no solo sean rápidas y escalables, sino también resilientes y fáciles de mantener. En la Parte 4 y final de nuestra serie, abordaremos temas más avanzados como Server-Sent Events, el concepto de Backpressure y el Contexto Reactivo, y, por supuesto, cómo probar eficazmente nuestras aplicaciones WebFlux.
¡Nos vemos en la última parte para solidificar aún más tu conocimiento en WebFlux!
Spring WebFlux 2: Alta Concurrencia sin Más Hilos
- Mauricio ECR
- Arquitectura
- 12 May, 2025
¡Bienvenido de nuevo a nuestra inmersión en Spring WebFlux! 👋 En la primera parte de esta serie, exploramos el "por qué" de la programación reactiva, entendiendo los problemas del bloqueo y descubri
Spring WebFlux 2: Alta Concurrencia sin Más Hilos
- Mauricio ECR
- Arquitectura
- 12 May, 2025
¡Bienvenido de nuevo a nuestra inmersión en Spring WebFlux! 👋
En la primera parte de esta serie, exploramos el "por qué" de la programación reactiva, entendiendo los problemas del bloqueo y descubriendo a Project Reactor como el motor que impulsa los flujos de datos asíncronos. Ahora que tenemos una base sólida sobre los principios reactivos y los tipos Mono/Flux, es momento de subir un nivel y entender cómo Spring WebFlux aplica estos conceptos para construir aplicaciones web eficientes y escalables.
En esta segunda entrega, nos centraremos en la arquitectura que diferencia a WebFlux de su predecesor, Spring MVC, y aprenderemos las dos formas principales de definir los endpoints de nuestra API reactiva.
3. Arquitectura de Spring WebFlux
Si Spring MVC se construyó sobre la API de Servlets (diseñada originalmente para un modelo síncrono de un hilo por petición), Spring WebFlux se construye sobre una pila completamente reactiva y no bloqueante. Esta diferencia fundamental es la clave de su capacidad para manejar alta concurrencia.
Teoría: Componentes Clave
La arquitectura de WebFlux se basa en:
- Servidores No Bloqueantes: A diferencia de depender de un Contenedor de Servlets (como Tomcat, Jetty) configurado de forma tradicional, WebFlux utiliza servidores web diseñados para manejar I/O no bloqueante. El servidor por defecto integrado con Spring Boot WebFlux es Netty, un framework asíncrono basado en eventos muy popular en la industria por su rendimiento. Sin embargo, WebFlux es flexible y también soporta otros servidores reactivos como Undertow o incluso Servlets 3.1+ API en modo no bloqueante (aunque el uso de Netty o Undertow es más común y eficiente para aprovechar plenamente el potencial reactivo).
- EventLoop: El corazón del procesamiento no bloqueante. En lugar de asignar un hilo por petición, WebFlux (y los servidores como Netty) utilizan un pequeño número de hilos llamados "Event Loop threads". Estos hilos no realizan operaciones de I/O bloqueantes directamente. En cambio, delegan la operación al sistema operativo y quedan libres para procesar otras tareas o peticiones. Cuando la operación de I/O se completa (por ejemplo, llega la respuesta de una base de datos o un servicio externo), el sistema operativo notifica al Event Loop, que entonces toma el resultado y continúa el procesamiento del flujo reactivo asociado a esa petición.
- Reactor Core: Como vimos en la Parte 1, Project Reactor proporciona los tipos
MonoyFluxy los operadores para componer la lógica asíncrona. WebFlux se integra estrechamente con Reactor. - Spring Web Reactive Framework: Capas por encima de Reactor y el servidor para proporcionar la funcionalidad web: manejo de peticiones, ruteo, serialización/deserialización, manejo de errores, etc.
Cómo WebFlux Maneja las Peticiones (El Pipeline Reactivo)
Cuando una petición HTTP llega a un servidor WebFlux:
- Uno de los Event Loop threads del servidor la recibe.
- La petición pasa a través de la cadena de procesamiento de WebFlux (filtros, ruteo).
- La petición llega al Handler (controlador o función manejadora) correspondiente.
- El Handler ejecuta la lógica de negocio, que típicamente involucra operaciones que devuelven
MonooFlux(ej: llamar a un servicio, acceder a una base de datos reactiva). - Estas operaciones, al ser reactivas y no bloqueantes, no detienen el Event Loop thread. El thread delega la tarea (ej: consulta a DB) y queda libre.
- Cuando la operación asíncrona finaliza (ej: la DB devuelve resultados), uno de los Event Loop threads recibe la notificación.
- Los resultados fluyen de vuelta a través de la cadena de operadores definida en el
Mono/Flux. - El resultado final del
Mono/Fluxse convierte en una respuesta HTTP y se envía de vuelta al cliente, de nuevo, utilizando los Event Loop threads de forma no bloqueante.
Todo el procesamiento, desde la recepción de la petición hasta el envío de la respuesta, se maneja sin bloquear los hilos principales, permitiendo que un pequeño número de hilos gestione una alta concurrencia.
Diferencias Arquitectónicas Fundamentales con Spring MVC
| Característica | Spring MVC (Tradicional) | Spring WebFlux (Reactivo) |
|---|---|---|
| Modelo de Hilos | Thread-per-request (Bloqueante) | Event Loop (No Bloqueante) |
| Contenedor/Servidor | Basado en Servlet API (Tomcat, Jetty, etc.) | Basado en servidores reactivos (Netty, Undertow) o Servlet 3.1+ no bloqueante |
| Manejo de I/O | Bloqueante (por defecto) | No Bloqueante |
| Dependencies Base | spring-webmvc |
spring-webflux |
| Tipos de Retorno | Objetos POJO, ResponseEntity, ModelAndView, etc. |
Mono<?>, Flux<?>, ResponseEntity<Mono<?>>, etc. |
| Backpressure | No aplica directamente | Soportado nativamente a través de Reactive Streams |
¿Puedes usar Spring MVC y Spring WebFlux en el mismo proyecto?
Generalmente no. Aunque es técnicamente posible tener ambas dependencias en el classpath, Spring Boot configurará automáticamente solo una de las dos pilas web (MVC o WebFlux) basándose en la que encuentre primero o una configuración explícita. Son dos arquitecturas de manejo de peticiones fundamentalmente diferentes que no están diseñadas para coexistir y procesar la misma petición dentro del mismo contexto de aplicación Spring de forma híbrida y coherente. Debes elegir una u otra para tu aplicación web principal.
Casos Típicos/Práctica
Flujo de una Petición Típica en WebFlux:
- Llega petición HTTP a Netty (Event Loop thread A la recibe).
- WebFlux la rutea a un
HandlerFunction(el mismo thread A). - El Handler llama a un
UserService.findById(id)que devuelveMono<User>. UserServiceusa unReactiveUserRepository.findById(id)(que usa un driver R2DBC no bloqueante).- El Event Loop thread A delega la consulta a la DB y queda libre.
- Cuando la DB responde, otro Event Loop thread (B) recibe la notificación.
- El thread B retoma el flujo del
Mono<User>. - El resultado
Userfluye de regreso al Handler. - El Handler devuelve el
Mono<User>, que WebFlux serializa a JSON. - El Event Loop thread B envía la respuesta HTTP de vuelta al cliente.
Modelo de Hilos de Spring MVC vs. WebFlux:
- MVC: Un pico de 1000 peticiones concurrentes esperando por una DB lenta podría requerir 1000 hilos (o el tamaño máximo del pool), muchos de ellos inactivos.
- WebFlux: Esas mismas 1000 peticiones podrían ser manejadas por 4-8 Event Loop threads, que nunca esperan, simplemente gestionan el estado de las operaciones asíncronas pendientes. Esto libera recursos para otras tareas.
4. Creación de Endpoints (Controladores y Endpoints Funcionales)
Spring WebFlux ofrece dos enfoques principales para definir los puntos finales de tu API: el modelo tradicional basado en anotaciones y un modelo más funcional.
Teoría: Dos Enfoques
- Basado en Anotaciones: Similar a Spring MVC, usas anotaciones como
@RestController,@RequestMapping,@GetMapping,@PostMapping,@RequestBody, etc. La diferencia clave es que los métodos del controlador deben devolver tipos reactivos (Mono<?>oFlux<?>). - Endpoints Funcionales: Un enfoque más funcional y declarativo. Defines las rutas usando
RouterFunctiony los manejadores de peticiones usandoHandlerFunction. No hay anotaciones a nivel de método o clase; es todo código Java.
Uso de Anotaciones con Tipos Reactivos
Es el enfoque más familiar si vienes de Spring MVC. Simplemente creas clases con @RestController y métodos con anotaciones de mapeo HTTP. La diferencia crucial es el tipo de retorno:
- Devuelve
Mono<T>si esperas 0 o 1 objetoTen la respuesta. - Devuelve
Flux<T>si esperas 0 a N objetosTen la respuesta (esto puede ser un array JSON o un stream de datos, por ejemplo, en Server-Sent Events). - Puedes envolver el tipo reactivo en
ResponseEntitypara tener control sobre el estado HTTP, cabeceras, etc.:Mono<ResponseEntity<T>>oResponseEntity<Flux<T>>.
Recibir datos en el cuerpo de la petición también se hace reactivamente: usas @RequestBody con Mono<T>.
Uso de Endpoints Funcionales
Este enfoque desacopla completamente la definición de la ruta de la lógica de manejo de la petición.
RouterFunction<ServerResponse>: Define cómo las peticiones se rutean a losHandlerFunctionbasándose en predicados (métodos HTTP, rutas, cabeceras, etc.). Usas la claseRouterFunctionspara construirlas (route(RequestPredicate, HandlerFunction)).HandlerFunction<ServerResponse>: Contiene la lógica de negocio para manejar una petición. Recibe unServerRequestcomo entrada y devuelve unMono<ServerResponse>. La claseServerResponsese usa para construir la respuesta (estado HTTP, cuerpo, cabeceras).
Ventajas del Enfoque Funcional:
- Mayor separación de preocupaciones (ruteo vs. manejo).
- Más fácil de testear unitariamente (HandlerFunction es solo una función pura).
- Permite una construcción de rutas más programática y dinámica.
- Evita el uso de reflexion asociado a las anotaciones (micro-optimización).
Desventajas del Enfoque Funcional:
- Puede ser menos conciso y legible para APIs REST simples comparado con las anotaciones.
- Menos familiar para desarrolladores acostumbrados al modelo de anotaciones.
Casos Típicos/Práctica
Endpoint GET que devuelva un
Mono<MyObject>(Anotaciones):Asumiendo una clase
MyObject { String message; }@RestController @RequestMapping("/api/greeting") public class GreetingController { @GetMapping("/{name}") public Mono<MyObject> getGreeting(@PathVariable String name) { // Simula una operación asíncrona que devuelve un solo objeto return Mono.just(new MyObject("Hello, " + name)) .delayElement(Duration.ofMillis(500)); // Simula latencia } }Endpoint GET que devuelva un
Flux<MyObject>(Stream de datos) (Anotaciones):@RestController @RequestMapping("/api/numbers") public class NumberStreamController { @GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE) // Importante: MediaType.TEXT_EVENT_STREAM_VALUE para SSE public Flux<String> streamNumbers() { // Emite un número cada segundo indefinidamente return Flux.interval(Duration.ofSeconds(1)) .map(sequence -> "Event: " + sequence); } @GetMapping("/list") // Devuelve como JSON array public Flux<MyObject> getObjectsList() { return Flux.just(new MyObject("one"), new MyObject("two"), new MyObject("three")) .delayElements(Duration.ofMillis(100)); } }Endpoint POST que reciba un
Mono<MyObject>en el body (Anotaciones):@RestController @RequestMapping("/api/objects") public class ObjectController { @PostMapping public Mono<String> createObject(@RequestBody Mono<MyObject> objectMono) { // Recibe un Mono<MyObject> del cuerpo de la petición // flatMap es necesario porque objectMono es un Publisher y save es otro Publisher return objectMono .flatMap(obj -> { System.out.println("Recibido objeto: " + obj.getMessage()); // Simula guardar el objeto asíncronamente y devolver un ID return Mono.just("Object saved with ID: " + obj.getMessage().hashCode()) .delayElement(Duration.ofMillis(300)); }); } }Definir una ruta y su manejador usando el enfoque funcional:
Primero, el
HandlerFunction:// En un archivo separado, por ejemplo, src/main/java/com/example/demo/handler/GreetingHandler.java @Component // Spring lo detecta como un Bean public class GreetingHandler { public Mono<ServerResponse> getGreeting(ServerRequest request) { String name = request.pathVariable("name"); return Mono.just(new MyObject("Hello, " + name)) .delayElement(Duration.ofMillis(500)) // Simula latencia .flatMap(obj -> ServerResponse.ok() // Construye la respuesta HTTP 200 .contentType(MediaType.APPLICATION_JSON) // Define el tipo de contenido .bodyValue(obj)); // Pone el objeto en el cuerpo de la respuesta } public Mono<ServerResponse> createObject(ServerRequest request) { return request.bodyToMono(MyObject.class) // Extrae el cuerpo a un Mono<MyObject> .flatMap(obj -> { System.out.println("Recibido objeto (Funcional): " + obj.getMessage()); // Simula guardar return Mono.just("Object saved (Funcional) with ID: " + obj.getMessage().hashCode()) .delayElement(Duration.ofMillis(300)); }) .flatMap(responseString -> ServerResponse.status(HttpStatus.CREATED) // Construye respuesta 201 Created .contentType(MediaType.TEXT_PLAIN) .bodyValue(responseString)); } }Luego, el
RouterFunction(en una clase de configuración, por ejemplo):// En una clase de configuración, por ejemplo, src/main/java/com/example/demo/config/RoutingConfig.java @Configuration public class RoutingConfig { @Bean public RouterFunction<ServerResponse> route(GreetingHandler greetingHandler) { return RouterFunctions.route(GET("/api/functional/greeting/{name}").and(accept(MediaType.APPLICATION_JSON)), greetingHandler::getGreeting) .andRoute(POST("/api/functional/objects").and(contentType(MediaType.APPLICATION_JSON)), greetingHandler::createObject); // Combina con otras rutas } }¿Cuándo elegirías anotaciones vs. endpoints funcionales?
- Anotaciones: Ideal para proyectos que migran de Spring MVC, equipos familiarizados con el modelo de anotaciones, o APIs REST con estructuras estándar. Es a menudo más rápido de implementar para casos simples o CRUDs.
- Funcionales: Preferible para APIs con lógica de ruteo compleja o dinámica, si buscas una mayor separación de preocupaciones para facilitar el testing unitario de la lógica del manejador, o si simplemente prefieres un estilo más funcional y programático. Puede tener una curva de aprendizaje inicial si no estás acostumbrado.
Conclusión
En esta segunda entrega, hemos explorado la arquitectura fundamental de Spring WebFlux, entendiendo cómo su modelo no bloqueante basado en EventLoop y servidores como Netty le permite manejar eficientemente la alta concurrencia, a diferencia del modelo tradicional de Spring MVC. También hemos aprendido las dos vías principales para construir endpoints: el familiar enfoque basado en anotaciones (adaptado para devolver tipos reactivos) y el modelo más programático y funcional de RouterFunction y HandlerFunction, comprendiendo las fortalezas de cada uno y cuándo considerar usarlos.
Con la arquitectura y la creación de endpoints cubiertas, estamos listos para abordar la interacción de nuestra aplicación WebFlux con el mundo exterior y el manejo de datos y errores. En la próxima parte, nos sumergiremos en el uso de WebClient para consumir servicios externos reactivamente, la integración con bases de datos reactivas (R2DBC, drivers NoSQL) y las estrategias para gestionar errores en los flujos reactivos.
¡Hasta la próxima entrega de nuestra serie sobre WebFlux!
Spring WebFlux 1: Fundamentos Reactivos y el Corazón de Reactor
- Mauricio ECR
- Arquitectura
- 08 May, 2025
¡Hola, entusiasta del desarrollo moderno! 👋 En el vertiginoso mundo de las aplicaciones web, donde la escalabilidad y la eficiencia son reyes, ha surgido un paradigma que desafía el modelo tradicion
Spring WebFlux 1: Fundamentos Reactivos y el Corazón de Reactor
- Mauricio ECR
- Arquitectura
- 08 May, 2025
¡Hola, entusiasta del desarrollo moderno! 👋
En el vertiginoso mundo de las aplicaciones web, donde la escalabilidad y la eficiencia son reyes, ha surgido un paradigma que desafía el modelo tradicional de solicitud-respuesta síncrono: la Programación Reactiva. Y si trabajas con Spring, inevitablemente te encontrarás con Spring WebFlux, la respuesta de este popular framework a este emocionante cambio.
Prepararte para una entrevista sobre WebFlux implica comprender no solo cómo usarlo, sino por qué existe y cómo funciona por dentro. En esta primera entrega de nuestra serie, sentaremos las bases, explorando los principios reactivos y conociendo a Project Reactor, la biblioteca que impulsa WebFlux.
1. Fundamentos de Programación Reactiva y el "Por Qué" de WebFlux
Imagínate un restaurante. En el modelo tradicional (síncrono), un camarero toma una orden (petición), va a la cocina y espera a que el plato esté listo para llevarlo a la mesa. Mientras espera, no puede atender a nadie más. Si el restaurante se llena, necesitas más camareros (hilos) esperando. Esto escala, pero llega un punto en que tener demasiados camareros se vuelve ineficiente (consumo de memoria, sobrecarga del planificador de hilos).
Ahora, imagina un modelo diferente. El camarero toma la orden, la lleva a la cocina y, en lugar de esperar, vuelve a tomar más órdenes. Cuando un plato está listo, el cocinero avisa, y el camarero que esté libre lo recoge y lo lleva. Este es el modelo reactivo/asíncrono/no bloqueante. Los camareros (hilos) no se quedan inactivos esperando; están constantemente haciendo algo útil.
Teoría: ¿Qué es la Programación Reactiva?
La Programación Reactiva es un paradigma de programación que se centra en trabajar con flujos de datos asíncronos que reaccionan a cambios. No es solo sobre asincronía; es sobre gestionar la propagación de cambios y el manejo de "eventos" de manera eficiente y no bloqueante.
Aunque existe un "Reactive Manifesto" que define los principios de sistemas reactivos (responsivos, resilientes, elásticos y basados en mensajes), en el contexto de la programación reactiva a nivel de código, nos enfocamos más en cómo manejamos esos flujos de datos asíncronos.
Programación Síncrona vs. Asíncrona vs. No Bloqueante vs. Reactiva
Es crucial entender estas diferencias:
- Síncrona: Las operaciones se ejecutan secuencialmente. Una operación debe completarse antes de que la siguiente pueda comenzar. Un hilo realiza una tarea de principio a fin.
- Asíncrona: Una operación se inicia y el programa continúa ejecutando otras tareas sin esperar a que la primera termine. Cuando la operación asíncrona finaliza, a menudo notifica al programa (por ejemplo, a través de un callback o una promesa).
- No Bloqueante: Un subconjunto importante de la programación asíncrona. Una llamada a una función no bloqueante regresa inmediatamente, incluso si la operación solicitada no se ha completado. Si el resultado no está disponible, a menudo devuelve un valor especial (como
nullo un indicador de "pendiente"). No bloquea el hilo llamador. - Reactiva: Un estilo de programación que utiliza flujos de datos asíncronos y no bloqueantes. Se basa en el patrón Observer, donde un "Publisher" emite elementos y un "Subscriber" los consume reaccionando a ellos. Permite componer operaciones complejas sobre estos flujos de manera declarativa.
El Problema del Bloqueo (Thread per Request):
En las arquitecturas web tradicionales (como Spring MVC sobre Servlet API), el modelo común es "un hilo por petición". Cuando una petición llega, se le asigna un hilo del pool. Si esa petición necesita interactuar con algo lento (una base de datos, un servicio externo, una espera de I/O), el hilo asignado se bloquea esperando. Mientras está bloqueado, no puede atender otras peticiones. En escenarios de alto tráfico o latencia, esto lleva a:
- Agotamiento del pool de hilos.
- Alta demanda de recursos del sistema (memoria, CPU por el cambio de contexto entre muchos hilos).
- Disminución del rendimiento y la capacidad de respuesta.
La programación reactiva y WebFlux resuelven esto utilizando un modelo basado en eventos y no bloqueante. Un pequeño número de hilos (a menudo llamados Event Loop threads) maneja muchas peticiones concurrentemente. Cuando una operación de I/O es necesaria, el hilo no espera; delega la operación al sistema operativo y se libera para manejar otras peticiones. Cuando el resultado de la operación de I/O está listo, el sistema operativo notifica a uno de los hilos del Event Loop, que entonces procesa la respuesta.
Ventajas de Usar WebFlux
- Escalabilidad: Maneja un gran número de conexiones concurrentes con un número reducido de hilos, lo que se traduce en una mejor utilización de recursos y mayor capacidad para escalar horizontalmente.
- Uso Eficiente de Recursos: Menos hilos significan menos consumo de memoria y menos sobrecarga del planificador de hilos.
- Manejo de Latencia: Al no bloquear hilos en operaciones de I/O, la aplicación sigue siendo receptiva incluso cuando depende de servicios lentos o tiene alta latencia.
- Composición de Flujos Asíncronos: El modelo reactivo basado en operadores facilita la construcción de lógica compleja que involucra múltiples operaciones asíncronas.
¿Cuándo NO Usar WebFlux?
WebFlux no es una bala de plata para todos los casos. Hay situaciones donde Spring MVC tradicional puede ser más adecuado:
- Aplicaciones CPU-Bound: Si tu aplicación realiza principalmente cálculos intensivos que consumen mucha CPU, un modelo reactivo no te dará grandes beneficios en términos de escalabilidad, ya que los hilos estarán ocupados computando, no esperando I/O. De hecho, la sobrecarga del modelo reactivo podría ser detrimental.
- Aplicaciones Simples con Bajo Tráfico: Para APIs sencillas o aplicaciones internas con poca carga, la complejidad adicional de la programación reactiva puede no justificarse. El modelo síncrono de Spring MVC es a menudo más rápido de desarrollar en estos casos.
- Ecosistema Bloqueante: Si dependes fuertemente de bibliotecas o tecnologías que son inherentemente bloqueantes y no tienen alternativas reactivas, adoptar WebFlux implicará wrappers o adaptadores que pueden complicar el código.
Casos Típicos/Práctica
Hilo Bloqueado vs. Hilo No Bloqueado:
- Hilo Bloqueado: Imagina un hilo pidiendo datos a una base de datos y esperando pasivamente hasta que todos los datos llegan. Durante ese tiempo, el hilo no puede hacer nada más.
- Hilo No Bloqueado: El hilo pide los datos y, en lugar de esperar, le dice a la base de datos "avísame cuando tengas los datos". Luego, el hilo queda libre para procesar otra petición. Cuando la base de datos termina, notifica a un hilo disponible para que procese los resultados.
Escenario donde WebFlux Brilla: Una API Gateway que recibe miles de peticiones por segundo, cada una de las cuales necesita hacer varias llamadas a microservicios internos (con latencia variable) y a bases de datos antes de agregar y devolver la respuesta. En este escenario, un modelo tradicional agotaría rápidamente los hilos, mientras que WebFlux, al no bloquear, puede manejar la concurrencia eficientemente con muchos menos hilos.
¿Por qué Spring creó WebFlux si ya existía Spring MVC? Spring MVC se basa en la API de Servlets, que es fundamentalmente síncrona y bloqueante en su diseño original (aunque ha evolucionado). Para ofrecer una solución de programación reactiva y no bloqueante de extremo a extremo que pudiera competir con frameworks como Node.js o Vert.x en escenarios de alta concurrencia y I/O-bound, Spring necesitaba una arquitectura desde cero que no dependiera del modelo Servlet. WebFlux nació para llenar ese vacío, proporcionando una pila web completamente reactiva construida sobre bibliotecas como Reactor y servidores no bloqueantes como Netty.
2. Project Reactor: El Corazón de WebFlux
WebFlux no implementa la programación reactiva desde cero; se apoya en una biblioteca especializada para ello: Project Reactor. Reactor es una biblioteca de programación reactiva para JVM, basada en la especificación Reactive Streams, que define un estándar para el procesamiento de flujos de datos asíncronos con "backpressure".
Teoría: Conceptos Clave de Reactor
Reactor proporciona dos tipos principales para representar flujos de datos asíncronos:
- Mono: Representa un flujo reactivo que emite 0 o 1 elemento y luego se completa (o emite un error). Ideal para operaciones que devuelven un único resultado o ninguna (como guardar un registro, buscar por ID si existe, o una operación de borrado).
- Flux: Representa un flujo reactivo que emite 0 a N elementos y luego se completa (o emite un error). Ideal para operaciones que pueden devolver múltiples resultados (como buscar todos los usuarios, un stream de eventos, o resultados de una consulta paginada).
Estos tipos implementan la interfaz Publisher de Reactive Streams.
El modelo de Reactor (y Reactive Streams) se basa en cuatro interfaces principales:
- Publisher: Produce elementos (eventos). Es el origen de la secuencia. Solo tiene un método:
subscribe(Subscriber s). - Subscriber: Consume elementos emitidos por el Publisher. Define métodos de callback:
onSubscribe(Subscription s): Se invoca una vez cuando el Subscriber se suscribe exitosamente al Publisher. Recibe un objetoSubscription.onNext(T t): Se invoca para cada elemento emitido por el Publisher.onError(Throwable t): Se invoca si el Publisher encuentra un error. La secuencia termina.onComplete(): Se invoca cuando el Publisher ha terminado de emitir elementos exitosamente. La secuencia termina.
- Subscription: Representa la relación entre un Publisher y un Subscriber. Permite al Subscriber gestionar el flujo de datos (pedir más elementos - backpressure) o cancelar la suscripción. Métodos clave:
request(long n)ycancel(). - Operator: Son funciones puras que transforman, filtran, combinan o manipulan flujos. Reciben un Publisher como entrada y devuelven un nuevo Publisher. Encadenar operadores crea un pipeline reactivo.
El Ciclo de Vida de un Stream Reactivo
El ciclo de vida es fundamental:
- Un Subscriber se suscribe a un Publisher llamando a
publisher.subscribe(subscriber). - El Publisher, si acepta la suscripción, llama a
subscriber.onSubscribe(subscription), pasándole un objetoSubscription. - El Subscriber utiliza el objeto
Subscriptionpara solicitar elementos llamando asubscription.request(n). Esto es backpressure: el consumidor le dice al productor cuántos elementos está listo para manejar. - El Publisher emite elementos llamando a
subscriber.onNext(element)hasta que se alcanzan losnelementos solicitados o se agotan los elementos disponibles. - Este proceso de
request(n)yonNext(element)se repite. - Eventualmente, el Publisher terminará la secuencia llamando a
subscriber.onComplete()osubscriber.onError(error). Una vez queonCompleteoonErrorson llamados, la secuencia termina y no se emitirán más eventos. El Subscriber también puede cancelar la suscripción prematuramente llamando asubscription.cancel().
Importante: La ejecución real del flujo (el pushing de datos a través del pipeline) solo comienza cuando hay un Subscriber. Esto se conoce como lazy execution.
Operadores: ¿Qué son y por qué son importantes?
Los operadores son el poder de Reactor. Permiten construir lógica compleja sobre flujos de datos de manera declarativa y componible. Cada operador toma un Publisher de entrada y devuelve un nuevo Publisher modificado. Puedes encadenar múltiples operadores para construir una secuencia de procesamiento.
Ejemplos de categorías de operadores:
- Transformación:
map,flatMap,concatMap. - Filtrado:
filter,take,skip. - Combinación:
merge,zip,concat. - Manejo de Errores:
onErrorReturn,onErrorResume,doOnError. - Utilidad:
doOnNext,doOnComplete,delayElements.
Casos Típicos/Práctica
Diferencia entre Mono y Flux con ejemplos:
// Mono: Representa 0 o 1 elemento Mono<String> greeting = Mono.just("Hola Mundo"); // Emite "Hola Mundo" Mono<String> noValue = Mono.empty(); // Emite 0 elementos // Flux: Representa 0 a N elementos Flux<Integer> numbers = Flux.just(1, 2, 3, 4, 5); // Emite 1, 2, 3, 4, 5 Flux<String> greetings = Flux.fromIterable(Arrays.asList("Hello", "World", "Reactor")); // Emite "Hello", "World", "Reactor" Flux<Long> infinite = Flux.interval(Duration.ofSeconds(1)); // Emite un número cada segundo (infinito)- Ejemplo de Uso: Usarías un
Mono<User>para obtener los detalles de un usuario por su ID, y unFlux<Product>para obtener una lista de productos de una categoría.
- Ejemplo de Uso: Usarías un
Demostrar el uso de operadores comunes:
Flux.just(1, 2, 3, 4, 5) .filter(n -> n % 2 == 0) // Filtra solo números pares .map(n -> "Número par: " + n) // Transforma cada número en un String .subscribe(System.out::println); // Suscriptor que imprime cada elemento // Salida: // Número par: 2 // Número par: 4 Mono.just("spring") .map(String::toUpperCase) // Transforma a mayúsculas .subscribe(System.out::println); // Suscriptor // Salida: // SPRINGEntender bien
flatMapvsmap: ¡Crucial!map: Transforma cada elemento emitido por el origen sincrónicamente en otro elemento. Si la función de mapeo devuelve un tipo reactivo (MonooFlux), el resultado será unFluxdeMonos oFluxs anidados (unFlux<Mono<T>>oFlux<Flux<T>>), lo cual rara vez es lo que quieres.flatMap: Transforma cada elemento emitido por el origen en un nuevo Publisher (MonooFlux) y luego aplana (fusiona) los elementos de estos Publishers resultantes en un únicoFlux. Es ideal para operaciones asíncronas. El orden de los elementos resultantes no está garantizado conflatMapsi las operaciones internas tardan tiempos variables.concatMap: Similar aflatMap, pero garantiza que los Publishers internos se suscriban y emitan sus elementos en el mismo orden en que llegaron los elementos originales. Esto es útil cuando el orden es importante, pero puede ser menos eficiente queflatMapya que espera a que cada Publisher interno termine antes de procesar el siguiente.
// Ejemplo flatMap vs map Flux.just("Alpha", "Beta") .flatMap(word -> Mono.just(word.length()) // Crea un Mono con la longitud de la palabra (asíncrono o síncrono envuelto en Mono) .delayElement(Duration.ofMillis(word.length() * 100))) // Simula una operación asíncrona con retraso .subscribe(length -> System.out.println("flatMap - Longitud: " + length)); // Posible salida (el orden puede variar debido a delayElement y flatMap): // flatMap - Longitud: 5 // flatMap - Longitud: 4 Flux.just("Alpha", "Beta") .map(word -> Mono.just(word.length()) // Crea un Mono con la longitud .delayElement(Duration.ofMillis(word.length() * 100))) .subscribe(monoLength -> monoLength.subscribe(length -> System.out.println("map - Longitud: " + length))); // Necesitas suscribirte al Mono interno! // Salida (después de 500ms y 400ms): // map - Longitud: 5 // map - Longitud: 4 // ¡Fíjate que map devolvió un Flux<Mono<Integer>>! Tuvimos que suscribirnos a cada Mono. flatMap lo hizo automáticamente y aplanó el resultado. Flux.just("Alpha", "Beta") .concatMap(word -> Mono.just(word.length()) // Crea un Mono con la longitud .delayElement(Duration.ofMillis(word.length() * 100))) // Simula operación asíncrona con retraso .subscribe(length -> System.out.println("concatMap - Longitud: " + length)); // Salida (el orden está garantizado por concatMap): // concatMap - Longitud: 5 (espera 500ms) // concatMap - Longitud: 4 (luego espera 400ms)Secuencia que emita números y luego los transforme:
Flux.range(1, 10) // Emite números del 1 al 10 .map(n -> n * 2) // Multiplica cada número por 2 .filter(n -> n > 10) // Mantiene solo los resultados mayores que 10 .subscribe(result -> System.out.println("Resultado transformado: " + result), // onNext error -> System.err.println("Ocurrió un error: " + error), // onError () -> System.out.println("Secuencia completada.")); // onComplete // Salida: // Resultado transformado: 12 // Resultado transformado: 14 // Resultado transformado: 16 // Resultado transformado: 18 // Resultado transformado: 20 // Secuencia completada.¿Qué sucede si un Flux emite un error? ¿Cómo lo manejas? Cuando un Publisher emite un error a través de
onError(Throwable t), la secuencia termina inmediatamente. Ningún elemento posterior será emitido. El Subscriber recibe la notificaciónonError, y el flujo se detiene en ese punto. Para manejar errores de forma elegante, se usan operadores de manejo de errores (los veremos en detalle en un artículo posterior), comoonErrorReturn(devuelve un valor por defecto y completa),onErrorResume(cambia a un Publisher alternativo), oretry(intenta la secuencia de nuevo).subscribeOnvspublishOn: ¡Otro concepto fundamental! Controlan la ejecución concurrente.subscribeOn(Scheduler scheduler): Afecta el contexto de ejecución del Publisher original y toda la cadena de operadores subsiguiente hasta que se encuentra otropublishOn. Define en quéScheduler(un ejecutor de tareas, similar a un Thread Pool) se ejecutará el trabajo del Publisher y dónde comenzará el pipeline. Si hay múltiplessubscribeOn, solo el primero (el más cercano al Publisher) tiene efecto.publishOn(Scheduler scheduler): Afecta el contexto de ejecución de los operadores que le siguen en la cadena, no los que están antes o el Publisher original. Es útil para cambiar de contexto de ejecución en medio de un pipeline, por ejemplo, para pasar del hilo rápido de I/O a un pool de hilos de trabajo para una operación intensiva en CPU. Puede haber múltiplespublishOnen una cadena, cada uno afectando a la parte del pipeline que le sigue.
Scheduler ioScheduler = Schedulers.boundedElastic(); // Scheduler adecuado para I/O Scheduler computationScheduler = Schedulers.parallel(); // Scheduler adecuado para CPU-bound Flux.range(1, 5) .map(i -> { System.out.println("Map 1 en hilo: " + Thread.currentThread().getName()); return i * 2; }) .publishOn(computationScheduler) // Los operadores que siguen se ejecutarán aquí .map(i -> { System.out.println("Map 2 en hilo: " + Thread.currentThread().getName()); return i + 1; }) .subscribeOn(ioScheduler) // El Publisher original y todo comienza aquí (si no hay publishOn antes) .subscribe(result -> System.out.println("Subscripción en hilo: " + Thread.currentThread().getName() + " - Resultado: " + result)); // Posible Salida (los nombres de hilos variarán): // Map 1 en hilo: boundedElastic-1 // Map 1 en hilo: boundedElastic-1 // Map 1 en hilo: boundedElastic-1 // Map 1 en hilo: boundedElastic-1 // Map 1 en hilo: boundedElastic-1 // Map 2 en hilo: parallel-1 // Subscripción en hilo: parallel-1 - Resultado: 3 // Map 2 en hilo: parallel-2 // Subscripción en hilo: parallel-2 - Resultado: 5 // Map 2 en hilo: parallel-3 // Subscripción en hilo: parallel-3 - Resultado: 7 // Map 2 en hilo: parallel-4 // Subscripción en hilo: parallel-4 - Resultado: 9 // Map 2 en hilo: parallel-1 // Subscripción en hilo: parallel-1 - Resultado: 11 // Observa cómo el primer map se ejecuta en el scheduler de subscribeOn (boundedElastic), // mientras que el segundo map y la subscripción se ejecutan en el scheduler de publishOn (parallel).- Cuándo usar cada uno:
- Usa
subscribeOncerca del origen de tu stream (elPublisherque quizás interactúa con una API bloqueante envuelta o realiza una operación de I/O inicial) para asegurar que esa parte del trabajo no bloquee tus hilos principales. - Usa
publishOnpara cambiar de contexto de ejecución en medio del pipeline, por ejemplo, si después de una operación de I/O (que se ejecuta en un scheduler de I/O), necesitas realizar cálculos intensivos en CPU y quieres usar un pool de hilos diferente dedicado a la computación para no saturar los hilos de I/O.
- Usa
Conclusión
En esta primera parte, hemos desempacado los conceptos fundamentales que motivaron la creación de Spring WebFlux: los desafíos del bloqueo en arquitecturas tradicionales y cómo la programación reactiva, basada en flujos de datos asíncronos y no bloqueantes, ofrece una solución elegante y escalable. Hemos introducido Project Reactor como la biblioteca clave detrás de WebFlux, explorando sus tipos principales (Mono y Flux), el modelo Publisher/Subscriber/Subscription y la importancia de los operadores. Conceptos como flatMap vs map y subscribeOn vs publishOn son esenciales para dominar la programación reactiva con Reactor.
Comprender estas bases es el primer paso crucial. En la próxima entrega de esta serie, nos adentraremos en la arquitectura específica de Spring WebFlux y cómo se construyen las aplicaciones sobre este modelo reactivo, explorando el EventLoop y las diferencias arquitectónicas con Spring MVC.
¡Mantente reactivo!