Redis Streams como reemplazo de colas de mensajes: guía práctica con ejemplos
Introducción: por qué Redis Streams no es simplemente Pub/Sub
La mayoría de los desarrolladores conocen Redis como una caché de alto rendimiento. Muchos utilizan Redis Pub/Sub para notificaciones simples. Pero Redis Streams es una estructura de datos independiente y significativamente más potente, introducida en Redis 5.0, que cambia radicalmente el enfoque para organizar colas de mensajes en arquitecturas de microservicios.
¿En qué se diferencia fundamentalmente Redis Streams de Pub/Sub? En Pub/Sub los mensajes no se persisten: si el suscriptor no está disponible en el momento de la publicación, perderá el mensaje para siempre. Redis Streams, por el contrario, almacena todos los mensajes en un registro ordenado con IDs únicos, admite consumer groups con confirmación de entrega (ACK) y permite reproducir el historial de mensajes, de forma similar a Apache Kafka, pero sin su complejidad operacional.
Este artículo está dirigido a desarrolladores backend que ya usan Redis como caché y a arquitectos de microservicios que buscan una alternativa pragmática a RabbitMQ o Kafka para tareas de complejidad media. Cubriremos todo: desde los conceptos básicos hasta código listo para producción en Go e integración en Laravel.
Conceptos fundamentales de Redis Streams
Stream — registro de mensajes
Un Stream en Redis es un append-only log: una estructura de datos a la que solo se pueden agregar nuevas entradas. Cada entrada tiene un ID único con el formato millisecondsTime-sequenceNumber (por ejemplo, 1700000000000-0) y un conjunto arbitrario de campos clave-valor.
# Agregar un mensaje al stream\nXADD orders * user_id 42 product_id 101 action purchase\n\n# Leer los últimos 10 mensajes\nXRANGE orders - + COUNT 10\n\n# Obtener la longitud del stream\nXLEN ordersEl símbolo * en el comando XADD significa «generar el ID automáticamente». También puede pasar un ID explícito, lo cual es útil al replicar datos desde sistemas externos.
Consumer Groups — procesamiento en paralelo
Un Consumer Group es un mecanismo que permite a varios consumidores procesar conjuntamente un stream, donde cada mensaje se entrega a un único consumidor del grupo. Esta es la diferencia clave respecto al XREAD estándar, donde cada lector recibe todos los mensajes.
# Crear un consumer group\nXGROUP CREATE orders processing-group $ MKSTREAM\n\n# Leer nuevos mensajes como consumidor worker-1\nXREADGROUP GROUP processing-group worker-1 COUNT 5 BLOCK 2000 STREAMS orders >\n\n# Confirmar el procesamiento de mensajes (ACK)\nXACK orders processing-group 1700000000000-0 1700000000001-0El símbolo > significa «dame solo los mensajes no procesados». El símbolo $ al crear el grupo significa «comenzar desde los mensajes nuevos»; 0 significa desde el principio del stream.
Pending Entries List (PEL)
Cuando un consumidor toma un mensaje mediante XREADGROUP, este pasa a la Pending Entries List: una lista de mensajes entregados pero aún no confirmados. Esto es fundamental para garantizar la entrega: si un worker falla antes de enviar el ACK, el mensaje permanecerá en el PEL y podrá ser transferido a otro worker.
# Ver mensajes pendientes\nXPENDING orders processing-group - + 10\n\n# Información detallada sobre entradas pendientes específicas\nXPENDING orders processing-group - + 10 worker-1Comparación de Redis Streams con RabbitMQ y Kafka
Antes de escribir código, es importante entender cuándo Redis Streams es la elección correcta y cuándo conviene usar un broker especializado.
| Criterio | Redis Streams | RabbitMQ | Apache Kafka |
|---|---|---|---|
| Persistencia | RDB/AOF, opcional | Disco (colas durables) | Disco (partition log) |
| Throughput | Alto (cientos de miles/s) | Medio | Muy alto (millones/s) |
| Latencia | Muy baja (<1 ms) | Baja | Media (batch) |
| Replay de mensajes | Sí (por offset) | No | Sí |
| Routing | Simple (por nombre de stream) | Complejo (exchanges, bindings) | Por tópicos/particiones |
| Complejidad operacional | Baja | Media | Alta |
| Garantías de entrega | At-least-once | At-least-once / exactly-once | At-least-once / exactly-once |
Use Redis Streams cuando:
- Ya tiene Redis en la infraestructura y quiere evitar agregar un nuevo componente
- La carga es de miles, no millones de mensajes por segundo
- Necesita una arquitectura de eventos simple para microservicios
- La baja latencia es importante
- No necesita enrutamiento complejo
Permanezca en Kafka si:
- Los volúmenes de datos son de cientos de gigabytes por día
- Necesita replicación multinivel y garantías exactly-once
- Requiere almacenamiento de eventos a largo plazo (event sourcing a escala)
Permanezca en RabbitMQ si:
- Necesita enrutamiento complejo a través de exchanges y bindings
- El soporte del protocolo AMQP es importante
- El equipo conoce bien RabbitMQ y no tiene sentido migrar
Práctica: implementación de producer/consumer en Go
Veamos un ejemplo real: un sistema de procesamiento de pedidos en una tienda en línea. El producer publica eventos de creación de pedidos, varios workers consumidores los procesan en paralelo.
Instalación de dependencias
go mod init orders-processor\ngo get github.com/redis/go-redis/v9Producer: publicación de eventos
package main\n\nimport (\n \"context\"\n \"fmt\"\n \"log\"\n \"time\"\n\n \"github.com/redis/go-redis/v9\"\n)\n\nfunc main() {\n rdb := redis.NewClient(&redis.Options{\n Addr: \"localhost:6379\",\n })\n ctx := context.Background()\n\n // Crear el consumer group al iniciar (ignorar error si ya existe)\n err := rdb.XGroupCreateMkStream(ctx, \"orders\", \"processing-group\", \"$\").Err()\n if err != nil && err.Error() != \"BUSYGROUP Consumer Group name already exists\" {\n log.Fatalf(\"Failed to create consumer group: %v\", err)\n }\n\n // Publicar eventos de pedidos\n for i := 1; i <= 100; i++ {\n msgID, err := rdb.XAdd(ctx, &redis.XAddArgs{\n Stream: \"orders\",\n Values: map[string]interface{}{\n \"order_id\": fmt.Sprintf(\"ORD-%04d\", i),\n \"user_id\": i * 10,\n \"amount\": float64(i) * 99.9,\n \"created_at\": time.Now().Unix(),\n },\n }).Result()\n if err != nil {\n log.Printf(\"Failed to publish order: %v\", err)\n continue\n }\n fmt.Printf(\"Published order ORD-%04d with ID: %s\\n\", i, msgID)\n time.Sleep(10 * time.Millisecond)\n }\n}Consumer: procesamiento de mensajes con ACK
package main\n\nimport (\n \"context\"\n \"fmt\"\n \"log\"\n \"os\"\n \"os/signal\"\n \"syscall\"\n \"time\"\n\n \"github.com/redis/go-redis/v9\"\n)\n\nconst (\n streamName = \"orders\"\n groupName = \"processing-group\"\n blockDuration = 2 * time.Second\n maxRetries = 3\n)\n\nfunc main() {\n consumerName := os.Getenv(\"CONSUMER_NAME\")\n if consumerName == \"\" {\n consumerName = \"worker-1\"\n }\n\n rdb := redis.NewClient(&redis.Options{\n Addr: \"localhost:6379\",\n })\n ctx, cancel := context.WithCancel(context.Background())\n\n // Graceful shutdown\n sigCh := make(chan os.Signal, 1)\n signal.Notify(sigCh, syscall.SIGTERM, syscall.SIGINT)\n go func() {\n <-sigCh\n fmt.Println(\"Shutting down...\")\n cancel()\n }()\n\n // Primero procesar mensajes pendientes (tras reinicio)\n processPendingMessages(ctx, rdb, consumerName)\n\n // Bucle principal de procesamiento de nuevos mensajes\n for {\n select {\n case <-ctx.Done():\n return\n default:\n messages, err := rdb.XReadGroup(ctx, &redis.XReadGroupArgs{\n Group: groupName,\n Consumer: consumerName,\n Streams: []string{streamName, \">\"},\n Count: 10,\n Block: blockDuration,\n }).Result()\n if err != nil {\n if err == redis.Nil {\n continue // timeout, no hay nuevos mensajes\n }\n log.Printf(\"Error reading from stream: %v\", err)\n time.Sleep(time.Second)\n continue\n }\n\n for _, stream := range messages {\n for _, msg := range stream.Messages {\n if err := processOrder(msg); err != nil {\n log.Printf(\"Failed to process message %s: %v\", msg.ID, err)\n continue\n }\n // Confirmar procesamiento exitoso\n if err := rdb.XAck(ctx, streamName, groupName, msg.ID).Err(); err != nil {\n log.Printf(\"Failed to ACK message %s: %v\", msg.ID, err)\n }\n }\n }\n }\n }\n}\n\nfunc processOrder(msg redis.XMessage) error {\n orderID := msg.Values[\"order_id\"]\n amount := msg.Values[\"amount\"]\n fmt.Printf(\"Processing order %s, amount: %s\\n\", orderID, amount)\n // Aquí va su lógica de negocio: guardar en BD, enviar email, etc.\n time.Sleep(50 * time.Millisecond) // emulación de trabajo\n return nil\n}\n\nfunc processPendingMessages(ctx context.Context, rdb *redis.Client, consumerName string) {\n pending, err := rdb.XPendingExt(ctx, &redis.XPendingExtArgs{\n Stream: streamName,\n Group: groupName,\n Start: \"-\",\n Stop: \"+\",\n Count: 100,\n Consumer: consumerName,\n }).Result()\n if err != nil {\n return\n }\n for _, p := range pending {\n if p.RetryCount >= maxRetries {\n log.Printf(\"Message %s exceeded retry limit, moving to DLQ\", p.ID)\n // Lógica de Dead Letter Queue\n continue\n }\n // Reclamar y procesar\n msgs, err := rdb.XClaim(ctx, &redis.XClaimArgs{\n Stream: streamName,\n Group: groupName,\n Consumer: consumerName,\n MinIdle: 30 * time.Second,\n Messages: []string{p.ID},\n }).Result()\n if err != nil {\n continue\n }\n for _, msg := range msgs {\n if err := processOrder(msg); err == nil {\n rdb.XAck(ctx, streamName, groupName, msg.ID)\n }\n }\n }\n}Manejo de errores y garantías de entrega
ACK y garantía at-least-once
Redis Streams ofrece garantía de entrega at-least-once: un mensaje se considera procesado únicamente tras un XACK explícito. Si el worker falla antes del ACK, el mensaje permanece en el PEL y puede ser reclamado (XCLAIM) por otro worker.
XPENDING: inspección de mensajes no procesados
# Resumen de mensajes pendientes\nXPENDING orders processing-group - + 100\n\n# Resultado:\n# 1) \"1700000001234-0\"\n# 2) \"worker-1\"\n# 3) (integer) 85000 <-- tiempo desde la última entrega en ms\n# 4) (integer) 2 <-- número de intentos de entregaXCLAIM: recuperación de mensajes bloqueados
XCLAIM permite que un worker tome un mensaje que otro worker lleva demasiado tiempo reteniendo. Esto es fundamental para garantizar la fiabilidad ante fallos de workers.
# Reclamar mensajes con idle superior a 60 segundos\nXCLAIM orders processing-group worker-2 60000 1700000001234-0\n\n# Reclamar automáticamente un lote de mensajes bloqueados (Redis 6.2+)\nXAUTOCLAIM orders processing-group worker-2 60000 0-0 COUNT 10Dead Letter Queue (DLQ)
Para los mensajes que no se pudieron procesar tras N intentos, se recomienda implementar una Dead Letter Queue: un stream separado para mensajes «envenenados»:
# Mover un mensaje problemático a la DLQ\nXADD orders-dlq * original_id 1700000001234-0 order_id ORD-0042 reason \"processing_failed\" attempts 3\nXACK orders processing-group 1700000001234-0\nXDEL orders 1700000001234-0Limitación del tamaño del stream
Redis almacena todos los mensajes en memoria (cuando no se usa descarga a disco). Para evitar desbordamientos, use el parámetro MAXLEN:
# Limitar el stream a 100 000 mensajes (aproximado)\nXADD orders MAXLEN ~ 100000 * order_id ORD-0001 ...\n\n# O recortar periódicamente de forma manual\nXTRIM orders MAXLEN ~ 100000Integración de Redis Streams en una aplicación Laravel
Laravel tiene un driver Redis integrado para colas, pero funciona a través de RPUSH/LPOP (List), no de Streams. A continuación se muestra cómo crear un driver Streams personalizado para Laravel Queue.
Creación del driver
<?php\n// app/Queue/RedisStreamConnector.php\nnamespace App\\Queue;\n\nuse Illuminate\\Queue\\Connectors\\ConnectorInterface;\n\nclass RedisStreamConnector implements ConnectorInterface\n{\n public function connect(array $config): RedisStreamQueue\n {\n return new RedisStreamQueue(\n app('redis')->connection($config['connection'] ?? 'default'),\n $config['queue'] ?? 'default',\n $config['group'] ?? 'laravel-workers',\n $config['consumer'] ?? gethostname(),\n );\n }\n}<?php\n// app/Queue/RedisStreamQueue.php\nnamespace App\\Queue;\n\nuse Illuminate\\Contracts\\Queue\\Queue;\nuse Illuminate\\Queue\\Queue as BaseQueue;\nuse Illuminate\\Redis\\Connections\\Connection;\n\nclass RedisStreamQueue extends BaseQueue implements Queue\n{\n public function __construct(\n protected Connection $redis,\n protected string $default,\n protected string $group,\n protected string $consumer,\n ) {}\n\n public function push($job, $data = '', $queue = null): mixed\n {\n $queue = $this->getQueue($queue);\n $payload = $this->createPayload($job, $queue, $data);\n\n return $this->redis->command('xadd', [\n $queue, '*',\n 'payload', $payload,\n 'attempts', 0,\n ]);\n }\n\n public function pop($queue = null): ?\\Illuminate\\Contracts\\Queue\\Job\n {\n $queue = $this->getQueue($queue);\n $this->ensureGroupExists($queue);\n\n $results = $this->redis->command('xreadgroup', [\n 'GROUP', $this->group, $this->consumer,\n 'COUNT', 1,\n 'BLOCK', 2000,\n 'STREAMS', $queue, '>',\n ]);\n\n if (empty($results[$queue])) {\n return null;\n }\n\n [$messageId, $values] = $results[$queue][0];\n $payload = json_decode($values['payload'], true);\n\n return new RedisStreamJob(\n $this->container, $this, $this->redis,\n $queue, $this->group, $messageId, $payload,\n );\n }\n\n public function ack(string $queue, string $messageId): void\n {\n $this->redis->command('xack', [$queue, $this->group, $messageId]);\n }\n\n protected function ensureGroupExists(string $queue): void\n {\n try {\n $this->redis->command('xgroup', ['CREATE', $queue, $this->group, '$', 'MKSTREAM']);\n } catch (\\Exception $e) {\n // El grupo ya existe — esto es normal\n }\n }\n\n protected function getQueue(?string $queue): string\n {\n return 'stream:' . ($queue ?? $this->default);\n }\n\n // Demás métodos obligatorios de la interfaz Queue...\n public function size($queue = null): int { return 0; }\n public function later($delay, $job, $data = '', $queue = null): mixed { return null; }\n public function bulk($jobs, $data = '', $queue = null): void {}\n}Registro del driver en AppServiceProvider
<?php\n// app/Providers/AppServiceProvider.php\npublic function boot(): void\n{\n Queue::extend('redis-stream', function () {\n return new \\App\\Queue\\RedisStreamConnector();\n });\n}Configuración en config/queue.php
'connections' => [\n 'redis-stream' => [\n 'driver' => 'redis-stream',\n 'connection' => 'default',\n 'queue' => 'default',\n 'group' => 'laravel-workers',\n 'consumer' => env('QUEUE_CONSUMER_NAME', gethostname()),\n 'retry_after' => 90,\n ],\n],Después de esto, todos los jobs estándar de Laravel (dispatch, Queue::push) funcionarán a través de Redis Streams, y obtendrá todas las ventajas: ACK, PEL, replay y consumer groups.
Escalado: múltiples Consumer Groups y particionamiento
Múltiples Consumer Groups para distintos propósitos
Una de las capacidades más potentes de Redis Streams: un stream puede ser leído por varios consumer groups independientes. Esto permite implementar el patrón fan-out sin duplicar datos:
# Los mismos eventos de pedidos son leídos por tres servicios distintos\nXGROUP CREATE orders notification-service $ MKSTREAM\nXGROUP CREATE orders analytics-service $ MKSTREAM \nXGROUP CREATE orders inventory-service $ MKSTREAMCada grupo recibe todos los mensajes y los procesa de forma independiente. Esto es fundamentalmente diferente del round-robin dentro de un mismo grupo, donde cada mensaje lo recibe un único worker.
Particionamiento mediante múltiples streams
Redis es single-threaded por defecto, por lo que para escalar horizontalmente se utilizan varios streams con nombres distintos (similar a las particiones en Kafka):
# El producer determina la partición por hash del user_id\nfunc getPartition(userID int, numPartitions int) string {\n return fmt.Sprintf(\"orders:partition:%d\", userID % numPartitions)\n}\n\n// Publicar en la partición\npartition := getPartition(userID, 8) // 8 particiones\nrdb.XAdd(ctx, &redis.XAddArgs{\n Stream: partition,\n Values: orderData,\n})Al usar Redis Cluster, las particiones se distribuyen automáticamente entre diferentes nodos, lo que proporciona un escalado lineal del throughput.
Escalado automático de workers
En Kubernetes se puede configurar un HPA (Horizontal Pod Autoscaler) basado en la longitud del PEL o el lag del consumer group. Las métricas se exportan mediante redis_exporter a Prometheus:
# Ejemplo de HPA basado en custom metrics\napiVersion: autoscaling/v2\nkind: HorizontalPodAutoscaler\nmetadata:\n name: order-worker-hpa\nspec:\n scaleTargetRef:\n apiVersion: apps/v1\n kind: Deployment\n name: order-worker\n minReplicas: 2\n maxReplicas: 20\n metrics:\n - type: External\n external:\n metric:\n name: redis_stream_pending_entries\n selector:\n matchLabels:\n stream: orders\n target:\n type: AverageValue\n averageValue: \"100\"Monitoreo de Redis Streams: XINFO, métricas y alertas
XINFO: diagnóstico integrado
# Información general del stream\nXINFO STREAM orders\n\n# Información sobre los consumer groups\nXINFO GROUPS orders\n\n# Información sobre los consumidores del grupo\nXINFO CONSUMERS orders processing-group\n\n# Información completa (Redis 7.0+)\nXINFO STREAM orders FULL COUNT 10Métricas clave para el monitoreo:
- pending-messages — cantidad de mensajes en el PEL (debe estar cerca de 0 en condiciones normales)
- lag — diferencia entre el último ID del stream y el último ID leído por el grupo
- consumers count — número de workers activos
- idle time — tiempo de inactividad del consumidor
Monitoreo con redis_exporter + Prometheus
Use redis_exporter para exportar métricas a Prometheus. Agregue el monitoreo de métricas de stream en la configuración:
# prometheus.yml\nscrape_configs:\n - job_name: 'redis'\n static_configs:\n - targets: ['redis-exporter:9121']\n params:\n stream-groups:\n - orders\n - paymentsAlertas en Alertmanager
# Alerta por lag elevado\n- alert: RedisStreamHighLag\n expr: redis_stream_group_lag{stream=\"orders\"} > 1000\n for: 5m\n labels:\n severity: warning\n annotations:\n summary: \"Redis Stream lag is too high\"\n description: \"Consumer group lag: {{ $value }} messages\"\n\n# Alerta por mensajes pendientes bloqueados\n- alert: RedisStreamStalePending \n expr: redis_stream_group_pending{stream=\"orders\"} > 100\n for: 10m\n labels:\n severity: criticalDashboard en Grafana
Paneles recomendados para el dashboard de Redis Streams:
- Longitud del stream (XLEN) en el tiempo
- Lag del consumer group por cada grupo
- Recuento de entradas pendientes
- Tasa de ACK (mensajes/seg)
- Antigüedad del mensaje pendiente más antiguo
- Número de consumidores por grupo
Conclusión
Redis Streams en 2026 es una herramienta madura, de alto rendimiento y operacionalmente sencilla para organizar colas de mensajes en arquitecturas de microservicios. Si ya tiene Redis en su stack, adoptar Streams no requiere incorporar un nuevo componente de infraestructura, y la curva de aprendizaje es significativamente menor que la de Kafka.
Conclusiones principales del artículo:
- Redis Streams se diferencia de Pub/Sub por la persistencia y las garantías de entrega mediante ACK/PEL
- Los Consumer Groups permiten escalar horizontalmente el procesamiento de mensajes
- XCLAIM y XAUTOCLAIM resuelven el problema de mensajes bloqueados ante fallos de workers
- La integración en Go lleva unas pocas horas; en Laravel se puede implementar un driver de colas personalizado
- Monitorear el lag y las entradas pendientes es clave para una operación fiable en producción
Redis Streams no es un reemplazo de Kafka para sistemas a escala de petabytes, pero es una excelente alternativa a RabbitMQ para la mayoría de las tareas en microservicios. Comience con un stream, un consumer group y dos workers, y se sorprenderá de lo lejos que puede llegar.
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í →