tutoriales.com

Explorando la Programación Reactiva en Java con Project Reactor: Un Enfoque Práctico

Este tutorial te guiará a través de los fundamentos de Project Reactor, la biblioteca clave para la programación reactiva en Java. Aprenderás a utilizar `Flux` y `Mono` para manejar flujos de datos asíncronos y eventos, construyendo aplicaciones más responsivas y resilientes. Exploraremos operadores esenciales y patrones de diseño reactivos con ejemplos prácticos.

Intermedio20 min de lectura21 views
Reportar error

La programación reactiva se ha convertido en un paradigma crucial para construir aplicaciones modernas que necesitan ser altamente escalables, responsivas y tolerantes a fallos. En el ecosistema Java, Project Reactor es la biblioteca de elección, especialmente si trabajas con Spring WebFlux.

En este tutorial, profundizaremos en Project Reactor, entendiendo sus conceptos fundamentales y aplicándolos a través de ejemplos prácticos. Prepárate para transformar tu forma de pensar sobre el manejo de datos asíncronos y eventos.

🚀 ¿Qué es la Programación Reactiva y por qué Reactor?

La programación reactiva es un paradigma que se centra en el manejo de flujos de datos asíncronos y la propagación del cambio. Imagina un flujo de eventos o datos que llega con el tiempo. En lugar de esperar activamente a que cada elemento llegue, el paradigma reactivo reacciona a medida que los elementos están disponibles. Esto es particularmente útil para:

  • Interfaces de usuario: Reaccionar a clics, entradas de teclado.
  • Eventos de red: Manejar respuestas de API, mensajes de websocket.
  • Procesamiento de datos: Procesar grandes volúmenes de datos que llegan de forma continua.

💡 Principios Fundamentales

Los sistemas reactivos se basan en cuatro principios clave, conocidos como el Manifiesto Reactivo:

  • Responsiveness (Capacidad de Respuesta): Los sistemas responden rápidamente, incluso bajo carga.
  • Resilience (Resiliencia): Los sistemas se mantienen responsivos frente a fallos.
  • Elasticity (Elasticidad): Los sistemas pueden escalar o reducirse fácilmente según la demanda.
  • Message Driven (Orientado a Mensajes): Los sistemas interactúan a través de mensajes asíncronos para desacoplar componentes.

Project Reactor: El Corazón Reactivo de Spring

Project Reactor es una biblioteca de programación reactiva para la JVM, implementando la especificación Reactive Streams. Es la base de Spring WebFlux, permitiendo construir APIs REST no bloqueantes y eficientes. Sus tipos principales son Flux y Mono.

🛠️ Configuración Inicial: Añadiendo Reactor a tu Proyecto

Para empezar a usar Project Reactor, solo necesitas añadir la dependencia a tu proyecto Maven o Gradle.

Maven

<dependency>
    <groupId>io.projectreactor</groupId>
    <artifactId>reactor-core</artifactId>
    <version>3.5.12</version> <!-- Usar la última versión estable -->
</dependency>
<dependency>
    <groupId>io.projectreactor</groupId>
    <artifactId>reactor-test</artifactId>
    <version>3.5.12</version> <!-- Para pruebas reactivas -->
    <scope>test</scope>
</dependency>

Gradle

implementation 'io.projectreactor:reactor-core:3.5.12'
testImplementation 'io.projectreactor:reactor-test:3.5.12'
📌 Nota: Siempre verifica la última versión estable en el repositorio de Maven Central o en la página oficial de Project Reactor.

📖 Entendiendo Flux y Mono: Los Fundamentos Reactivos

Project Reactor introduce dos tipos de publicadores principales:

  • Flux<T>: Emite 0...N elementos (cero a muchos) y luego completa o emite un error. Es ideal para flujos de datos continuos o colecciones.
  • Mono<T>: Emite 0...1 elemento (cero o uno) y luego completa o emite un error. Es ideal para operaciones que devuelven un solo resultado o ninguna.

Ambos implementan la interfaz Publisher de Reactive Streams.

Publisher Flux<T> (0...N elementos) Mono<T> (0...1 elemento)

Creando Flux y Mono

Hay muchas maneras de crear estos publicadores. Veamos algunos ejemplos.

Creando Flux

import reactor.core.publisher.Flux;

public class FluxCreation {

    public static void main(String[] args) {
        // Desde elementos
        Flux<String> fluxFromElements = Flux.just("apple", "banana", "orange");
        fluxFromElements.subscribe(System.out::println);
        // Salida:
        // apple
        // banana
        // orange

        System.out.println("------------------");

        // Desde un array o Iterable
        Flux<Integer> fluxFromArray = Flux.fromArray(new Integer[]{1, 2, 3, 4, 5});
        fluxFromArray.subscribe(System.out::println);
        // Salida:
        // 1
        // 2
        // 3
        // 4
        // 5

        System.out.println("------------------");

        // Generando un rango
        Flux<Long> fluxFromRange = Flux.range(1, 5).map(i -> (long) i * 10);
        fluxFromRange.subscribe(System.out::println);
        // Salida:
        // 10
        // 20
        // 30
        // 40
        // 50

        System.out.println("------------------");

        // Flux vacío
        Flux<String> emptyFlux = Flux.empty();
        emptyFlux.subscribe(System.out::println, 
                            error -> System.err.println("Error: " + error.getMessage()), 
                            () -> System.out.println("Flux vacío completado"));
        // Salida:
        // Flux vacío completado

        System.out.println("------------------");
        
        // Flux que emite un error
        Flux<String> errorFlux = Flux.error(new RuntimeException("Algo salió mal!"));
        errorFlux.subscribe(System.out::println, 
                            error -> System.err.println("Error: " + error.getMessage()), 
                            () -> System.out.println("Completado"));
        // Salida:
        // Error: Algo salió mal!
    }
}

Creando Mono

import reactor.core.publisher.Mono;

public class MonoCreation {

    public static void main(String[] args) {
        // Desde un solo elemento
        Mono<String> monoFromValue = Mono.just("Hello Reactor!");
        monoFromValue.subscribe(System.out::println);
        // Salida:
        // Hello Reactor!

        System.out.println("------------------");

        // Mono vacío
        Mono<String> emptyMono = Mono.empty();
        emptyMono.subscribe(System.out::println, 
                            error -> System.err.println("Error: " + error.getMessage()), 
                            () -> System.out.println("Mono vacío completado"));
        // Salida:
        // Mono vacío completado

        System.out.println("------------------");

        // Mono que emite un error
        Mono<Integer> errorMono = Mono.error(new IllegalArgumentException("Argumento inválido"));
        errorMono.subscribe(System.out::println,
                            error -> System.err.println("Error: " + error.getMessage()),
                            () -> System.out.println("Completado"));
        // Salida:
        // Error: Argumento inválido

        System.out.println("------------------");

        // Mono diferido (la lógica se ejecuta al suscribirse)
        Mono<String> deferredMono = Mono.fromSupplier(() -> {
            System.out.println("Generando valor para Mono");
            return "Valor Diferido";
        });
        System.out.println("Antes de suscribirse al Mono diferido");
        deferredMono.subscribe(System.out::println);
        // Salida:
        // Antes de suscribirse al Mono diferido
        // Generando valor para Mono
        // Valor Diferido
    }
}

🌊 Operadores Reactivos: Transformando y Combinando Flujos

La verdadera potencia de Reactor reside en su vasta colección de operadores. Estos te permiten transformar, filtrar, combinar y gestionar flujos de datos de manera declarativa y funcional. Veamos algunos de los más comunes.

Operadores de Transformación

  • map(): Transforma cada elemento del flujo uno a uno.
  • flatMap(): Transforma cada elemento en un nuevo Publisher y aplana los resultados en un solo flujo. Es crucial para operaciones asíncronas encadenadas.
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

public class TransformationOperators {

    public static void main(String[] args) {
        // map: convierte strings a mayúsculas
        Flux.just("hello", "world")
            .map(String::toUpperCase)
            .subscribe(System.out::println);
        // Salida:
        // HELLO
        // WORLD

        System.out.println("------------------");

        // flatMap: simula una operación asíncrona por cada elemento
        // En este caso, cada palabra se convierte a mayúsculas y se retrasa un poco
        Flux.just("alpha", "beta", "gamma")
            .flatMap(s -> Mono.just(s.toUpperCase())
                              .delayElement(java.time.Duration.ofMillis(100)))
            .subscribe(System.out::println);

        // Esperar un poco para ver los resultados asíncronos
        try {
            Thread.sleep(500);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
        // Salida (el orden puede variar ligeramente debido al delay asíncrono):
        // BETA
        // ALPHA
        // GAMMA
    }
}
🔥 Importante: La diferencia clave entre `map` y `flatMap` es que `map` opera sincrónicamente sobre el valor, mientras que `flatMap` opera asincrónicamente, esperando por el `Publisher` interno.

Operadores de Filtrado

  • filter(): Emite solo los elementos que cumplen una condición.
  • take(): Toma un número específico de elementos y luego completa.
  • skip(): Omite un número específico de elementos al principio.
import reactor.core.publisher.Flux;

public class FilteringOperators {

    public static void main(String[] args) {
        // filter: solo números pares
        Flux.range(1, 10)
            .filter(i -> i % 2 == 0)
            .subscribe(System.out::println);
        // Salida:
        // 2
        // 4
        // 6
        // 8
        // 10

        System.out.println("------------------");

        // take: toma los primeros 3 elementos
        Flux.just("a", "b", "c", "d", "e")
            .take(3)
            .subscribe(System.out::println);
        // Salida:
        // a
        // b
        // c

        System.out.println("------------------");

        // skip: omite los primeros 2 elementos
        Flux.range(10, 5)
            .skip(2)
            .subscribe(System.out::println);
        // Salida:
        // 12
        // 13
        // 14
    }
}

Operadores de Combinación

  • zip(): Combina elementos de múltiples Publishers basándose en su orden, emitiendo una tupla cuando todos los Publishers han emitido un elemento en esa posición.
  • merge(): Combina flujos entrelazando sus elementos tan pronto como llegan. El orden no está garantizado si llegan al mismo tiempo de diferentes fuentes.
  • concat(): Combina flujos secuencialmente; el segundo Publisher no se suscribe hasta que el primero ha completado.
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

import java.time.Duration;

public class CombinationOperators {

    public static void main(String[] args) {
        // zip: combina elementos por índice
        Flux<String> names = Flux.just("Alice", "Bob");
        Flux<Integer> ages = Flux.just(30, 25);

        Flux.zip(names, ages, (name, age) -> name + " is " + age + " years old")
            .subscribe(System.out::println);
        // Salida:
        // Alice is 30 years old
        // Bob is 25 years old

        System.out.println("------------------");

        // merge: intercala elementos (el orden puede variar)
        Flux<String> flux1 = Flux.just("A", "B").delayElements(Duration.ofMillis(10));
        Flux<String> flux2 = Flux.just("X", "Y").delayElements(Duration.ofMillis(15));

        Flux.merge(flux1, flux2)
            .subscribe(System.out::println);
        // Esperar un poco
        try { Thread.sleep(100); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
        // Salida (ejemplo, puede variar):
        // A
        // X
        // B
        // Y

        System.out.println("------------------");

        // concat: el segundo flujo espera al primero
        Flux<String> fluxA = Flux.just("1", "2").delayElements(Duration.ofMillis(10));
        Flux<String> fluxB = Flux.just("3", "4").delayElements(Duration.ofMillis(10));

        Flux.concat(fluxA, fluxB)
            .subscribe(System.out::println);
        // Esperar un poco
        try { Thread.sleep(200); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
        // Salida:
        // 1
        // 2
        // 3
        // 4
    }
}

Operadores de Error Handling

  • onErrorReturn(): Retorna un valor por defecto si ocurre un error.
  • onErrorResume(): Retorna un Publisher alternativo si ocurre un error.
  • retry(): Reintenta la suscripción al Publisher si ocurre un error.
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

public class ErrorHandlingOperators {

    public static void main(String[] args) {
        // onErrorReturn: devuelve un valor predeterminado en caso de error
        Flux.range(1, 5)
            .map(i -> {
                if (i == 3) throw new RuntimeException("Error en el elemento 3");
                return i;
            })
            .onErrorReturn(0) // Retorna 0 si hay un error y completa
            .subscribe(System.out::println,
                       error -> System.err.println("Error capturado por subscribe: " + error.getMessage()));
        // Salida:
        // 1
        // 2
        // 0

        System.out.println("------------------");

        // onErrorResume: devuelve un Publisher alternativo en caso de error
        Flux.range(1, 5)
            .map(i -> {
                if (i == 4) throw new IllegalArgumentException("Error en el elemento 4");
                return i;
            })
            .onErrorResume(e -> {
                System.err.println("Error detectado: " + e.getMessage());
                return Flux.just(99, 100); // Continúa con un nuevo flujo
            })
            .subscribe(System.out::println);
        // Salida:
        // 1
        // 2
        // 3
        // Error detectado: Error en el elemento 4
        // 99
        // 100

        System.out.println("------------------");

        // retry: reintenta el flujo 2 veces si hay un error
        Flux.defer(() -> {
                System.out.println("Subscribing...");
                return Flux.range(1, 3)
                           .map(i -> {
                               if (i == 2) {
                                   System.out.println("Generating error...");
                                   throw new IllegalStateException("Simulated error");
                               }
                               return i;
                           });
            })
            .retry(2) // Intenta 2 veces más después del primer fallo
            .subscribe(System.out::println,
                       error -> System.err.println("Error final: " + error.getMessage()),
                       () -> System.out.println("Completado"));
        // Salida (ejemplo):
        // Subscribing...
        // 1
        // Generating error...
        // Subscribing...
        // 1
        // Generating error...
        // Subscribing...
        // 1
        // Generating error...
        // Error final: Simulated error
    }
}

🔄 Backpressure: Controlando el Flujo de Datos

Uno de los pilares de Reactive Streams (y por ende, de Project Reactor) es el backpressure. Es un mecanismo que permite al suscriptor notificar al publicador cuántos elementos está listo para recibir. Esto evita que el publicador inunde al suscriptor con más datos de los que puede manejar, previniendo problemas de memoria y rendimiento.

En Reactor, los operadores manejan el backpressure automáticamente. Cuando usas subscribe(), estás pidiendo una cantidad infinita de elementos por defecto. Puedes controlar esto si implementas un Subscriber personalizado.

import org.reactivestreams.Subscription;
import reactor.core.publisher.BaseSubscriber;
import reactor.core.publisher.Flux;

public class BackpressureExample {

    public static void main(String[] args) {
        Flux.range(1, 100)
            .subscribe(new BaseSubscriber<Integer>() {
                private int count = 0;
                private final int BATCH_SIZE = 10;

                @Override
                protected void hookOnSubscribe(Subscription subscription) {
                    System.out.println("Suscrito. Pidiendo " + BATCH_SIZE + " elementos.");
                    request(BATCH_SIZE); // Pedimos los primeros 10 elementos
                }

                @Override
                protected void hookOnNext(Integer value) {
                    System.out.println("Recibido: " + value);
                    count++;
                    if (count % BATCH_SIZE == 0) {
                        System.out.println("Procesados " + BATCH_SIZE + " elementos. Pidiendo otros " + BATCH_SIZE + "...");
                        request(BATCH_SIZE); // Pedimos otros 10 elementos
                    }
                }

                @Override
                protected void hookOnError(Throwable throwable) {
                    System.err.println("Error: " + throwable.getMessage());
                }

                @Override
                protected void hookOnComplete() {
                    System.out.println("Flujo completado.");
                }
            });
    }
}

En este ejemplo, el suscriptor solo pide 10 elementos a la vez. Una vez que ha procesado esos 10, pide los siguientes 10. Esto ilustra cómo el backpressure permite al suscriptor controlar el ritmo del publicador.

🧪 Probando Códigos Reactivos con StepVerifier

Probar código reactivo puede ser un desafío debido a su naturaleza asíncrona. Project Reactor ofrece StepVerifier, una herramienta fantástica para probar Flux y Mono de manera determinista.

StepVerifier te permite definir una secuencia esperada de eventos (onNext, onError, onComplete) y verificar que el Publisher los emite correctamente.

import org.junit.jupiter.api.Test;
import reactor.core.publisher.Flux;
import reactor.test.StepVerifier;

import java.time.Duration;

public class ReactorTest {

    // Simula un servicio que devuelve un flujo de nombres con un pequeño retraso
    Flux<String> getNamesAsync() {
        return Flux.just("Alice", "Bob", "Charlie")
                   .delayElements(Duration.ofMillis(50));
    }

    // Simula un servicio que devuelve un error después de un elemento
    Flux<String> getErrorFlux() {
        return Flux.just("data1")
                   .concatWith(Flux.error(new RuntimeException("Error simulado")));
    }

    @Test
    void testGetNamesAsync() {
        StepVerifier.create(getNamesAsync())
            .expectNext("Alice")
            .expectNext("Bob")
            .expectNext("Charlie")
            .verifyComplete(); // Espera la señal de completado
    }

    @Test
    void testErrorFlux() {
        StepVerifier.create(getErrorFlux())
            .expectNext("data1")
            .expectError(RuntimeException.class) // Espera una RuntimeException
            .verify(); // Verifica y espera errores o completado
    }

    @Test
    void testTransformAndFilter() {
        Flux<Integer> numbers = Flux.just(1, 2, 3, 4, 5);
        StepVerifier.create(numbers.filter(i -> i % 2 == 0).map(i -> i * 10))
            .expectNext(20)
            .expectNext(40)
            .verifyComplete();
    }
}
💡 Consejo: Para ejecutar estos tests, asegúrate de tener JUnit 5 (o tu framework de pruebas preferido) y `reactor-test` en tu classpath de pruebas.

🎯 Patrones Comunes y Mejores Prácticas

Encadenamiento de Operadores

Los operadores de Reactor están diseñados para ser encadenados, lo que permite una sintaxis fluida y declarativa.

Flux.range(1, 10)
    .filter(i -> i % 2 == 0) // Solo números pares
    .map(i -> "Número par: " + i) // Transformar a String
    .delayElements(Duration.ofMillis(100)) // Introducir un retraso artificial
    .subscribe(System.out::println, 
               error -> System.err.println("Error: " + error.getMessage()), 
               () -> System.out.println("Flujo de pares completado"));

try { Thread.sleep(1200); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }

Composición con zipWith y then

  • zipWith: Combina este Publisher con otro, emitiendo un Tuple2 de sus elementos.
  • then: Ejecuta un Mono después de que el Publisher actual haya completado, descartando sus elementos.
Mono<String> getUserById(String id) {
    return Mono.just("User-" + id).delayElement(Duration.ofMillis(50));
}

Mono<Integer> getOrderCountByUserId(String userId) {
    return Mono.just(userId.length() * 5).delayElement(Duration.ofMillis(70));
}

// Combinar información de usuario y conteo de órdenes
Mono<String> userAndOrderInfo = getUserById("123")
    .zipWith(getOrderCountByUserId("123"), 
             (user, orderCount) -> "Usuario: " + user + ", Órdenes: " + orderCount);

userAndOrderInfo.subscribe(System.out::println);
// Salida (después de ~120ms):
// Usuario: User-123, Órdenes: 15

// Realizar una acción después de que un flujo haya completado
Mono.just("Tarea A completada")
    .delayElement(Duration.ofMillis(200))
    .then(Mono.just("Tarea B iniciada"))
    .subscribe(System.out::println);
// Salida (después de ~200ms):
// Tarea B iniciada

try { Thread.sleep(300); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }

Concurrencia y Schedulers

Project Reactor utiliza Schedulers para controlar en qué hilo se ejecutan los operadores. Esto es fundamental para operaciones que no son bloqueantes por naturaleza (como cálculos intensivos o I/O).

  • Schedulers.immediate(): Ejecuta en el hilo actual.

  • Schedulers.boundedElastic(): Utiliza un pool de hilos elástico y con límites, bueno para operaciones de bloqueo (I/O).

  • Schedulers.parallel(): Utiliza un pool de hilos para tareas computacionalmente intensivas.

  • Schedulers.single(): Utiliza un único hilo reutilizable.

  • publishOn(): Afecta la ejecución de los operadores subsiguientes en el flujo.

  • subscribeOn(): Afecta la ejecución de la suscripción y, por ende, el origen del flujo.

Contexto: Scheduler A Contexto: Scheduler B Origen de Datos subscribeOn(Scheduler A) Operador 1 (Sch. A) publishOn(Scheduler B) Operador 2 (Sch. B) Suscripción (Sch. B)
import reactor.core.publisher.Flux;
import reactor.core.scheduler.Schedulers;

import java.time.Duration;

public class SchedulersExample {

    public static void main(String[] args) {
        Flux.range(1, 5)
            .map(i -> {
                System.out.println("Map 1 en hilo: " + Thread.currentThread().getName());
                return i * 2;
            })
            .subscribeOn(Schedulers.boundedElastic()) // La suscripción y la primera parte del flujo
            .publishOn(Schedulers.parallel()) // Los operadores posteriores en un pool paralelo
            .map(i -> {
                System.out.println("Map 2 en hilo: " + Thread.currentThread().getName());
                return "Resultado: " + i;
            })
            .blockLast(); // Bloquea hasta que el último elemento sea procesado (solo para ejemplos/tests)

        System.out.println("------------------");

        // Un ejemplo más práctico de un Mono que realiza una operación de bloqueo
        Mono.fromCallable(() -> {
                System.out.println("Operación bloqueante en hilo: " + Thread.currentThread().getName());
                Thread.sleep(100); // Simula una operación de I/O o cálculo largo
                return "Datos de I/O";
            })
            .subscribeOn(Schedulers.boundedElastic()) // Ejecuta la Callable en un hilo de I/O
            .map(s -> {
                System.out.println("Transformando datos en hilo: " + Thread.currentThread().getName());
                return s.toUpperCase();
            })
            .block(); // Bloquea para ver el resultado
    }
}

La salida mostrará cómo los hilos cambian según subscribeOn y publishOn.

✅ Conclusión

Project Reactor es una herramienta poderosa para construir aplicaciones Java asíncronas, no bloqueantes y resilientes. Al dominar Flux y Mono, junto con su rica colección de operadores, puedes escribir código más eficiente y fácil de razonar para manejar flujos de datos y eventos.

La curva de aprendizaje puede ser pronunciada al principio, pero la inversión vale la pena para desarrollar sistemas modernos que cumplan con las demandas de rendimiento y escalabilidad de hoy en día.

¡Experimenta con los operadores, prueba tus flujos con StepVerifier, y comienza a pensar reactivamente! La programación reactiva no es solo una moda, es una evolución necesaria para las arquitecturas de software modernas.

Esperamos que este tutorial te haya proporcionado una base sólida para comenzar tu viaje con Project Reactor.

Tutoriales relacionados

Comentarios (0)

Aún no hay comentarios. ¡Sé el primero!