Saltar a contenido

← Volver al índice | Anterior: Decisión Dapr vs Pekko

Investigación Completa: Quarkus Mutiny — Framework Reactivo

Tipo: Investigación Técnica Exhaustiva — Reactive Streams
Audiencia: Ingeniería (Arquitectura, Backend)
Fecha: 05 de Abril de 2026
Stack: Quarkus 3.32 · Java 25 · SmallRye Reactive Messaging · RabbitMQ · Dapr
Fuentes: smallrye.io/mutiny, quarkus.io/guides, investigación original


0. Contexto: ¿Por Qué Mutiny?

En la investigación de Dapr Building Blocks vs Pekko se concluyó que Dapr no ofrece Reactive Streams integrado como Pekko Streams. Sin embargo, Mutiny —la librería reactiva nativa de Quarkus— puede cubrir esa funcionalidad combinándose con Dapr.

Mutiny NO es un clon de RxJava ni de Project Reactor. Su filosofía es: - API fluida e imperativa, no funcional-académica - Dos tipos fundamentales: Uni<T> (0..1 item) y Multi<T> (0..N items) - Lazy por defecto: nada se ejecuta hasta que hay un suscriptor - Integración nativa con Vert.x, SmallRye Reactive Messaging y Quarkus REST


1. Catálogo Completo de Operadores

1.1 Operadores de Uni<T>

Categoría Operador Descripción
Creación Uni.createFrom().item(T) Crea un Uni con un valor fijo
Uni.createFrom().completionStage(Supplier) Bridge desde CompletableFuture
Uni.createFrom().emitter(Consumer) Bridge imperativo con callback
Uni.createFrom().failure(Throwable) Uni que falla inmediatamente
Uni.createFrom().nullItem() Uni que emite null
Uni.createFrom().voidItem() Uni que emite Void
Transformación .onItem().transform(Function) Mapeo síncrono de T → R
.onItem().transformToUni(Function) FlatMap asíncrono (encadena otro Uni)
.onItem().transformToMulti(Function) Convierte item a stream
.onItem().castTo(Class) Cast de tipo
Acciones laterales .onItem().invoke(Consumer) Side-effect sin alterar el item
.onItem().call(Function<T, Uni<?>>) Side-effect asíncrono (espera finalización)
Delays .onItem().delayIt().by(Duration) Retarda la emisión
.onItem().delayIt().until(Function) Retarda hasta que un Uni completa
Failure .onFailure().recoverWithItem(T) Fallback con valor fijo
.onFailure().recoverWithUni(Function) Fallback asíncrono
.onFailure().retry().atMost(n) Reintentos simples
.onFailure().retry().withBackOff(init, max) Exponential backoff
.onFailure().transform(Function) Transforma la excepción
Timeout .ifNoItem().after(Duration).fail() Falla si no hay respuesta
.ifNoItem().after(Duration).recoverWithItem(T) Fallback por timeout
Threading .emitOn(Executor) Emite el resultado en otro thread
.runSubscriptionOn(Executor) Ejecuta la suscripción en otro thread
Combinación Uni.combine().all().unis(u1, u2).asTuple() Espera todos los Unis
Uni.join().all(u1, u2).andCollectFailures() Paralelo con recolección de fallos
Uni.join().first(u1, u2).withItem() Primero que emita
Conversión .await().atMost(Duration) Bloquea (solo tests/worker threads)
.await().indefinitely() Bloquea indefinidamente (⚠️ tests only)
.subscribeAsCompletionStage() Convierte a CompletionStage

1.2 Operadores de Multi<T>

Categoría Operador Descripción
Creación Multi.createFrom().items(T...) Stream desde varargs
Multi.createFrom().iterable(Iterable) Stream desde colección
Multi.createFrom().range(start, end) Rango de enteros
Multi.createFrom().ticks().every(Duration) Ticker periódico
Multi.createFrom().emitter(Consumer) Bridge imperativo
Multi.createFrom().generator(Supplier, BiFunction) Generador con estado
Transformación .onItem().transform(Function) Map síncrono
.onItem().transformToUniAndMerge(f, concurrency) FlatMap paralelo con concurrencia máxima
.onItem().transformToUniAndConcatenate(f) FlatMap secuencial (preserva orden)
.onItem().transformToMultiAndMerge(f) Sub-streams mergeados
.onItem().transformToMultiAndConcatenate(f) Sub-streams concatenados
.onItem().castTo(Class) Cast de tipos
Filtrado .select().where(Predicate) Filtro síncrono
.select().when(Function<T, Uni<Boolean>>) Filtro asíncrono
.select().distinct() Eliminar duplicados (⚠️ memoria)
.select().first(n) Tomar primeros N
.select().last(n) Tomar últimos N
.skip().first(n) Saltar primeros N
.skip().first(Duration) Saltar durante un período
.skip().repetitions() Dedup consecutivo (safe para streams infinitos)
Agrupación .group().intoLists().of(n) Batch por cantidad
.group().intoLists().every(Duration) Batch por tiempo
.group().intoMultis().of(n) Sub-streams por cantidad
.group().by(Function) Agrupación por clave → GroupedMulti
Combinación Multi.createBy().merging().streams(m1, m2) Merge intercalando
Multi.createBy().concatenating().streams(m1, m2) Concatenación secuencial
Multi.createBy().combining().streams(m1, m2).using(f) Zip/Combine
Acciones laterales .onItem().invoke(Consumer) Side-effect por item
.onItem().call(Function<T, Uni<?>>) Side-effect asíncrono por item
.onCompletion().invoke(Runnable) Al completar el stream
.onSubscription().invoke(Consumer) Al suscribirse
.log() Debug logging
Back-Pressure .onOverflow().buffer(n) Buffer de N items
.onOverflow().drop() Descartar exceso
.onOverflow().dropPreviousItems() Mantener solo el más reciente
.onOverflow().invoke(Consumer) Callback en overflow
.capDemandsTo(long) Limitar demanda del subscriber
.capDemandsUsing(Function) Demanda custom
.paceDemand() Control temporal de demanda
Recolección .collect().asList() Uni<List<T>>
.collect().asMap(keyMapper) Uni<Map<K, T>>
.collect().with(Collector) Usar collector de Java Streams
.collect().first() Uni<T> primer item
.collect().last() Uni<T> último item
Error Handling .onFailure().recoverWithItem(T) Fallback y termina
.onFailure().recoverWithMulti(f) Switch a otro stream
.onFailure().retry().atMost(n) Reintentos
.onFailure().retry().withBackOff(init, max) Exponential backoff
Threading .emitOn(Executor) Emite items en otro pool
.runSubscriptionOn(Executor) Ejecuta subscription upstream en otro pool
Broadcast .broadcast().toAllSubscribers() Hot stream multicasting
.broadcast().toAtLeast(n) Espera N suscriptores antes de emitir

2. Back-Pressure en Multi

2.1 Modelo Conceptual

Mutiny implementa Reactive Streams (ahora java.util.concurrent.Flow). Esto significa: 1. El subscriber indica cuántos items puede procesar (request(n))
2. El publisher no emite más de lo solicitado
3. Si el publisher genera items más rápido, se activa la estrategia de overflow

2.2 Estrategias de Overflow

// Strategy 1: Buffer (default limitado)
multi.onOverflow().buffer(256)          // Buffer de 256 items
     .subscribe().with(this::process);

// Strategy 2: Drop silencioso  
multi.onOverflow().drop()               // Descarta items no consumidos
     .subscribe().with(this::process);

// Strategy 3: Mantener solo el último
multi.onOverflow().dropPreviousItems()  // Solo el más reciente

// Strategy 4: Fail (default sin overflow)
// Si no se especifica onOverflow, falla con BackPressureFailure

// Strategy 5: Custom callback
multi.onOverflow().invoke(item -> metrics.incrementDropped())
     .drop()
     .subscribe().with(this::process);

2.3 Comparativa: Back-Pressure Pekko Streams vs Mutiny

Aspecto Pekko Streams Mutiny Multi
Modelo Materializer + Graph stages Reactive Streams / Flow API
Buffer .buffer(n, OverflowStrategy) .onOverflow().buffer(n)
Drop OverflowStrategy.dropHead/dropTail/dropBuffer .onOverflow().drop() / .dropPreviousItems()
Fail OverflowStrategy.fail Comportamiento default
Backpressure OverflowStrategy.backpressure Default con request(n)
Demand control Fused stages automáticos .capDemandsTo(n) / .paceDemand()
Complejidad Alta (materializer, graph DSL) Baja (API fluida)
Bounded buffer bufferSize en materializer settings .onOverflow().buffer(n)

Veredicto MVP Backlog: Mutiny cubre el 90% de los casos de back-pressure de Pekko Streams con una API significativamente más simple. Para los casos restantes (graph topologies complejas), no los necesitamos en el MVP Backlog.


3. Integración con SmallRye Reactive Messaging

3.1 Flujo de Mensajes: RabbitMQ → Pipeline Reactivo

RabbitMQ Queue (pet-flows-events)
  → SmallRye RabbitMQ Connector  
    → @Incoming("encounter-events") method  
      → Mutiny pipeline (transform, filter, retry, etc.)  
        → Business logic / Dapr Workflow trigger

3.2 Patrón Processor: Entrada → Transformación → Salida

@ApplicationScoped
public class EventProcessor {

    @Incoming("encounter-events")
    @Outgoing("processed-events")
    public Multi<ProcessedEvent> process(Multi<RawEvent> events) {
        return events
            .select().where(e -> e.type().equals("encounter.abandoned"))
            .onItem().transform(this::enrichEvent)
            .onFailure().retry().withBackOff(
                Duration.ofMillis(100), Duration.ofSeconds(1))
            .atMost(3);
    }
}

3.3 Patrón Consumer con Uni (Back-Pressure Natural)

@ApplicationScoped
public class EncounterEventConsumer {

    @Incoming("encounter-events")
    public Uni<Void> consume(Message<JsonObject> message) {
        // Retornar Uni<Void> = back-pressure natural
        // SmallRye NO pide el siguiente mensaje hasta que este Uni completa
        return processEvent(message.getPayload())
            .onItem().invoke(() -> message.ack())
            .onFailure().invoke(ex -> {
                log.error("Failed to process event", ex);
                message.nack(ex);
            })
            .replaceWithVoid();
    }
}

Clave: Retornar Uni<Void> desde un @Incoming = back-pressure automática. SmallRye espera a que finalice antes de pedir el siguiente mensaje al broker.

3.4 Failure Strategies en SmallRye RabbitMQ

# Si el procesamiento falla:
mp.messaging.incoming.encounter-events.failure-strategy=reject
# Opciones: fail (default), reject, requeue
Estrategia Comportamiento Equivalente Pekko
fail Detiene el stream completo Supervision.stop
reject Rechaza el mensaje (dead letter) Manual DLQ forwarding
requeue Devuelve al broker para reintento N/A (manual)

4. Concurrencia y Virtual Threads

4.1 Modelo de Threading en Quarkus

┌──────────────────────────────────────────┐
│  Event Loop Threads (Vert.x)             │
│  • Pocos threads (= cores)               │
│  • NUNCA bloquear aquí                   │
│  • Manejan I/O no bloqueante            │
├──────────────────────────────────────────┤
│  Worker Thread Pool                       │
│  • Para operaciones @Blocking            │
│  • Thread pool fijo (default 20)         │
├──────────────────────────────────────────┤
│  Virtual Threads (Java 25)               │
│  • @RunOnVirtualThread                   │
│  • Baratos de crear y bloquear           │
│  • Complementarios a Mutiny              │
└──────────────────────────────────────────┘

4.2 emitOn vs runSubscriptionOn

// emitOn: DÓNDE se procesan los items downstream
multi.emitOn(Infrastructure.getDefaultExecutor())  
    .onItem().transform(this::heavyComputation); // Ejecuta en worker pool

// runSubscriptionOn: DÓNDE se ejecuta la suscripción upstream
Uni.createFrom().item(() -> blockingDatabaseCall())
    .runSubscriptionOn(Infrastructure.getDefaultExecutor()); // Ejecuta la query en worker

// Con Virtual Threads (Java 25)
var vtExecutor = Executors.newVirtualThreadPerTaskExecutor();
multi.emitOn(vtExecutor)
    .onItem().transform(this::blockingLegacyCall); 

4.3 Cuándo Usar Cada Enfoque

Escenario Enfoque Recomendado
REST endpoint que devuelve datos reactivos Retornar Uni<T> directamente (Quarkus suscribe)
REST con lógica bloqueante (legacy JDBC) @RunOnVirtualThread en el endpoint
Pipeline reactivo que llama API bloqueante .runSubscriptionOn(vtExecutor) o .emitOn(vtExecutor)
Stream processing puro (transform, filter) Event loop directo (sin thread switching)
Batch processing masivo (50K notificaciones) transformToUniAndMerge(f, concurrency) con virtual threads

4.4 ⚠️ Regla de Oro

// ❌ NUNCA hacer esto en el event loop
Uni.createFrom().item(() -> {
    Thread.sleep(5000);  // BLOQUEA EL EVENT LOOP
    return result;
});

// ✅ Correcto: offload a virtual threads
Uni.createFrom().item(() -> {
    Thread.sleep(5000); 
    return result;
}).runSubscriptionOn(Executors.newVirtualThreadPerTaskExecutor());

// ✅ Aún mejor: usar API no bloqueante
Uni.createFrom().nullItem()
    .onItem().delayIt().by(Duration.ofSeconds(5))
    .onItem().transform(v -> result);

5. Error Handling

5.1 Anatomía Completa de onFailure()

// === Retry Simple ===
uni.onFailure().retry().atMost(3);

// === Retry con Exponential Backoff ===
uni.onFailure()
   .retry()
   .withBackOff(Duration.ofMillis(100), Duration.ofSeconds(2))  // initial, max
   .withJitter(0.2)  // ±20% aleatorio (evita thundering herd)
   .atMost(5);       // máximo 5 intentos

// === Retry con Deadline (no por intentos, por tiempo) ===
uni.onFailure()
   .retry()
   .withBackOff(Duration.ofMillis(200))
   .expireIn(Duration.ofSeconds(30));  // Reintentar durante máximo 30s

// === Retry Condicional ===
uni.onFailure(IOException.class)  // Solo para IOException
   .retry().atMost(3);

uni.onFailure(ex -> ex instanceof TimeoutException)
   .retry().atMost(3);

// === Recovery (Fallback) ===
uni.onFailure().recoverWithItem("valor-por-defecto");
uni.onFailure().recoverWithUni(ex -> callSecondaryService());

// === Transform Exception ===
uni.onFailure().transform(ex -> new BusinessException("Error procesando", ex));

5.2 Circuit Breaker con SmallRye Fault Tolerance

Mutiny NO tiene circuit breaker nativo. Se usa SmallRye Fault Tolerance (declarativo):

import org.eclipse.microprofile.faulttolerance.CircuitBreaker;
import org.eclipse.microprofile.faulttolerance.Retry;
import org.eclipse.microprofile.faulttolerance.Timeout;

@ApplicationScoped
public class NotificationService {

    @CircuitBreaker(
        requestVolumeThreshold = 4,    // Evaluar cada 4 requests
        failureRatio = 0.5,            // Abrir si >50% fallan
        delay = 5000,                  // Esperar 5s antes de half-open
        successThreshold = 2           // 2 éxitos para cerrar
    )
    @Retry(maxRetries = 3, delay = 200)
    @Timeout(2000)  // 2 segundos máximo
    public Uni<NotificationResult> sendPush(PushRequest request) {
        return gorushClient.send(request);
    }
}

5.3 Dead Letter Pattern

@Incoming("encounter-events")
public Uni<Void> consume(Message<JsonObject> msg) {
    return processEvent(msg.getPayload())
        .onItem().invoke(() -> msg.ack())
        .onFailure().invoke(ex -> {
            log.error("Moving to DLQ", ex);
            // SmallRye + RabbitMQ: nack con requeue=false envía a DLX
            msg.nack(ex);  
        })
        .replaceWithVoid();
}
# En application.properties — configurar Dead Letter Exchange
mp.messaging.incoming.encounter-events.dead-letter-exchange=pet.events.dlx
mp.messaging.incoming.encounter-events.dead-letter-routing-key=encounter.failed
mp.messaging.incoming.encounter-events.failure-strategy=reject

6. Batching y Windowing

6.1 Equivalencias Directas

Pekko Streams Mutiny Multi Notas
groupedWithin(n, duration) .group().intoLists().of(n) + .group().intoLists().every(d) Mutiny separa por count y por tiempo (no ambos en un solo operador)
grouped(n) .group().intoLists().of(n) Equivalente directo
groupedWeightedWithin(...) No disponible nativo Requiere implementación custom

6.2 Batch por Cantidad

Multi<List<Recipient>> batches = recipients
    .group().intoLists().of(100);  // Batches de 100

batches.onItem().transformToUniAndMerge(batch -> 
    sendBatchNotification(batch), 
    5  // máximo 5 batches en paralelo
);

6.3 Batch por Tiempo

Multi<List<Event>> windows = events
    .group().intoLists().every(Duration.ofSeconds(5));
    // Emite lo acumulado cada 5 segundos (puede ser lista vacía)

6.4 Batch por Tiempo Y Cantidad (Workaround)

Mutiny NO tiene groupedWithin(n, duration) nativo. Solución con timer + count:

/**
 * Custom groupedWithin: emite cuando se acumulan {@code maxSize} items
 * O cuando pasan {@code maxDuration} desde el último batch.
 */
public static <T> Multi<List<T>> groupedWithin(
        Multi<T> source, int maxSize, Duration maxDuration) {

    return Multi.createFrom().emitter(emitter -> {
        var buffer = new CopyOnWriteArrayList<T>();
        var lock = new ReentrantLock();

        // Timer periódico para flush por tiempo
        var ticker = Multi.createFrom().ticks().every(maxDuration)
            .subscribe().with(tick -> {
                lock.lock();
                try {
                    if (!buffer.isEmpty()) {
                        emitter.emit(new ArrayList<>(buffer));
                        buffer.clear();
                    }
                } finally {
                    lock.unlock();
                }
            });

        source.subscribe().with(
            item -> {
                lock.lock();
                try {
                    buffer.add(item);
                    if (buffer.size() >= maxSize) {
                        emitter.emit(new ArrayList<>(buffer));
                        buffer.clear();
                    }
                } finally {
                    lock.unlock();
                }
            },
            emitter::fail,
            emitter::complete
        );
    });
}

Recomendación MVP Backlog: Usar .group().intoLists().of(n) para el caso de 50K notificaciones (batch por cantidad). El windowing por tiempo + count se reserva para producción si surge la necesidad.


7. Fan-Out / Fan-In

7.1 Paralelismo Controlado (Fan-Out)

// El operador CLAVE para paralelismo controlado:
Multi<NotificationResult> results = recipients
    .onItem().transformToUniAndMerge(
        recipient -> sendNotification(recipient),  // Cada item → Uni
        10  // máximo 10 en paralelo
    );

Este es el equivalente directo a mapAsync(10, f) de Pekko Streams.

7.2 Fan-In: Merge de Múltiples Streams

// Merge intercalado (items llegan en cualquier orden)
Multi<Event> merged = Multi.createBy().merging()
    .streams(pushStream, emailStream, smsStream);

// Concatenación (primero completa uno, luego el siguiente)
Multi<Event> sequential = Multi.createBy().concatenating()
    .streams(pushStream, emailStream, smsStream);

7.3 Combinación de Unis (Join)

// Esperar que TODOS completen
Uni<Tuple3<A, B, C>> combined = Uni.combine().all()
    .unis(fetchUserUni, fetchPrefsUni, fetchTokenUni)
    .asTuple();

// Usar resultado combinado
combined.onItem().transform(tuple -> {
    var user = tuple.getItem1();
    var prefs = tuple.getItem2();
    var token = tuple.getItem3();
    return new NotificationPayload(user, prefs, token);
});

// Esperar el PRIMERO que complete
Uni<String> fastest = Uni.join()
    .first(primaryServiceUni, fallbackServiceUni)
    .withItem();

7.4 Comparativa Fan-Out/Fan-In

Patrón Pekko Streams Mutiny
mapAsync (paralelo) Flow[T].mapAsync(n)(f) .transformToUniAndMerge(f, n)
mapAsync (ordenado) Flow[T].mapAsync(n)(f) (default) .transformToUniAndConcatenate(f)
mapAsyncUnordered Flow[T].mapAsyncUnordered(n)(f) .transformToUniAndMerge(f, n)
Merge streams Source.merge(s1, s2) Multi.createBy().merging().streams(...)
Concat streams Source.concat(s1, s2) Multi.createBy().concatenating().streams(...)
Zip s1.zip(s2) Multi.createBy().combining().streams(s1, s2).using(f)
Join Unis N/A (Futures) Uni.combine().all().unis(...)

8. Conversión Imperativo ↔ Reactivo

8.1 De Imperativo a Reactivo

// Desde valor
Uni<String> uni = Uni.createFrom().item("hello");

// Desde CompletableFuture (LAZY - usa Supplier)
Uni<String> uni = Uni.createFrom().completionStage(
    () -> httpClient.sendAsync(request, BodyHandlers.ofString())
);

// Desde callback/listener
Uni<String> uni = Uni.createFrom().emitter(em -> {
    legacyService.doAsync(result -> em.complete(result));
});

// Multi desde iterador
Multi<User> multi = Multi.createFrom().iterable(userList);

8.2 De Reactivo a Imperativo

// ⚠️ SOLO en worker threads o virtual threads, NUNCA en event loop
String result = uni.await().atMost(Duration.ofSeconds(5));

// Conversión a CompletionStage
CompletionStage<String> future = uni.subscribeAsCompletionStage();

// Multi → List
List<String> list = multi.collect().asList()
    .await().atMost(Duration.ofSeconds(10));

8.3 ⚠️ Cómo NO Bloquear el Event Loop de Vert.x

// ❌ INCORRECTO — bloquea event loop
@GET
@Path("/users")
public List<User> getUsers() {
    return userService.findAll()  // Uni<List<User>>
        .await().indefinitely();  // 💀 DEADLOCK POTENCIAL
}

// ✅ CORRECTO — devuelve Uni, Quarkus suscribe
@GET
@Path("/users")
public Uni<List<User>> getUsers() {
    return userService.findAll();
}

// ✅ ALTERNATIVA — Virtual Threads para código imperativo
@GET
@Path("/users")
@RunOnVirtualThread
public List<User> getUsers() {
    return userService.findAllBlocking();  // OK en virtual thread
}

9. Testing de Pipelines Reactivos

9.1 Testing de Uni con UniAssertSubscriber

import io.smallrye.mutiny.helpers.test.UniAssertSubscriber;

@Test
void testUniTransformation() {
    Uni<String> uni = Uni.createFrom().item("hello")
        .onItem().transform(String::toUpperCase);

    UniAssertSubscriber<String> subscriber = uni
        .subscribe().withSubscriber(UniAssertSubscriber.create());

    subscriber
        .assertCompleted()
        .assertItem("HELLO");
}

@Test
void testUniFailure() {
    Uni<String> uni = Uni.createFrom()
        .failure(new RuntimeException("boom"));

    UniAssertSubscriber<String> subscriber = uni
        .subscribe().withSubscriber(UniAssertSubscriber.create());

    subscriber.assertFailedWith(RuntimeException.class, "boom");
}

9.2 Testing de Multi con AssertSubscriber

import io.smallrye.mutiny.helpers.test.AssertSubscriber;

@Test
void testMultiFiltering() {
    Multi<Integer> multi = Multi.createFrom().items(1, 2, 3, 4, 5)
        .select().where(n -> n % 2 == 0);

    AssertSubscriber<Integer> subscriber = multi
        .subscribe().withSubscriber(AssertSubscriber.create(10));

    subscriber
        .assertCompleted()
        .assertItems(2, 4);
}

@Test
void testMultiBatching() {
    Multi<List<Integer>> batched = Multi.createFrom().range(1, 11)
        .group().intoLists().of(3);

    AssertSubscriber<List<Integer>> subscriber = batched
        .subscribe().withSubscriber(AssertSubscriber.create(10));

    subscriber.assertCompleted();
    // Batches: [1,2,3], [4,5,6], [7,8,9], [10]
    assertEquals(4, subscriber.getItems().size());
}

9.3 Testing Asíncrono (con delays)

@Test
void testDelayedUni() {
    Uni<String> delayed = Uni.createFrom().item("result")
        .onItem().delayIt().by(Duration.ofMillis(100));

    UniAssertSubscriber<String> subscriber = delayed
        .subscribe().withSubscriber(UniAssertSubscriber.create());

    // Esperar a que complete
    subscriber.awaitItem(Duration.ofSeconds(1))
        .assertItem("result");
}

9.4 Testing con Quarkus @QuarkusTest

@QuarkusTest
class NotificationServiceTest {

    @Inject
    NotificationService service;

    @Test
    void testSendBatch() {
        var recipients = List.of(
            new Recipient("token1"), new Recipient("token2")
        );

        var result = service.sendBatch(recipients)
            .await().atMost(Duration.ofSeconds(5));

        assertEquals(2, result.successCount());
    }
}

10. Patrones Avanzados

10.1 Hot vs Cold Streams

// COLD: cada subscriber inicia su propia ejecución
Multi<Long> cold = Multi.createFrom().ticks().every(Duration.ofSeconds(1));
// Subscriber A ve: 0, 1, 2, 3...
// Subscriber B ve: 0, 1, 2, 3... (independiente)

// HOT: una sola ejecución compartida
Multi<Long> hot = Multi.createFrom().ticks().every(Duration.ofSeconds(1))
    .broadcast().toAllSubscribers();
// Subscriber A (conecta en t=0): 0, 1, 2, 3...
// Subscriber B (conecta en t=2): 2, 3, 4... (pierde 0, 1)

10.2 BroadcastProcessor (Event Bus interno)

@ApplicationScoped
public class ProgressBroadcaster {

    private final BroadcastProcessor<ProgressEvent> processor = 
        BroadcastProcessor.create();

    // Publicar (side que genera progreso)
    public void emit(ProgressEvent event) {
        processor.onNext(event);
    }

    // Suscribirse (SSE endpoint)
    public Multi<ProgressEvent> stream() {
        return processor;
    }
}

10.3 Emitter Imperativo

@ApplicationScoped
public class EventBridge {

    private final Multi<Event> stream;
    private final MultiEmitter<? super Event> emitter;

    public EventBridge() {
        // Crear un Multi con emitter
        AtomicReference<MultiEmitter<? super Event>> ref = new AtomicReference<>();
        this.stream = Multi.createFrom().emitter(em -> ref.set(em))
            .broadcast().toAllSubscribers();
        this.emitter = ref.get();
    }

    // Método imperativo
    public void push(Event event) {
        emitter.emit(event);
    }

    // Para SSE
    public Multi<Event> subscribe() {
        return stream;
    }
}

10.4 Overflow Strategies para Hot Streams

// Hot stream con diferentes estrategias
Multi<SensorData> sensorFeed = sensorSource
    .onOverflow().buffer(1000)        // Bufferear hasta 1000
    .broadcast().toAllSubscribers();

// Para subscribers lentos: drop
sensorFeed
    .onOverflow().drop()              // Si no puedo procesar, descarto
    .onItem().transform(this::enrichData);

// Para UI: solo el último valor importa
sensorFeed
    .onOverflow().dropPreviousItems() // Solo el más reciente
    .onItem().transform(this::formatForUI);

11. Integración con JAX-RS/RESTEasy Reactive

11.1 Endpoints que Devuelven Uni/Multi

@Path("/notifications")
@ApplicationScoped
public class NotificationResource {

    @Inject NotificationService service;

    // Uni → Response individual
    @POST
    @Path("/send")
    public Uni<Response> send(NotificationRequest request) {
        return service.send(request)
            .onItem().transform(result -> 
                Response.ok(result).build())
            .onFailure().recoverWithItem(ex -> 
                Response.serverError().entity(ex.getMessage()).build());
    }

    // Multi → JSON Array
    @GET
    @Path("/history")
    @Produces(MediaType.APPLICATION_JSON)
    public Multi<NotificationLog> getHistory() {
        return service.getHistory();  // Streamed como JSON array
    }
}

11.2 Server-Sent Events (SSE) con Multi

@Path("/progress")
@ApplicationScoped
public class ProgressResource {

    @Inject ProgressBroadcaster broadcaster;

    // SSE endpoint — cada item del Multi = un evento SSE
    @GET
    @Path("/stream")
    @RestStreamElementType(MediaType.APPLICATION_JSON)
    public Multi<ProgressEvent> streamProgress() {
        return broadcaster.stream();
    }

    // SSE con custom event types
    @GET
    @Path("/custom-stream")
    @Produces(MediaType.SERVER_SENT_EVENTS)
    public Multi<OutboundSseEvent> customStream(@Context Sse sse) {
        return broadcaster.stream()
            .onItem().transform(event -> sse.newEventBuilder()
                .name(event.type())                    // event: type
                .id(event.id())                        // id: xxx
                .data(event.toJson())                  // data: {...}
                .reconnectDelay(3000)                  // retry: 3000
                .build());
    }
}

11.3 SSE + Progress para Batch de 50K

@GET
@Path("/batch-progress/{jobId}")
@RestStreamElementType(MediaType.APPLICATION_JSON)
public Multi<BatchProgress> batchProgress(@PathParam("jobId") String jobId) {
    return broadcaster.stream()
        .select().where(p -> p.jobId().equals(jobId))
        .onCompletion().invoke(() -> log.info("SSE client disconnected"));
}

12. Métricas y Observabilidad

12.1 OpenTelemetry con Pipelines Reactivos

<!-- pom.xml -->
<dependency>
    <groupId>io.quarkus</groupId>
    <artifactId>quarkus-opentelemetry</artifactId>
</dependency>
# application.properties
quarkus.otel.exporter.otlp.endpoint=http://localhost:4317
quarkus.otel.metrics.enabled=true
quarkus.otel.traces.enabled=true

12.2 Spans Automáticos y Manuales

import io.opentelemetry.instrumentation.annotations.WithSpan;
import io.opentelemetry.instrumentation.annotations.SpanAttribute;

@ApplicationScoped
public class NotificationService {

    // Span automático
    @WithSpan("send-notification")
    public Uni<NotificationResult> send(
            @SpanAttribute("notification.type") String type,
            @SpanAttribute("notification.recipient") String recipient) {
        return gorushClient.send(type, recipient);
    }
}

12.3 Métricas Custom en Pipelines

@ApplicationScoped
public class BatchMetrics {

    @Inject MeterRegistry registry;

    private final AtomicLong processed = new AtomicLong();
    private final AtomicLong failed = new AtomicLong();

    public <T> Multi<T> instrument(Multi<T> pipeline, String name) {
        return pipeline
            .onItem().invoke(item -> {
                processed.incrementAndGet();
                registry.counter("batch.processed", "name", name).increment();
            })
            .onFailure().invoke(ex -> {
                failed.incrementAndGet();
                registry.counter("batch.failed", "name", name).increment();
            });
    }
}

12.4 Context Propagation

Quarkus propaga automáticamente el contexto de tracing a través de operaciones reactivas usando MicroProfile Context Propagation. Esto incluye: - Trace context entre emitOn() / runSubscriptionOn() - Headers de tracing en mensajes SmallRye Reactive Messaging - Spans parent/child correctos en pipelines Uni/Multi

<!-- La dependencia viene incluida en quarkus-opentelemetry -->
<dependency>
    <groupId>io.quarkus</groupId>
    <artifactId>quarkus-smallrye-context-propagation</artifactId>
</dependency>

13. Anti-Patrones y Gotchas

13.1 Catálogo de Anti-Patrones

# Anti-Patrón Impacto Solución
1 Bloquear el Event Loop Deadlock completo de la app @RunOnVirtualThread, emitOn(), runSubscriptionOn()
2 Olvidar suscribirse Código nunca se ejecuta (pipeline lazy) Retornar Uni/Multi al framework, o .subscribe() explícito
3 await().indefinitely() en producción Bloqueo indefinido, stall await().atMost(Duration) solo en tests
4 No manejar failure en Multi Stream termina silenciosamente .onFailure().recoverWith... o aislar con transformToUniAndMerge
5 select().distinct() en streams infinitos OutOfMemoryError Usar .skip().repetitions() (solo compara con el anterior)
6 Transformar sin efecto Pipeline crea nuevo Multi que nadie consume Siempre retornar o re-asignar la nueva referencia
7 Thread pool propio para I/O Over-threading, contention Usar virtual thread executor o Infrastructure.getDefaultExecutor()
8 Mezclar blocking JDBC con Mutiny Event loop bloqueado Usar Hibernate Reactive o @RunOnVirtualThread
9 Retry sin backoff Thundering herd, sobrecarga del servicio .withBackOff().withJitter() siempre
10 Ignorar back-pressure BackPressureFailure en producción .onOverflow().buffer(n) explícito

13.2 Gotchas Específicos de Quarkus + Mutiny

Gotcha 1: El pipeline lazy que "no funciona"

// ❌ El transform se crea pero NADIE lo consume
void doSomething() {
    uni.onItem().transform(x -> x.toUpperCase()); // Resultado perdido
}

// ✅ Correcto
Uni<String> doSomething() {
    return uni.onItem().transform(x -> x.toUpperCase());
}

Gotcha 2: Failure terminal en Multi

// ❌ Si un item falla, TODO el stream se detiene
multi.onItem().transformToUniAndMerge(item -> 
    riskyOperation(item), 10);  // Un fallo = stream muerto

// ✅ Aislar fallos por item
multi.onItem().transformToUniAndMerge(item -> 
    riskyOperation(item)
        .onFailure().recoverWithItem(
            ex -> new ErrorResult(item, ex)), // Encapsular error
    10);

Gotcha 3: Uni.combine no es realmente paralelo por sí solo

// ⚠️ Si callA() y callB() son síncronos, no hay paralelismo real
Uni.combine().all().unis(callA(), callB()).asTuple();
// callA() se CREATE primero, luego callB(), pero la ejecución real
// depende de si son async internamente (HTTP reactivo, etc.)

14. Caso de Uso: Pipeline de 50.000 Notificaciones

14.1 Requisitos

  • Leer 50.000 destinatarios desde base de datos
  • Agrupar en batches de 100
  • Throttle: 200ms de pausa entre batches
  • Concurrencia: máximo 5 batches simultáneos
  • Retry: exponential backoff por item fallido
  • Recolección: contar éxitos, fallos, pendientes
  • Progreso: SSE en tiempo real al frontend

14.2 Diseño del Pipeline

graph LR
    A[DB: 50K Recipients] --> B[Multi stream]
    B --> C[group().intoLists().of 100]
    C --> D[delay 200ms entre batches]
    D --> E[transformToUniAndMerge concurrency=5]
    E --> F[sendBatch con retry]
    F --> G[collect results]
    G --> H[emit progress via SSE]

14.3 Implementación Completa

@ApplicationScoped
public class MassNotificationService {

    @Inject NotificationClient notificationClient;
    @Inject ProgressBroadcaster progress;
    @Inject MeterRegistry metrics;

    private static final int BATCH_SIZE = 100;
    private static final int MAX_CONCURRENCY = 5;
    private static final Duration THROTTLE_DELAY = Duration.ofMillis(200);
    private static final int MAX_RETRIES = 3;

    /**
     * Envía notificaciones masivas con throttling, batching,
     * retry y progreso en tiempo real.
     */
    public Uni<BatchJobResult> sendMassNotification(
            String jobId, 
            List<Recipient> recipients, 
            NotificationTemplate template) {

        var totalCount = recipients.size();
        var stats = new AtomicBatchStats(totalCount);

        return Multi.createFrom().iterable(recipients)

            // 1. BATCHING: agrupar en lotes de 100
            .group().intoLists().of(BATCH_SIZE)

            // 2. THROTTLING: delay entre batches  
            .onItem().call(batch -> 
                Uni.createFrom().nullItem()
                   .onItem().delayIt().by(THROTTLE_DELAY))

            // 3. FAN-OUT: procesar hasta 5 batches en paralelo
            .onItem().transformToUniAndMerge(
                batch -> processBatch(jobId, batch, template, stats),
                MAX_CONCURRENCY
            )

            // 4. COLLECT: acumular resultados
            .collect().asList()

            // 5. AGGREGATE: generar resultado final
            .onItem().transform(results -> {
                var finalResult = new BatchJobResult(
                    jobId, 
                    stats.successCount(), 
                    stats.failureCount(),
                    stats.total()
                );
                progress.emit(ProgressEvent.completed(jobId, finalResult));
                return finalResult;
            })

            // 6. ERROR GLOBAL: si algo catastrófico falla
            .onFailure().invoke(ex -> {
                log.error("Mass notification job {} failed", jobId, ex);
                progress.emit(ProgressEvent.failed(jobId, ex.getMessage()));
            });
    }

    /**
     * Procesa un batch individual con retry por item.
     */
    private Uni<BatchResult> processBatch(
            String jobId,
            List<Recipient> batch,
            NotificationTemplate template,
            AtomicBatchStats stats) {

        return Multi.createFrom().iterable(batch)
            // Enviar cada notificación con retry
            .onItem().transformToUniAndMerge(
                recipient -> sendSingleWithRetry(recipient, template),
                10  // 10 envíos paralelos dentro del batch
            )
            // Recolectar resultados del batch
            .collect().asList()
            .onItem().invoke(results -> {
                long ok = results.stream().filter(SendResult::success).count();
                long fail = results.size() - ok;

                stats.addSuccess(ok);
                stats.addFailure(fail);

                // Emitir progreso
                progress.emit(new ProgressEvent(
                    jobId, 
                    "batch_complete",
                    stats.processedCount(), 
                    stats.total(),
                    stats.successCount(),
                    stats.failureCount()
                ));

                metrics.counter("notifications.sent", 
                    "status", "success").increment(ok);
                metrics.counter("notifications.sent", 
                    "status", "failed").increment(fail);
            })
            .onItem().transform(results -> 
                new BatchResult(results.size(), 
                    results.stream().filter(SendResult::success).count()));
    }

    /**
     * Envía una notificación individual con exponential backoff.
     */
    private Uni<SendResult> sendSingleWithRetry(
            Recipient recipient, NotificationTemplate template) {

        return notificationClient.send(recipient, template)
            .onItem().transform(resp -> SendResult.success(recipient.id()))
            .onFailure().retry()
                .withBackOff(Duration.ofMillis(100), Duration.ofSeconds(2))
                .withJitter(0.2)
                .atMost(MAX_RETRIES)
            .onFailure().recoverWithItem(ex -> {
                log.warn("Failed after {} retries: {}", 
                    MAX_RETRIES, recipient.id(), ex);
                return SendResult.failure(recipient.id(), ex.getMessage());
            });
    }
}

14.4 SSE Endpoint para Progreso

@Path("/jobs")
@ApplicationScoped
public class JobProgressResource {

    @Inject ProgressBroadcaster broadcaster;

    @GET
    @Path("/{jobId}/progress")
    @RestStreamElementType(MediaType.APPLICATION_JSON)
    public Multi<ProgressEvent> streamProgress(
            @PathParam("jobId") String jobId) {
        return broadcaster.stream()
            .select().where(e -> e.jobId().equals(jobId));
    }
}

14.5 Frontend: Consumir SSE

const eventSource = new EventSource('/jobs/job-123/progress');

eventSource.onmessage = (event) => {
    const data = JSON.parse(event.data);
    updateProgressBar(data.processed, data.total);
    updateStats(data.successCount, data.failureCount);

    if (data.type === 'completed' || data.type === 'failed') {
        eventSource.close();
    }
};

14.6 Records de Soporte

// Records para el pipeline
record Recipient(String id, String token, String channel) {}
record SendResult(String recipientId, boolean success, String error) {
    static SendResult success(String id) { return new SendResult(id, true, null); }
    static SendResult failure(String id, String error) { return new SendResult(id, false, error); }
}
record BatchResult(long total, long successful) {}
record BatchJobResult(String jobId, long success, long failures, long total) {}
record ProgressEvent(String jobId, String type, long processed, long total,
                     long successCount, long failureCount) {
    static ProgressEvent completed(String jobId, BatchJobResult r) {
        return new ProgressEvent(jobId, "completed", r.total(), r.total(), r.success(), r.failures());
    }
    static ProgressEvent failed(String jobId, String error) {
        return new ProgressEvent(jobId, "failed", 0, 0, 0, 0);
    }
}
record AtomicBatchStats(long total) {
    // Implementación con AtomicLong para thread-safety
    private static final AtomicLong success = new AtomicLong();
    private static final AtomicLong failure = new AtomicLong();
    void addSuccess(long n) { success.addAndGet(n); }
    void addFailure(long n) { failure.addAndGet(n); }
    long successCount() { return success.get(); }
    long failureCount() { return failure.get(); }
    long processedCount() { return success.get() + failure.get(); }
}

15. Recomendaciones: MVP Backlog vs Producción

Capacidad MVP Backlog (Ahora) Producción (Futuro)
Batching .group().intoLists().of(n) Custom groupedWithin(n, duration)
Back-Pressure Uni<Void> return (natural) .onOverflow().buffer() + métricas
Retry .retry().withBackOff().atMost(3) + Circuit breaker + DLQ
Concurrency transformToUniAndMerge(f, 5) + Rate limiter externo (Dapr/Redis)
Throttling .onItem().call(delay) Token bucket con Vert.x timer
Monitoring .log() + logging OpenTelemetry spans + Micrometer gauges
Testing UniAssertSubscriber + Integration tests con Testcontainers
SSE Progress BroadcastProcessor + WebSocket para bidireccional
Dead Letter nack() + RabbitMQ DLX + Retry queue escalonada

16. Dependencias Maven Requeridas

<!-- Ya incluidas con quarkus-resteasy-reactive-jackson -->
<!-- Mutiny viene automáticamente con Quarkus -->

<!-- Para Reactive Messaging + RabbitMQ -->
<dependency>
    <groupId>io.quarkus</groupId>
    <artifactId>quarkus-messaging-rabbitmq</artifactId>
</dependency>

<!-- Para Fault Tolerance (Circuit Breaker, Retry declarativo) -->
<dependency>
    <groupId>io.quarkus</groupId>
    <artifactId>quarkus-smallrye-fault-tolerance</artifactId>
</dependency>

<!-- Para OpenTelemetry -->
<dependency>
    <groupId>io.quarkus</groupId>
    <artifactId>quarkus-opentelemetry</artifactId>
</dependency>

<!-- Para testing -->
<dependency>
    <groupId>io.smallrye.reactive</groupId>
    <artifactId>smallrye-mutiny-vertx-core</artifactId>
    <scope>test</scope>
</dependency>

Documentos Relacionados

Nivel Documento Descripción
Complementario Investigación Dapr Building Blocks vs Pekko 12+ building blocks, State Management, Virtual Actors, Resiliencia, Throttling
Decisión Dapr vs Pekko & Spring Decisión de arquitectura original
Concepto Motor de Flujos Dinámicos Patrón Intérprete sobre Dapr Workflows
Spike Spike Dapr Workflows Informe de viabilidad Workflow-as-Code