Construcción de un sistema event-driven en Go con Redis Pub/Sub y PostgreSQL
Introducción a la arquitectura event-driven
La arquitectura orientada a eventos (EDA, por sus siglas en inglés) es un enfoque de diseño de sistemas en el que los componentes se comunican mediante eventos en lugar de llamadas directas. En vez de que el servicio A invoque directamente un método del servicio B, publica un evento «algo ha ocurrido», y el servicio B (o varios servicios) reacciona a él de forma independiente.
Las ventajas de este enfoque son evidentes para cualquier desarrollador con experiencia:
- Bajo acoplamiento: el publicador no conoce a los suscriptores ni depende de su disponibilidad.
- Escalabilidad horizontal: los suscriptores escalan de forma independiente a los publicadores.
- Tolerancia a fallos: la indisponibilidad temporal de un componente no bloquea al resto.
- Auditoría y reproducibilidad: los eventos almacenados permiten restaurar el estado del sistema en cualquier momento.
En 2026, la arquitectura event-driven en Go se ha convertido en el estándar de facto para sistemas de microservicios de alto rendimiento. En este artículo construiremos un sistema completo desde cero, utilizando Redis Pub/Sub como bus de eventos y PostgreSQL para la persistencia.
Visión general de las herramientas: Go, Redis Pub/Sub, PostgreSQL
Cada componente de nuestra pila cumple un rol estrictamente definido.
Go — el núcleo del sistema
Go es ideal para sistemas event-driven gracias al soporte nativo de concurrencia mediante goroutines y canales, su biblioteca estándar minimalista y su bajo consumo de memoria. Su naturaleza compilada, la tipificación estricta y el arranque rápido lo convierten en una excelente elección tanto para el publisher como para el subscriber.
Redis Pub/Sub — el bus de eventos
Redis en modo Pub/Sub proporciona un mecanismo «publicar y olvidar» (fire-and-forget) con latencia mínima. No almacena mensajes: si el suscriptor no está disponible en el momento de la publicación, el mensaje se pierde. Esta es la limitación clave que tendremos en cuenta en la arquitectura.
PostgreSQL — almacén de eventos
PostgreSQL actúa como almacén persistente y confiable: registraremos cada evento en una tabla de audit log antes de publicarlo en Redis. Esto nos permite recuperar eventos perdidos y garantizar la entrega.
Diseño de eventos: estructura, esquema y versionado
Un evento es un hecho inmutable que ha ocurrido en el sistema. Un evento bien diseñado contiene suficientes datos para ser procesado sin necesidad de consultas adicionales.
Estructura base de un evento en Go:
package events
import (
"time"
"github.com/google/uuid"
)
// EventType define el tipo de evento
type EventType string
const (
UserCreated EventType = "user.created"
UserUpdated EventType = "user.updated"
OrderPlaced EventType = "order.placed"
PaymentFailed EventType = "payment.failed"
)
// BaseEvent — envoltorio común para todos los eventos
type BaseEvent struct {
ID string `json:"id"` // identificador único del evento
Type EventType `json:"type"` // tipo de evento
Version int `json:"version"` // versión del esquema
OccurredAt time.Time `json:"occurred_at"` // momento en que ocurrió
Source string `json:"source"` // servicio origen
Payload interface{} `json:"payload"` // datos del evento
}
// NewEvent crea un nuevo evento con los metadatos rellenos
func NewEvent(eventType EventType, source string, version int, payload interface{}) BaseEvent {
return BaseEvent{
ID: uuid.New().String(),
Type: eventType,
Version: version,
OccurredAt: time.Now().UTC(),
Source: source,
Payload: payload,
}
}
// UserCreatedPayload — payload para el evento de creación de usuario
type UserCreatedPayload struct {
UserID string `json:"user_id"`
Email string `json:"email"`
Name string `json:"name"`
CreatedAt string `json:"created_at"`
}
El versionado de eventos es crítico para sistemas de larga vida. Utilice el campo version y gestione la migración de esquemas explícitamente en los manejadores. Se recomienda seguir el principio de cambios aditivos: añada nuevos campos sin eliminar los existentes.
Implementación del publisher en Go
El publisher es responsable de publicar eventos en el canal de Redis. El patrón clave consiste en guardar primero el evento en PostgreSQL (outbox pattern) y luego publicarlo en Redis. Esto garantiza la durabilidad.
package publisher
import (
"context"
"encoding/json"
"fmt"
"log"
"github.com/go-redis/redis/v9"
"github.com/jmoiron/sqlx"
"myapp/events"
)
// Publisher publica eventos en Redis y los guarda en PostgreSQL
type Publisher struct {
redis *redis.Client
db *sqlx.DB
}
func NewPublisher(redisClient *redis.Client, db *sqlx.DB) *Publisher {
return &Publisher{
redis: redisClient,
db: db,
}
}
// Publish guarda el evento en la BD y lo publica en Redis
func (p *Publisher) Publish(ctx context.Context, event events.BaseEvent) error {
// 1. Serializamos el evento
payload, err := json.Marshal(event)
if err != nil {
return fmt.Errorf("marshal event: %w", err)
}
// 2. Guardamos en PostgreSQL (outbox / audit log)
if err := p.saveEventToDB(ctx, event, payload); err != nil {
return fmt.Errorf("save event to db: %w", err)
}
// 3. Publicamos en Redis Pub/Sub
channel := string(event.Type)
if err := p.redis.Publish(ctx, channel, payload).Err(); err != nil {
// No es fatal — el evento ya está en la BD, el relay-worker lo reenviará
log.Printf("WARN: failed to publish to Redis channel %s: %v", channel, err)
return nil
}
log.Printf("INFO: published event %s (id=%s) to channel %s", event.Type, event.ID, channel)
return nil
}
// saveEventToDB guarda el evento en la tabla event_outbox
func (p *Publisher) saveEventToDB(ctx context.Context, event events.BaseEvent, payload []byte) error {
query := `
INSERT INTO event_outbox (
id, event_type, version, source, occurred_at, payload, published
) VALUES (
$1, $2, $3, $4, $5, $6, false
)
`
_, err := p.db.ExecContext(ctx, query,
event.ID,
string(event.Type),
event.Version,
event.Source,
event.OccurredAt,
payload,
)
return err
}
Relay Worker: del outbox a Redis
Para garantizar que los eventos que no llegaron a Redis por un fallo sean entregados igualmente, implementamos un worker en segundo plano que consulta el outbox y retransmite los eventos no enviados:
package publisher
import (
"context"
"log"
"time"
)
// RelayWorker reenvía eventos del outbox a Redis
func (p *Publisher) RelayWorker(ctx context.Context) {
ticker := time.NewTicker(5 * time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
if err := p.relayPendingEvents(ctx); err != nil {
log.Printf("ERROR: relay worker: %v", err)
}
}
}
}
func (p *Publisher) relayPendingEvents(ctx context.Context) error {
rows, err := p.db.QueryContext(ctx, `
SELECT id, event_type, payload
FROM event_outbox
WHERE published = false
ORDER BY occurred_at
LIMIT 100
FOR UPDATE SKIP LOCKED
`)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var id, eventType string
var payload []byte
if err := rows.Scan(&id, &eventType, &payload); err != nil {
continue
}
if err := p.redis.Publish(ctx, eventType, payload).Err(); err != nil {
log.Printf("WARN: relay failed for event %s: %v", id, err)
continue
}
// Marcamos como publicado
_, _ = p.db.ExecContext(ctx,
`UPDATE event_outbox SET published = true, published_at = NOW() WHERE id = $1`,
id,
)
}
return rows.Err()
}
Implementación del subscriber y manejadores de eventos
El subscriber se suscribe a los canales de Redis y despacha los eventos a los manejadores registrados. Utilizamos el patrón registry para una registro flexible de manejadores:
package subscriber
import (
"context"
"encoding/json"
"fmt"
"log"
"github.com/go-redis/redis/v9"
"myapp/events"
)
// HandlerFunc — tipo de función manejadora de eventos
type HandlerFunc func(ctx context.Context, event events.BaseEvent) error
// Subscriber gestiona las suscripciones a eventos
type Subscriber struct {
redis *redis.Client
handlers map[events.EventType][]HandlerFunc
}
func NewSubscriber(redisClient *redis.Client) *Subscriber {
return &Subscriber{
redis: redisClient,
handlers: make(map[events.EventType][]HandlerFunc),
}
}
// Register registra un manejador para un tipo de evento específico
func (s *Subscriber) Register(eventType events.EventType, handler HandlerFunc) {
s.handlers[eventType] = append(s.handlers[eventType], handler)
}
// Listen inicia la escucha de los canales de Redis
func (s *Subscriber) Listen(ctx context.Context) error {
channels := make([]string, 0, len(s.handlers))
for eventType := range s.handlers {
channels = append(channels, string(eventType))
}
if len(channels) == 0 {
return fmt.Errorf("no channels registered")
}
pubsub := s.redis.Subscribe(ctx, channels...)
defer pubsub.Close()
log.Printf("INFO: subscribed to channels: %v", channels)
msgCh := pubsub.Channel()
for {
select {
case <-ctx.Done():
log.Println("INFO: subscriber shutting down")
return nil
case msg, ok := <-msgCh:
if !ok {
return fmt.Errorf("subscription channel closed")
}
go s.dispatch(ctx, msg.Channel, []byte(msg.Payload))
}
}
}
// dispatch deserializa el evento e invoca los manejadores
func (s *Subscriber) dispatch(ctx context.Context, channel string, payload []byte) {
var event events.BaseEvent
if err := json.Unmarshal(payload, &event); err != nil {
log.Printf("ERROR: unmarshal event on channel %s: %v", channel, err)
return
}
handlers, ok := s.handlers[event.Type]
if !ok {
log.Printf("WARN: no handlers for event type %s", event.Type)
return
}
for _, handler := range handlers {
if err := handler(ctx, event); err != nil {
log.Printf("ERROR: handler for %s failed: %v", event.Type, err)
}
}
}
Ejemplo de un manejador concreto
package handlers
import (
"context"
"encoding/json"
"log"
"myapp/events"
)
// EmailHandler envía un correo de bienvenida al crear un usuario
type EmailHandler struct {
emailService EmailService
}
func NewEmailHandler(es EmailService) *EmailHandler {
return &EmailHandler{emailService: es}
}
func (h *EmailHandler) Handle(ctx context.Context, event events.BaseEvent) error {
var payload events.UserCreatedPayload
if err := remarshal(event.Payload, &payload); err != nil {
return fmt.Errorf("parse UserCreatedPayload: %w", err)
}
log.Printf("INFO: sending welcome email to %s", payload.Email)
return h.emailService.SendWelcome(ctx, payload.Email, payload.Name)
}
// remarshal convierte el payload interface{} en un tipo concreto
func remarshal(src interface{}, dst interface{}) error {
data, err := json.Marshal(src)
if err != nil {
return err
}
return json.Unmarshal(data, dst)
}
Garantías de entrega: at-least-once e idempotencia
Redis Pub/Sub ofrece la garantía at-most-once: el mensaje se entrega cero o una vez. Para lograr at-least-once utilizamos la combinación del outbox pattern y el relay worker descritos anteriormente. Sin embargo, al reintentar pueden aparecer eventos duplicados, por lo que los manejadores deben ser idempotentes.
Patrón para garantizar la idempotencia — tabla de eventos procesados:
-- Tabla para deduplicación
CREATE TABLE IF NOT EXISTS processed_events (
event_id UUID PRIMARY KEY,
handler VARCHAR(255) NOT NULL,
processed_at TIMESTAMPTZ DEFAULT NOW()
);
-- Índice para búsqueda rápida
CREATE INDEX IF NOT EXISTS idx_processed_events_handler
ON processed_events (handler, event_id);
package handlers
import (
"context"
"database/sql"
"errors"
"fmt"
"github.com/jmoiron/sqlx"
"myapp/events"
)
// IdempotentHandler envuelve un manejador para garantizar idempotencia
type IdempotentHandler struct {
db *sqlx.DB
handlerName string
inner func(ctx context.Context, event events.BaseEvent) error
}
func NewIdempotentHandler(
db *sqlx.DB,
name string,
inner func(ctx context.Context, event events.BaseEvent) error,
) *IdempotentHandler {
return &IdempotentHandler{db: db, handlerName: name, inner: inner}
}
func (h *IdempotentHandler) Handle(ctx context.Context, event events.BaseEvent) error {
// Verificamos si el evento ya fue procesado
var exists bool
err := h.db.QueryRowContext(ctx,
`SELECT EXISTS(SELECT 1 FROM processed_events WHERE event_id = $1 AND handler = $2)`,
event.ID, h.handlerName,
).Scan(&exists)
if err != nil && !errors.Is(err, sql.ErrNoRows) {
return fmt.Errorf("check idempotency: %w", err)
}
if exists {
// El evento ya fue procesado — lo omitimos
return nil
}
// Ejecutamos la lógica de negocio
if err := h.inner(ctx, event); err != nil {
return err
}
// Marcamos como procesado
_, err = h.db.ExecContext(ctx,
`INSERT INTO processed_events (event_id, handler) VALUES ($1, $2) ON CONFLICT DO NOTHING`,
event.ID, h.handlerName,
)
return err
}
Almacenamiento de eventos en PostgreSQL como audit log
Esquema de tablas para el outbox y el audit log:
-- Tabla para el Outbox Pattern
CREATE TABLE IF NOT EXISTS event_outbox (
id UUID PRIMARY KEY,
event_type VARCHAR(255) NOT NULL,
version INTEGER NOT NULL DEFAULT 1,
source VARCHAR(255) NOT NULL,
occurred_at TIMESTAMPTZ NOT NULL,
payload JSONB NOT NULL,
published BOOLEAN NOT NULL DEFAULT false,
published_at TIMESTAMPTZ,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE INDEX idx_event_outbox_unpublished
ON event_outbox (occurred_at)
WHERE published = false;
-- Audit log para análisis y restauración de estado
CREATE TABLE IF NOT EXISTS event_audit_log (
id BIGSERIAL PRIMARY KEY,
event_id UUID NOT NULL,
event_type VARCHAR(255) NOT NULL,
version INTEGER NOT NULL,
source VARCHAR(255) NOT NULL,
occurred_at TIMESTAMPTZ NOT NULL,
payload JSONB NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE INDEX idx_event_audit_log_event_type
ON event_audit_log (event_type, occurred_at DESC);
CREATE INDEX idx_event_audit_log_event_id
ON event_audit_log (event_id);
La función de escritura en el audit log se invoca desde el subscriber tras el procesamiento exitoso del evento:
func SaveToAuditLog(ctx context.Context, db *sqlx.DB, event events.BaseEvent) error {
payload, err := json.Marshal(event.Payload)
if err != nil {
return err
}
_, err = db.ExecContext(ctx, `
INSERT INTO event_audit_log
(event_id, event_type, version, source, occurred_at, payload)
VALUES ($1, $2, $3, $4, $5, $6)
ON CONFLICT DO NOTHING
`, event.ID, event.Type, event.Version, event.Source, event.OccurredAt, payload)
return err
}
Arranque y orquestación con Docker Compose
Docker Compose permite levantar todo el entorno con un solo comando. A continuación se muestra un archivo completo para el entorno de desarrollo:
version: '3.9'
services:
postgres:
image: postgres:16-alpine
container_name: eda_postgres
environment:
POSTGRES_USER: eda_user
POSTGRES_PASSWORD: eda_secret
POSTGRES_DB: eda_db
ports:
- "5432:5432"
volumes:
- postgres_data:/var/lib/postgresql/data
- ./migrations:/docker-entrypoint-initdb.d
healthcheck:
test: ["CMD-SHELL", "pg_isready -U eda_user -d eda_db"]
interval: 5s
timeout: 5s
retries: 5
redis:
image: redis:7-alpine
container_name: eda_redis
ports:
- "6379:6379"
command: redis-server --save 60 1 --loglevel warning
volumes:
- redis_data:/data
healthcheck:
test: ["CMD", "redis-cli", "ping"]
interval: 5s
timeout: 3s
retries: 5
publisher:
build:
context: .
dockerfile: ./cmd/publisher/Dockerfile
container_name: eda_publisher
environment:
DATABASE_URL: postgres://eda_user:eda_secret@postgres:5432/eda_db?sslmode=disable
REDIS_URL: redis:6379
SERVICE_NAME: publisher-service
depends_on:
postgres:
condition: service_healthy
redis:
condition: service_healthy
subscriber:
build:
context: .
dockerfile: ./cmd/subscriber/Dockerfile
container_name: eda_subscriber
environment:
DATABASE_URL: postgres://eda_user:eda_secret@postgres:5432/eda_db?sslmode=disable
REDIS_URL: redis:6379
SERVICE_NAME: subscriber-service
depends_on:
postgres:
condition: service_healthy
redis:
condition: service_healthy
deploy:
replicas: 2
volumes:
postgres_data:
redis_data:
Pruebas del sistema event-driven
Probar una EDA requiere un enfoque especial. Dividimos las pruebas en tres niveles: pruebas unitarias para los manejadores, pruebas de integración con Redis y PostgreSQL reales, y pruebas end-to-end de los flujos de eventos.
Pruebas unitarias de los manejadores
package handlers_test
import (
"context"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/mock"
"myapp/events"
"myapp/handlers"
)
type MockEmailService struct {
mock.Mock
}
func (m *MockEmailService) SendWelcome(ctx context.Context, email, name string) error {
args := m.Called(ctx, email, name)
return args.Error(0)
}
func TestEmailHandler_Handle(t *testing.T) {
mockES := new(MockEmailService)
mockES.On("SendWelcome", mock.Anything, "user@example.com", "Alice").Return(nil)
handler := handlers.NewEmailHandler(mockES)
payload := events.UserCreatedPayload{
UserID: "123",
Email: "user@example.com",
Name: "Alice",
}
event := events.NewEvent(events.UserCreated, "user-service", 1, payload)
err := handler.Handle(context.Background(), event)
assert.NoError(t, err)
mockES.AssertExpectations(t)
}
func TestIdempotentHandler_SkipsDuplicate(t *testing.T) {
// Prueba de integración con testcontainers-go o sqlmock
// Primera llamada — procesa el evento
// Segunda llamada con el mismo event.ID — lo omite
// La implementación detallada depende de la infraestructura de pruebas
t.Log("See integration tests for full idempotency coverage")
}
Pruebas de integración con testcontainers-go
package integration_test
import (
"context"
"testing"
"time"
"github.com/stretchr/testify/require"
"github.com/testcontainers/testcontainers-go"
"github.com/testcontainers/testcontainers-go/modules/redis"
"github.com/testcontainers/testcontainers-go/modules/postgres"
)
func TestPublishSubscribeFlow(t *testing.T) {
ctx := context.Background()
// Levantamos el contenedor de Redis
redisContainer, err := redis.RunContainer(ctx,
testcontainers.WithImage("redis:7-alpine"),
)
require.NoError(t, err)
defer redisContainer.Terminate(ctx)
// Levantamos el contenedor de PostgreSQL
pgContainer, err := postgres.RunContainer(ctx,
testcontainers.WithImage("postgres:16-alpine"),
postgres.WithDatabase("test_db"),
postgres.WithUsername("test"),
postgres.WithPassword("test"),
)
require.NoError(t, err)
defer pgContainer.Terminate(ctx)
// ... inicialización del publisher y subscriber
// ... publicación del evento de prueba
// ... verificación del procesamiento mediante canal o WaitGroup con timeout
received := make(chan events.BaseEvent, 1)
// subscriber.Register(events.UserCreated, func(ctx context.Context, e events.BaseEvent) error {
// received <- e
// return nil
// })
select {
case event := <-received:
require.Equal(t, events.UserCreated, event.Type)
case <-time.After(5 * time.Second):
t.Fatal("timeout waiting for event")
}
}
Limitaciones de Redis Pub/Sub y cuándo considerar Kafka
Redis Pub/Sub es una excelente herramienta para determinados escenarios, pero tiene limitaciones fundamentales que es necesario comprender.
- Sin persistencia de mensajes: si el suscriptor no está conectado en el momento de la publicación, el mensaje se pierde de forma irrecuperable. Por eso utilizamos el outbox pattern.
- Sin consumer groups: todos los suscriptores de un canal reciben una copia de cada mensaje. No es posible distribuir la carga entre varias instancias de un mismo servicio sin lógica adicional.
- Sin replay: no es posible reproducir eventos desde un offset determinado, como en Kafka.
- Throughput limitado: bajo cargas muy elevadas (millones de mensajes por segundo), Redis se convierte en un cuello de botella.
Use Redis Streams (XADD/XREADGROUP) si necesita consumer groups y persistencia básica manteniéndose en el ecosistema Redis. Migre a Apache Kafka cuando requiera entrega garantizada, replay de eventos, procesamiento de millones de mensajes por segundo o almacenamiento a largo plazo del event log.
En una arquitectura de microservicios, Redis Pub/Sub es ideal para: notificaciones internas de baja latencia, invalidación de caché, notificaciones en tiempo real y casos en que la pérdida ocasional de un mensaje no es crítica gracias a la compensación del outbox. Para transacciones financieras, eventos de negocio críticos y sistemas con requisitos de compliance, elija Kafka u otros brokers con entrega garantizada.
Conclusión
Hemos construido un sistema event-driven completo en Go, combinando Redis Pub/Sub como bus de eventos ligero con PostgreSQL para una persistencia confiable mediante el outbox pattern. Los patrones clave que hemos aplicado son:
- Outbox Pattern — guardamos el evento en la BD antes de publicarlo en Redis, garantizando la entrega at-least-once.
- Idempotent Handlers — la tabla de eventos procesados protege contra la duplicación en los reintentos.
- Event Registry — registro flexible de múltiples manejadores para un mismo tipo de evento.
- Relay Worker — proceso en segundo plano que retransmite eventos cuando Redis se recupera.
El procesamiento asíncrono de eventos en Go en 2026 no es algo exótico, sino una necesidad práctica para sistemas de microservicios escalables. Comience con Redis Pub/Sub + PostgreSQL y migre a Kafka solo cuando realmente encuentre las limitaciones del sistema actual.
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í →