Arquitectura

Redis como broker de eventos: patrones Pub/Sub y Streams para microservicios sin Kafka

Ruslan Ismailov Publicado 12 min de lectura
R

Introducción: cuándo Kafka es excesivo y por qué Redis se ha vuelto una opción popular en 2026

Kafka es una herramienta poderosa, pero su despliegue y operación requieren recursos considerables: un clúster separado de ZooKeeper o KRaft, configuración de replicación, monitoreo de brokers y capacitación del equipo. Para startups, equipos pequeños o sistemas de microservicios con carga moderada, estos costos a menudo no se justifican.

En 2026, Redis se utiliza cada vez más como broker de eventos en arquitecturas event-driven. Las razones son varias: Redis ya está presente en la mayoría de los stacks de producción como caché, es bien conocido por los desarrolladores, sencillo de operar y ofrece dos mecanismos maduros de mensajería: Pub/Sub y Streams. Si tu equipo no cuenta con experiencia de infraestructura dedicada y el volumen de mensajes se mide en miles y no en millones por segundo, Redis como alternativa a Kafka merece una consideración seria.

Redis Pub/Sub: mecánica de funcionamiento y garantías de entrega

Cómo funciona Pub/Sub

Redis Pub/Sub es el clásico modelo de publicación-suscripción en memoria. El publicador envía un mensaje a un canal mediante el comando PUBLISH, y todos los suscriptores activos de ese canal lo reciben al instante a través de una conexión persistente. Sin buffering: si el suscriptor no está conectado en el momento de la publicación, el mensaje se pierde de forma irrecuperable.

Redis admite tanto canales exactos como suscripción por patrón mediante PSUBSCRIBE, lo que permite suscribirse a grupos de canales por máscara, por ejemplo orders.*.

Garantías de entrega — su ausencia

Pub/Sub en Redis funciona bajo el principio fire-and-forget: sin persistencia, sin confirmaciones, sin reentregas. No es un bug, sino una característica de diseño. Si el consumidor cae, se reinicia o simplemente es más lento que el productor, los mensajes se pierden. Esta semántica se denomina «at most once».

Cuándo es apropiado usar Pub/Sub

  • Invalidación de caché: notificaciones de cambios de datos donde la pérdida de un evento no es crítica
  • Notificaciones en tiempo real en la UI: estado en línea de usuarios, chats
  • Eventos de difusión sin garantías: métricas del sistema, mensajes de heartbeat
  • Prototipado de flujos event-driven antes de implementar una solución más compleja

Redis Streams: arquitectura, consumer groups y mecanismo ACK

Qué son Redis Streams

Redis Streams, introducidos en Redis 5.0, son una estructura de datos persistente y ordenada similar a un log append-only. Cada entrada tiene un ID único con el formato 1704067200000-0 (milisegundos + número secuencial). Los streams se almacenan en memoria y pueden replicarse y persistirse mediante los mecanismos estándar RDB/AOF.

Consumer Groups

La diferencia clave respecto a Pub/Sub es el soporte de consumer groups. Un grupo de consumidores procesa el stream de forma colaborativa: cada mensaje es entregado a exactamente un miembro del grupo. Esto permite escalar horizontalmente el procesamiento: varios workers leen del mismo stream sin duplicar trabajo. Los diferentes grupos reciben todos los mensajes de forma independiente, similar a los topics de Kafka con múltiples grupos de consumidores.

Mecanismo ACK y pending entries

Tras recibir un mensaje, el consumidor debe confirmar explícitamente su procesamiento con el comando XACK. Hasta la confirmación, el mensaje permanece en la Pending Entries List (PEL). Si el consumidor cae, los mensajes no procesados pueden reasignarse a otro worker mediante XCLAIM. Esto garantiza la semántica «at least once».

Persistencia y retención

Los streams admiten limitación de longitud mediante MAXLEN: se pueden conservar los últimos N mensajes o recortar por tiempo. Con AOF o RDB habilitados, los datos sobreviven a los reinicios de Redis. Sin embargo, es importante entender que Redis es una base de datos in-memory, y ante una pérdida total de datos (sin replicación ni persistencia) el stream se perderá.

Comparativa: Redis Pub/Sub vs Redis Streams vs Kafka

A continuación, los parámetros clave de las tres soluciones que importan al elegir un broker de eventos para una arquitectura de microservicios.

  • Persistencia: Pub/Sub — no; Streams — sí (RDB/AOF); Kafka — sí (log en disco)
  • Garantías de entrega: Pub/Sub — at most once; Streams — at least once; Kafka — at least once / exactly once
  • Consumer groups: Pub/Sub — no; Streams — sí; Kafka — sí
  • Reprocesamiento: Pub/Sub — imposible; Streams — sí, mediante XCLAIM; Kafka — sí, reset de offset
  • Orden de mensajes: Pub/Sub — no garantizado; Streams — garantizado dentro del stream; Kafka — garantizado dentro de la partición
  • Complejidad operativa: Pub/Sub — mínima; Streams — baja; Kafka — alta
  • Throughput: Pub/Sub — alto; Streams — alto (cientos de miles/seg); Kafka — muy alto (millones/seg)
  • Almacenamiento a largo plazo: Pub/Sub — no; Streams — limitado (RAM); Kafka — prácticamente ilimitado
  • Dependencias de infraestructura: Pub/Sub — solo Redis; Streams — solo Redis; Kafka — Kafka + ZooKeeper/KRaft

Patrón Event-Driven con Redis Streams: implementación en Go

Productor

Ejemplo de productor en Go utilizando la biblioteca go-redis/v9. El productor publica un evento de creación de pedido en el stream orders.

package main

import (
    "context"
    "fmt"
    "log"
    "time"

    "github.com/redis/go-redis/v9"
)

func main() {
    ctx := context.Background()

    rdb := redis.NewClient(&redis.Options{
        Addr: "localhost:6379",
    })
    defer rdb.Close()

    // Publicamos el evento en el stream
    id, err := rdb.XAdd(ctx, &redis.XAddArgs{
        Stream: "orders",
        MaxLen: 10000,    // límite de longitud del stream
        Approx: true,     // recorte aproximado para mayor rendimiento
        Values: map[string]interface{}{
            "order_id":   "ORD-12345",
            "user_id":    "USR-67890",
            "amount":     "1500.00",
            "currency":   "EUR",
            "created_at": time.Now().Unix(),
        },
    }).Result()

    if err != nil {
        log.Fatalf("XAdd error: %v", err)
    }

    fmt.Printf("Event published, ID: %s\n", id)
}

Consumidor con consumer group

El consumidor crea el grupo (si no existe), lee los mensajes y confirma el procesamiento mediante XACK.

package main

import (
    "context"
    "fmt"
    "log"
    "os"
    "os/signal"
    "syscall"
    "time"

    "github.com/redis/go-redis/v9"
)

const (
    streamName = "orders"
    groupName  = "order-processor"
)

func main() {
    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()

    rdb := redis.NewClient(&redis.Options{
        Addr: "localhost:6379",
    })
    defer rdb.Close()

    consumerName := fmt.Sprintf("worker-%s", os.Getenv("POD_NAME"))

    // Creamos el grupo si no existe
    err := rdb.XGroupCreateMkStream(ctx, streamName, groupName, "$").Err()
    if err != nil && err.Error() != "BUSYGROUP Consumer Group name already exists" {
        log.Fatalf("XGroupCreate error: %v", err)
    }

    // Manejo de señales para graceful shutdown
    sigCh := make(chan os.Signal, 1)
    signal.Notify(sigCh, syscall.SIGTERM, syscall.SIGINT)

    go func() {
        <-sigCh
        cancel()
    }()

    log.Printf("Consumer %s started, group: %s", consumerName, groupName)

    for {
        select {
        case <-ctx.Done():
            log.Println("Shutting down gracefully")
            return
        default:
        }

        // Leemos hasta 10 mensajes con timeout de 2 segundos
        streams, err := rdb.XReadGroup(ctx, &redis.XReadGroupArgs{
            Group:    groupName,
            Consumer: consumerName,
            Streams:  []string{streamName, ">"},
            Count:    10,
            Block:    2 * time.Second,
        }).Result()

        if err != nil {
            if err == redis.Nil || err.Error() == "redis: nil" {
                continue // timeout, no hay nuevos mensajes
            }
            if ctx.Err() != nil {
                return
            }
            log.Printf("XReadGroup error: %v", err)
            time.Sleep(time.Second)
            continue
        }

        for _, stream := range streams {
            for _, msg := range stream.Messages {
                if err := processOrder(msg); err != nil {
                    log.Printf("Processing error for %s: %v", msg.ID, err)
                    // El mensaje permanece en PEL para reprocesamiento
                    continue
                }

                // Confirmamos el procesamiento exitoso
                if err := rdb.XAck(ctx, streamName, groupName, msg.ID).Err(); err != nil {
                    log.Printf("XAck error for %s: %v", msg.ID, err)
                }
            }
        }
    }
}

func processOrder(msg redis.XMessage) error {
    fmt.Printf("Processing order: %s, amount: %s\n",
        msg.Values["order_id"],
        msg.Values["amount"],
    )
    // Lógica de negocio para procesar el pedido
    return nil
}

Manejo de errores y Dead Letter: qué hacer con los mensajes no procesados

Uno de los problemas típicos al trabajar con Redis Streams son los mensajes «bloqueados» en la PEL. Un mensaje puede acabar ahí si el consumidor cayó antes de llamar a XACK o si el procesamiento falla repetidamente.

Estrategia de gestión de la PEL

El patrón recomendado es ejecutar periódicamente una tarea que, mediante XPENDING, verifique los mensajes «bloqueados» por más tiempo del umbral definido (por ejemplo, 5 minutos) y los transfiera a otro consumidor mediante XCLAIM. Tras N intentos fallidos, el mensaje se mueve a un stream separado «papelera» — el dead letter stream.

func reclaimStalePendingMessages(ctx context.Context, rdb *redis.Client) {
    minIdleTime := 5 * time.Minute

    pending, err := rdb.XPendingExt(ctx, &redis.XPendingExtArgs{
        Stream: streamName,
        Group:  groupName,
        Start:  "-",
        End:    "+",
        Count:  100,
    }).Result()

    if err != nil {
        log.Printf("XPendingExt error: %v", err)
        return
    }

    for _, p := range pending {
        if p.Idle < minIdleTime {
            continue
        }

        if p.RetryCount >= 5 {
            // Movemos al dead letter stream
            moveToDeadLetter(ctx, rdb, p.ID)
            rdb.XAck(ctx, streamName, groupName, p.ID)
            continue
        }

        // Reasignamos al worker actual
        rdb.XClaim(ctx, &redis.XClaimArgs{
            Stream:   streamName,
            Group:    groupName,
            Consumer: "reclaim-worker",
            MinIdle:  minIdleTime,
            Messages: []string{p.ID},
        })
    }
}

func moveToDeadLetter(ctx context.Context, rdb *redis.Client, msgID string) {
    rdb.XAdd(ctx, &redis.XAddArgs{
        Stream: "orders:dead-letter",
        Values: map[string]interface{}{
            "original_id": msgID,
            "moved_at":    time.Now().Unix(),
        },
    })
    log.Printf("Message %s moved to dead letter stream", msgID)
}

Escalado de consumidores en Kubernetes

Despliegue de workers

Cada Pod-worker en Kubernetes utiliza la variable de entorno POD_NAME como identificador único del consumidor en el grupo. Redis balancea la carga entre los miembros del grupo automáticamente, sin necesidad de coordinación adicional.

apiVersion: apps/v1
kind: Deployment
metadata:
  name: order-processor
spec:
  replicas: 3
  selector:
    matchLabels:
      app: order-processor
  template:
    metadata:
      labels:
        app: order-processor
    spec:
      containers:
        - name: worker
          image: myapp/order-processor:latest
          env:
            - name: POD_NAME
              valueFrom:
                fieldRef:
                  fieldPath: metadata.name
            - name: REDIS_ADDR
              value: redis-service:6379
          resources:
            requests:
              cpu: 100m
              memory: 64Mi
            limits:
              cpu: 500m
              memory: 256Mi

Autoescalado según la longitud del stream

KEDA (Kubernetes Event-Driven Autoscaling) soporta Redis Streams de forma nativa. El scaler lee la longitud de la PEL o el número de mensajes no procesados y escala el Deployment automáticamente.

apiVersion: keda.sh/v1alpha1
kind: ScaledObject
metadata:
  name: order-processor-scaler
spec:
  scaleTargetRef:
    name: order-processor
  minReplicaCount: 1
  maxReplicaCount: 20
  triggers:
    - type: redis-streams
      metadata:
        address: redis-service:6379
        stream: orders
        consumerGroup: order-processor
        pendingEntriesCount: "50"  # escalamos si PEL > 50

Monitoreo de Redis como broker: métricas clave e instrumentos

Métricas clave

  • Longitud del stream (XLEN orders) — un crecimiento continuo indica que los consumidores no dan abasto
  • Tamaño de la PEL (XPENDING orders order-processor - + 1) — cantidad de mensajes no procesados
  • Lag del grupo — diferencia entre el último ID en el stream y el último ID confirmado por el grupo
  • used_memory — Redis almacena todo en RAM; el desbordamiento de memoria es crítico
  • connected_clients — número de conexiones activas de consumidores
  • evicted_keys — si Redis expulsa claves por falta de memoria, los datos pueden perderse

Herramientas

Para monitorear Redis como broker de eventos, utiliza la combinación Redis Exporter + Prometheus + Grafana. Redis Exporter (oliver006/redis_exporter) exporta las métricas redis_stream_length, redis_stream_groups y redis_stream_pending_entries_count. Configura alertas para: longitud del stream > 10 000, PEL > 500, lag del grupo > 60 segundos, used_memory > 80% del maxmemory.

Limitaciones de Redis como broker: cuándo sí es necesario Kafka

La honestidad es importante: Redis como broker de eventos tiene limitaciones fundamentales que conviene conocer de antemano.

  • El almacenamiento está limitado por la RAM. Kafka guarda datos en disco y puede retener mensajes durante semanas. Redis, ante la falta de memoria, comienza a expulsar datos — y los mensajes se pierden si no se ha configurado MAXLEN ni se controla la memoria.
  • No hay particionamiento nativo. Un stream es una sola secuencia. Para escalar hay que crear manualmente varios streams y distribuir la carga entre ellos.
  • La replicación no es síncrona. En Redis Sentinel y Cluster, la replicación es asíncrona. En caso de failover puede producirse la pérdida de los últimos mensajes.
  • No hay semántica exactly-once. Redis Streams garantiza at least once. Para exactly-once se requiere idempotencia a nivel de aplicación.
  • No hay Schema Registry integrado. La compatibilidad de formatos de mensajes es responsabilidad del desarrollador.

Elige Kafka si: necesitas almacenar eventos durante meses, el volumen supera varios millones de mensajes por segundo, se requiere una semántica strictly exactly-once, o dispones de un equipo de infraestructura dedicado para su mantenimiento.

Considera RabbitMQ si necesitas un enrutamiento de mensajes más rico (exchanges, routing keys, dead letter queues de serie) con volúmenes moderados.

Caso real: migración de Redis Pub/Sub a Redis Streams en producción

Contexto

El equipo de una plataforma de e-commerce utilizaba Redis Pub/Sub para notificaciones sobre cambios de estado de pedidos. El microservicio de notificaciones se suscribía al canal order.status.changed y enviaba notificaciones push y emails. El sistema funcionaba, pero surgían incidentes periódicos: al reiniciarse el Pod del consumidor en Kubernetes, parte de las notificaciones se perdía — los usuarios no recibían los correos de confirmación de entrega.

El problema

Pub/Sub no almacena mensajes en buffer. Mientras el Pod se reiniciaba (30–60 segundos), todos los eventos se perdían. El rolling update con cero downtime no ayudaba: en el momento del cambio de Pods siempre había una breve interrupción de la suscripción.

La solución

El equipo migró a Redis Streams con el consumer group notification-service. El productor (servicio de gestión de pedidos) comenzó a escribir en el stream order-events en lugar de publicar en el canal. El consumidor lee mediante XREADGROUP y confirma con XACK solo después de enviar la notificación con éxito. Un worker de PEL verifica cada minuto los mensajes bloqueados y reintenta su procesamiento.

El resultado

Durante los tres meses posteriores a la migración, el equipo no registró ningún incidente de pérdida de notificaciones. El tiempo de migración fue de dos días, incluyendo pruebas. No surgieron dependencias de infraestructura adicionales — Redis ya estaba en el stack. Los costos operativos de mantenimiento se mantuvieron igual.

Conclusión y recomendaciones para la elección

Redis como broker de eventos es una elección pragmática para equipos que quieren adoptar una arquitectura event-driven sin la complejidad operativa de Kafka.

  • Usa Redis Pub/Sub para notificaciones broadcast, invalidación de caché y eventos en tiempo real donde la pérdida de uno o dos mensajes no sea crítica.
  • Usa Redis Streams para el procesamiento asíncrono fiable de tareas, intercambio de eventos entre servicios y cualquier escenario donde sea importante la garantía de entrega at least once.
  • Añade KEDA para el autoescalado de consumidores en Kubernetes según la longitud del stream o la PEL.
  • Configura siempre MAXLEN en los streams y monitorea used_memory — esto protegerá contra la pérdida de datos por falta de RAM.
  • Implementa un dead letter stream y la lógica de XCLAIM para gestionar correctamente los fallos.
  • Si tu carga ha crecido a millones de eventos por segundo, se requiere almacenamiento a largo plazo o semántica exactly-once — migra a Kafka. Redis mismo te hará saber cuándo llegue ese momento.

Redis Streams no es una «versión pobre de Kafka». Es una herramienta diferente con trade-offs distintos, perfectamente adecuada para la mayoría de los sistemas de microservicios con carga moderada.

Tecnologías

Etiquetas

Ruslan Ismailov

Desarrollador Senior Web / Backend. Desarrollador senior web/backend con 9 años de experiencia. Stack: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, microservicios, CI/CD. Más sobre mí →