← 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 pipelinesUni/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 |