Programación Reactiva: Flujos de Datos Asíncronos, Observables y Subjects
Este artículo es una introducción completa a la Programación Reactiva, incluyendo flujos de datos asíncronos, Observables, Subjects y ejemplos prácticos.
En Resumen
La Programación Reactiva se centra en flujos de datos asíncronos. En lugar de esperar, te suscribes a los flujos y reaccionas a los eventos conforme llegan.
Descripción Técnica Compacta
Programación Reactiva es un paradigma para procesar flujos de datos asíncronos mediante programas no bloqueantes orientados a eventos.
Principios Fundamentales:
- Asincronía: Las tareas se ejecutan en segundo plano sin bloquear el programa principal
- Flujos de Datos (Streams): Todo se considera como una secuencia de eventos
- Observer Pattern: Los flujos se suscriben en lugar de consultarse
- Reactive Streams: Estándar para procesar flujos asíncronos con backpressure
Conceptos Clave:
- Observable: Flujo de datos que puede emitir 0 a n valores
- Observer: Receptor de los datos del Observable
- Subscription: Conexión entre Observable y Observer
- Operators: Transformaciones y filtrado de flujos de datos
- Scheduler: Control sobre la ejecución en threads
- Backpressure: Protección contra sobrecarga cuando el productor es rápido
Puntos Clave para Examen
- Programación Reactiva: Procesamiento asincrónico y no bloqueante
- Observable: Flujo de datos con 0 a n valores
- Observer: Receptor de eventos de datos
- Subscription: Gestión de conexión entre Observable y Observer
- Operators: map, filter, flatMap para transformar flujos
- Backpressure: Protección contra sobrecarga, control de flujo
- Scheduler: Control de threads para operaciones asincrónicas
- Relevancia Profesional: Paradigma arquitectónico moderno para sistemas escalables
Componentes Principales
- Observable: Fuente de eventos de datos
- Observer: Receptor de datos
- Subscription: Conexión y gestión de recursos
- Operators: Transformaciones de flujos
- Subjects: Simultáneamente Observable y Observer
- Scheduler: Control de ejecución en threads
- Backpressure: Mecanismos de control de flujo
- Reactive Streams: Especificación estándar
Ejemplos Prácticos
1. Programación Reactiva con RxJava
import io.reactivex.rxjava3.core.*;
import io.reactivex.rxjava3.schedulers.Schedulers;
import java.util.concurrent.TimeUnit;
public class ReactiveProgrammingDemo {
public static void main(String[] args) throws InterruptedException {
// Observable simple
Observable<String> observable = Observable.just("Hola", "Reactiva", "World");
// Suscribirse a un Observable
observable.subscribe(
value -> System.out.println("Next: " + value),
error -> System.err.println("Error: " + error),
() -> System.out.println("Completed!")
);
// Operaciones asincrónicas
asyncOperations();
// Demostración de Operators
operatorsDemo();
// Backpressure con Flowable
backpressureDemo();
// Subjects para multicasting
subjectsDemo();
}
private static void asyncOperations() {
System.out.println("\n=== Operaciones Asincrónicas ===");
// Observable con Scheduler
Observable.fromArray("Task 1", "Task 2", "Task 3")
.subscribeOn(Schedulers.io()) // Ejecutar en thread IO
.observeOn(Schedulers.single()) // Observar en thread único
.subscribe(
task -> System.out.println("Procesando: " + task + " en " + Thread.currentThread().getName()),
error -> System.err.println("Error: " + error),
() -> System.out.println("Todas las tareas completadas")
);
// Observable basado en tiempo
Observable.interval(1, TimeUnit.SECONDS)
.take(5)
.map(tick -> "Tick " + (tick + 1))
.subscribe(
tick -> System.out.println(tick + " en " + Thread.currentThread().getName())
);
// Simular llamadas a API asincrónicas
Observable<String> apiCall = Observable.fromCallable(() -> {
Thread.sleep(1000); // Llamada simulada a API
return "Datos de API";
}).subscribeOn(Schedulers.io());
apiCall
.observeOn(Schedulers.single())
.subscribe(
data -> System.out.println("Resultado de API: " + data),
error -> System.err.println("Error de API: " + error)
);
}
private static void operatorsDemo() {
System.out.println("\n=== Demostración de Operators ===");
Observable.range(1, 10)
.filter(n -> n % 2 == 0) // Filtrar números pares
.map(n -> n * n) // Elevar al cuadrado
.take(3) // Solo los primeros 3 elementos
.subscribe(
result -> System.out.println("Resultado: " + result),
error -> System.err.println("Error: " + error),
() -> System.out.println("Demostración de operators completada")
);
// flatMap para transformación asincrónica
Observable.just("user1", "user2", "user3")
.flatMap(userId -> getUserData(userId)
.subscribeOn(Schedulers.io()) // Cada llamada en su propio thread
)
.subscribe(
userData -> System.out.println("Datos de Usuario: " + userData),
error -> System.err.println("Error de Usuario: " + error)
);
}
private static void backpressureDemo() {
System.out.println("\n=== Demo de Backpressure ===");
// Productor rápido con Flowable
Flowable.range(1, 1000)
.onBackpressureBuffer(100) // Buffer de tamaño 100
.observeOn(Schedulers.io())
.subscribe(
item -> {
try {
Thread.sleep(10); // Consumidor lento
System.out.println("Procesado: " + item);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
},
error -> System.err.println("Error de Backpressure: " + error)
);
}
private static void subjectsDemo() {
System.out.println("\n=== Demo de Subjects ===");
// PublishSubject para multicasting
PublishSubject<String> publishSubject = PublishSubject.create();
// Múltiples observadores suscritos
Observer<String> observer1 = new Observer<String>() {
@Override
public void onNext(String value) {
System.out.println("Observer 1: " + value);
}
@Override
public void onError(Throwable e) {
System.err.println("Observer 1 Error: " + e);
}
@Override
public void onComplete() {
System.out.println("Observer 1 Completado");
}
};
Observer<String> observer2 = new Observer<String>() {
@Override
public void onNext(String value) {
System.out.println("Observer 2: " + value);
}
@Override
public void onError(Throwable e) {
System.err.println("Observer 2 Error: " + e);
}
@Override
public void onComplete() {
System.out.println("Observer 2 Completado");
}
};
publishSubject.subscribe(observer1);
publishSubject.subscribe(observer2);
// Emitir datos
publishSubject.onNext("Mensaje 1");
publishSubject.onNext("Mensaje 2");
publishSubject.onComplete();
// BehaviorSubject para retener último valor
BehaviorSubject<Integer> behaviorSubject = BehaviorSubject.createDefault(0);
behaviorSubject.subscribe(
value -> System.out.println("Behavior Subject: " + value)
);
behaviorSubject.onNext(10);
behaviorSubject.onNext(20);
// Suscriptor tardío recibe el último valor
behaviorSubject.subscribe(
value -> System.out.println("Suscriptor Tardío: " + value)
);
}
// Simulación de llamada a API
private static Observable<String> getUserData(String userId) {
return Observable.fromCallable(() -> {
Thread.sleep(500);
return "UserData for " + userId;
});
}
}
2. Programación Reactiva con Project Reactor (Spring)
import reactor.core.publisher.*;
import reactor.core.scheduler.Schedulers;
import java.time.Duration;
import java.util.List;
public class ReactorDemo {
public static void main(String[] args) throws InterruptedException {
// Mono para 0..1 valores
Mono<String> mono = Mono.just("Hello Reactor");
mono.subscribe(
value -> System.out.println("Mono: " + value),
error -> System.err.println("Error: " + error)
);
// Flux para 0..n valores
Flux<String> flux = Flux.just("A", "B", "C", "D", "E");
flux
.filter(letter -> !"C".equals(letter))
.map(String::toUpperCase)
.subscribe(
letter -> System.out.println("Flux: " + letter),
error -> System.err.println("Flux Error: " + error),
() -> System.out.println("Flux Completed")
);
// Simulación de Web API asincrónica
webApiSimulation();
// Manejo de errores
errorHandling();
// Backpressure
backpressureHandling();
}
private static void webApiSimulation() {
System.out.println("\n=== Web API Simulation ===");
// Simula una solicitud web
Mono<String> webResponse = Mono.fromCallable(() -> {
Thread.sleep(1000); // Network latency
return "Response Data";
}).subscribeOn(Schedulers.boundedElastic());
webResponse
.map(response -> "Processed: " + response)
.timeout(Duration.ofSeconds(2))
.onErrorResume(error -> Mono.just("Fallback Response"))
.subscribe(
result -> System.out.println("Web Result: " + result),
error -> System.err.println("Web Error: " + error)
);
}
private static void errorHandling() {
System.out.println("\n=== Error Handling ===");
Flux<Integer> numbers = Flux.range(1, 10)
.map(n -> {
if (n == 5) {
throw new RuntimeException("Error at " + n);
}
return n * n;
});
numbers
.onErrorContinue((error, item) ->
System.err.println("Error processing " + item + ": " + error))
.subscribe(
result -> System.out.println("Result: " + result),
error -> System.err.println("Final Error: " + error),
() -> System.out.println("Error Handling Completed")
);
}
private static void backpressureHandling() {
System.out.println("\n=== Backpressure Handling ===");
// Productor rápido
Flux<Long> fastProducer = Flux.interval(Duration.ofMillis(10))
.take(100);
// Consumidor lento con backpressure
fastProducer
.onBackpressureBuffer(50) // Buffer con estrategia de descarte
.publishOn(Schedulers.boundedElastic())
.delayElements(Duration.ofMillis(50)) // Consumidor lento
.subscribe(
item -> System.out.println("Processed: " + item),
error -> System.err.println("Backpressure Error: " + error)
);
}
}
// Ejemplo de servicio reactivo
@Service
public class ReactiveUserService {
public Flux<User> getAllUsers() {
return Flux.fromIterable(getUserList())
.subscribeOn(Schedulers.parallel());
}
public Mono<User> getUserById(String id) {
return Mono.fromCallable(() -> findUserById(id))
.subscribeOn(Schedulers.boundedElastic())
.filter(user -> user != null)
.switchIfEmpty(Mono.error(new UserNotFoundException("User not found: " + id)));
}
public Flux<User> searchUsers(String query) {
return getAllUsers()
.filter(user -> user.getName().toLowerCase().contains(query.toLowerCase()))
.take(10); // Limita los resultados
}
// Fuente de datos simulada
private List<User> getUserList() {
return List.of(
new User("1", "Alice", "alice@example.com"),
new User("2", "Bob", "bob@example.com"),
new User("3", "Charlie", "charlie@example.com")
);
}
private User findUserById(String id) {
return getUserList().stream()
.filter(user -> user.getId().equals(id))
.findFirst()
.orElse(null);
}
}
class User {
private String id;
private String name;
private String email;
public User(String id, String name, String email) {
this.id = id;
this.name = name;
this.email = email;
}
// Getters
public String getId() { return id; }
public String getName() { return name; }
public String getEmail() { return email; }
}
3. Programación Reactiva con JavaScript (RxJS)
// RxJS Programación Reactiva
const { Observable, Subject, BehaviorSubject, fromEvent, interval, of } = require('rxjs');
const { map, filter, switchMap, take, debounceTime, distinctUntilChanged } = require('rxjs/operators');
// Observable simple
const helloObservable = of('Hello', 'Reactive', 'JavaScript');
helloObservable.subscribe({
next: value => console.log('Next:', value),
error: error => console.error('Error:', error),
complete: () => console.log('Complete!')
});
// Eventos del DOM como Observable
const button = document.createElement('button');
button.textContent = 'Click me!';
document.body.appendChild(button);
const clickObservable = fromEvent(button, 'click');
clickObservable
.pipe(
map(event => ({ type: 'click', timestamp: Date.now() })),
debounceTime(300),
distinctUntilChanged()
)
.subscribe(event => {
console.log('Button clicked:', event);
// Simula una llamada API
simulateApiCall();
});
// Llamadas API con patrón reactivo
function simulateApiCall() {
return new Observable(subscriber => {
console.log('API Call started...');
setTimeout(() => {
const success = Math.random() > 0.3;
if (success) {
subscriber.next({ data: 'API Response', status: 'success' });
subscriber.complete();
} else {
subscriber.error(new Error('API Error'));
}
}, 1000);
});
}
// Búsqueda con patrón reactivo
const searchInput = document.createElement('input');
searchInput.placeholder = 'Search...';
document.body.appendChild(searchInput);
const searchObservable = fromEvent(searchInput, 'input')
.pipe(
map(event => event.target.value),
filter(query => query.length >= 3),
debounceTime(500),
distinctUntilChanged(),
switchMap(query => searchApi(query))
);
searchObservable.subscribe({
next: results => console.log('Search results:', results),
error: error => console.error('Search error:', error)
});
function searchApi(query) {
return new Observable(subscriber => {
console.log('Searching for:', query);
setTimeout(() => {
const results = [`Result 1 for ${query}`, `Result 2 for ${query}`];
subscriber.next(results);
subscriber.complete();
}, 300);
});
}
// Subject para multicasting
const subject = new Subject();
// Múltiples suscriptores
const subscriber1 = {
next: value => console.log('Subscriber 1:', value),
error: error => console.error('Subscriber 1 Error:', error),
complete: () => console.log('Subscriber 1 Complete')
};
const subscriber2 = {
next: value => console.log('Subscriber 2:', value),
error: error => console.error('Subscriber 2 Error:', error),
complete: () => console.log('Subscriber 2 Complete')
};
subject.subscribe(subscriber1);
subject.subscribe(subscriber2);
// Emite datos
subject.next('Message 1');
subject.next('Message 2');
subject.complete();
// BehaviorSubject para retener el último valor
const behaviorSubject = new BehaviorSubject(0);
behaviorSubject.subscribe(value => console.log('Initial:', value));
behaviorSubject.next(10);
behaviorSubject.next(20);
// Suscriptor tardío recibe el último valor
behaviorSubject.subscribe(value => console.log('Late Subscriber:', value));
// WebSocket con patrón reactivo
class ReactiveWebSocket {
constructor(url) {
this.url = url;
this.messageSubject = new Subject();
this.connectionStatus = new BehaviorSubject('disconnected');
}
connect() {
this.connectionStatus.next('connecting');
// Simula una conexión WebSocket
this.socket = {
send: (data) => console.log('Sending:', data),
close: () => console.log('WebSocket closed')
};
// Simula mensajes entrantes
setInterval(() => {
const message = { type: 'data', payload: Math.random() };
this.messageSubject.next(message);
}, 2000);
this.connectionStatus.next('connected');
return this.connectionStatus.asObservable();
}
disconnect() {
if (this.socket) {
this.socket.close();
this.connectionStatus.next('disconnected');
}
}
sendMessage(message) {
if (this.socket) {
this.socket.send(JSON.stringify(message));
}
}
getMessages() {
return this.messageSubject.asObservable();
}
getConnectionStatus() {
return this.connectionStatus.asObservable();
}
}
// Usa WebSocket
const webSocket = new ReactiveWebSocket('ws://localhost:8080');
webSocket.connect().subscribe(status => {
console.log('Connection Status:', status);
});
webSocket.getMessages().subscribe(message => {
console.log('Received Message:', message);
});
// Envía un mensaje
setTimeout(() => {
webSocket.sendMessage({ type: 'ping', data: 'Hello Server' });
}, 5000);
4. Programación Reactiva con Python (RxPY)
import rx
from rx import operators as ops
import time
import threading
from datetime import datetime
# Observable simple
def simple_observable():
source = rx.of("Python", "Reactive", "Programming")
source.subscribe(
on_next=lambda value: print(f"Next: {value}"),
on_error=lambda error: print(f"Error: {error}"),
on_completed=lambda: print("Completed!")
)
# Operaciones asincrónicas
def async_operations():
print("\n=== Operaciones asincrónicas ===")
# Observable basado en tiempo
rx.interval(1.0).pipe(
ops.take(5),
ops.map(lambda i: f"Tick {i + 1}")
).subscribe(
on_next=lambda value: print(f"{value} at {datetime.now().second}s"),
on_completed=lambda: print("Timer completed!")
)
# Simulación de llamadas asincrónicas a API
def simulate_api_call(item):
def observer(observer):
def worker():
time.sleep(1) # Simula latencia de red
if item % 3 == 0:
observer.on_next(f"Data for {item}")
observer.on_completed()
else:
observer.on_error(Exception(f"Error for {item}"))
thread = threading.Thread(target=worker)
thread.start()
return rx.create(observer)
rx.range(1, 6).pipe(
ops.map(simulate_api_call),
ops.merge_all()
).subscribe(
on_next=lambda value: print(f"API Result: {value}"),
on_error=lambda error: print(f"API Error: {error}")
)
# Demostración de operadores
def operators_demo():
print("\n=== Demostración de operadores ===")
rx.range(1, 11).pipe(
ops.filter(lambda x: x % 2 == 0),
ops.map(lambda x: x * x),
ops.take(3)
).subscribe(
on_next=lambda value: print(f"Result: {value}"),
on_completed=lambda: print("Operators completed!")
)
# flatMap para transformación asincrónica
def get_user_data(user_id):
return rx.of(f"UserData for {user_id}").pipe(
ops.delay(0.5) # Simula latencia
)
rx.of("user1", "user2", "user3").pipe(
ops.flat_map(get_user_data)
).subscribe(
on_next=lambda data: print(f"User: {data}")
)
# Subject para multicasting
def subjects_demo():
print("\n=== Demostración de Subject ===")
# PublishSubject
subject = rx.Subject()
# Múltiples suscriptores
def observer1(value):
print(f"Observer 1: {value}")
def observer2(value):
print(f"Observer 2: {value}")
subject.subscribe(observer1)
subject.subscribe(observer2)
# Emitir datos
subject.on_next("Message 1")
subject.on_next("Message 2")
subject.on_completed()
# BehaviorSubject
behavior_subject = rx.BehaviorSubject(0)
behavior_subject.subscribe(
on_next=lambda value: print(f"Behavior Subject: {value}")
)
behavior_subject.on_next(10)
behavior_subject.on_next(20)
# Suscriptor tardío
behavior_subject.subscribe(
on_next=lambda value: print(f"Late Subscriber: {value}")
)
# Observable Hot vs Cold
def hot_cold_demo():
print("\n=== Observable Hot vs Cold ===")
# Observable Cold (cada suscriptor recibe todos los valores)
cold = rx.of("A", "B", "C")
print("Cold Observable:")
cold.subscribe(on_next=lambda x: print(f"Subscriber 1: {x}"))
time.sleep(1)
cold.subscribe(on_next=lambda x: print(f"Subscriber 2: {x}"))
# Observable Hot (valores compartidos)
hot = rx.Subject()
print("\nHot Observable:")
hot.subscribe(on_next=lambda x: print(f"Subscriber 1: {x}"))
hot.on_next("X")
hot.on_next("Y")
hot.subscribe(on_next=lambda x: print(f"Subscriber 2: {x}"))
hot.on_next("Z")
if __name__ == "__main__":
simple_observable()
async_operations()
operators_demo()
subjects_demo()
hot_cold_demo()
# Mantén el programa en ejecución para operaciones asincrónicas
time.sleep(10)
Backpressure en Reactive Streams
Estrategias de Backpressure
// RxJava Backpressure
Flowable.range(1, 1000)
.onBackpressureBuffer() // Buffer (estándar)
.onBackpressureDrop() // Descartar elementos excedentes
.onBackpressureLatest() // Mantener solo el elemento más reciente
.onBackpressureError() // Lanzar error si hay sobrecarga
.subscribe(item -> processItem(item));
// Project Reactor Backpressure
Flux.range(1, 1000)
.onBackpressureBuffer(100) // Buffer con tamaño
.onBackpressureDrop() // Descartar elementos
.onBackpressureLatest() // Solo el último elemento
.limitRate(100) // Limitación de velocidad
.subscribe(item -> processItem(item));
Scheduler para Control de Threads
RxJava Scheduler
Observable.just("data")
.subscribeOn(Schedulers.io()) // Ejecución en thread de IO
.observeOn(Schedulers.single()) // Observación en single thread
.observeOn(Schedulers.computation()) // Cálculos en thread de CPU
.subscribe(result -> handleResult(result));
Project Reactor Scheduler
Mono.just("data")
.subscribeOn(Schedulers.boundedElastic()) // Operaciones de IO
.publishOn(Schedulers.parallel()) // Procesamiento paralelo
.publishOn(Schedulers.single()) // Single thread
.subscribe(result -> handleResult(result));
Ventajas e Inconvenientes
Ventajas de la Programación Reactiva
- Asincronía: Procesamiento no bloqueante
- Escalabilidad: Mejor uso de recursos del sistema
- Responsividad: Tiempos de respuesta más rápidos
- Flexibilidad: Composición sencilla de operaciones
- Manejo de errores: Gestión centralizada de excepciones
Inconvenientes
- Complejidad: Curva de aprendizaje pronunciada
- Debugging: Búsqueda de errores más difícil
- Overhead: Capa adicional de abstracción
- Gestión de recursos: Lifecycle más complejo
Preguntas frecuentes en exámenes
-
¿Cuál es la diferencia entre Observable y Subject? Observable es solo una fuente, mientras que Subject puede actuar tanto como fuente como receptor.
-
Explica qué es Backpressure. Mecanismo para proteger contra sobrecarga cuando un productor rápido trabaja con un consumidor lento.
-
¿Cuándo se utiliza Programación Reactiva? En operaciones asincrónicas, aplicaciones en tiempo real y sistemas con altos requisitos de escalabilidad.
-
¿Cuál es el propósito de los Scheduler? Controlar la ejecución en threads y el procesamiento paralelo en sistemas reactivos.



