Skip to content
IRC-CodingIRC-Coding
Reactive ProgrammingObservablesSubjectsRxJavaProject ReactorAsynchronous Streams

Reactive Programming: Observables, Subjects y Streams

Guía de Reactive Programming con asynchronous streams. Observables, Subjects, operadores, backpressure y ejemplos en RxJava, Project Reactor y JavaScript.

S

schutzgeist

11 min read
Reactive Programming: Observables, Subjects y Streams

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

  1. Observable: Fuente de eventos de datos
  2. Observer: Receptor de datos
  3. Subscription: Conexión y gestión de recursos
  4. Operators: Transformaciones de flujos
  5. Subjects: Simultáneamente Observable y Observer
  6. Scheduler: Control de ejecución en threads
  7. Backpressure: Mecanismos de control de flujo
  8. 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

  1. ¿Cuál es la diferencia entre Observable y Subject? Observable es solo una fuente, mientras que Subject puede actuar tanto como fuente como receptor.

  2. Explica qué es Backpressure. Mecanismo para proteger contra sobrecarga cuando un productor rápido trabaja con un consumidor lento.

  3. ¿Cuándo se utiliza Programación Reactiva? En operaciones asincrónicas, aplicaciones en tiempo real y sistemas con altos requisitos de escalabilidad.

  4. ¿Cuál es el propósito de los Scheduler? Controlar la ejecución en threads y el procesamiento paralelo en sistemas reactivos.

Fuentes principales

  1. https://reactivex.io/
  2. https://projectreactor.io/
  3. https://github.com/ReactiveX/RxJava
Volver al blog
Share:

Entradas relacionadas