Arquitectura

Construcción de un sistema event-driven en Go con Redis Pub/Sub y PostgreSQL

Ruslan Ismailov Publicado 18 min de lectura
C

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í →