Arquitectura

Implementación de CQRS en Go con PostgreSQL y Redis: separación de comandos y consultas en servicios de alta carga

Ruslan Ismailov Publicado 14 min de lectura
I

Qué es CQRS y por qué es relevante en 2026

CQRS (Command Query Responsibility Segregation) es un patrón arquitectónico en el que las operaciones de escritura (comandos) y de lectura (consultas) se separan a nivel de modelo, almacenamiento y flujo de procesamiento. No debe confundirse con la simple separación en capas (controlador / servicio / repositorio): en ese caso sigue existiendo un único modelo de datos que sirve tanto para lectura como para escritura.

En 2026, la carga sobre las API públicas continúa creciendo: la relación read/write en servicios de consumo típicos alcanza 100:1. Una única tabla PostgreSQL con índices para OLTP y complejas consultas JOIN para dashboards deja de ser suficiente. CQRS permite escalar el read side y el write side de forma independiente, elegir el almacenamiento óptimo para cada tarea y simplificar el modelo de dominio.

Ventajas clave de CQRS en arquitecturas de microservicios:

  • El write side está optimizado para escritura transaccional y consistencia.
  • El read side está optimizado para velocidad de lectura y desnormalización.
  • Los cambios en el read model no afectan la lógica de dominio.
  • Escalado horizontal del read side sin carga adicional sobre la base de datos principal.

Arquitectura de la solución

Consideremos un servicio de gestión de pedidos. La arquitectura consta de tres capas:

  1. Write side — PostgreSQL. Almacena agregados (pedidos, líneas de pedido) en forma normalizada. Todos los comandos pasan por validación de dominio y se guardan en una transacción.
  2. Read side — Redis. Almacena proyecciones desnormalizadas: objetos JSON listos para entregar al cliente sin JOINs adicionales.
  3. Sincronización — eventos de dominio. Tras el commit exitoso de la transacción en PostgreSQL, el Command Handler publica un evento que actualiza la proyección en Redis. En casos más complejos se utiliza el Outbox Pattern o Kafka.

Esquema del flujo de datos: HTTP Request → Command Handler → PostgreSQL → Domain Event → Projection Updater → Redis → Query Handler → HTTP Response.

Implementación del Command Handler en Go

Definamos las interfaces base. Cada comando es un value object sin métodos que transporta los datos de entrada. El Command Handler recibe el comando y devuelve un error.

// command.go
package cqrs

import "context"

// Command — interfaz marcadora para todos los comandos.
type Command interface {
	CommandName() string
}

// CommandHandler procesa un comando específico.
type CommandHandler[C Command] interface {
	Handle(ctx context.Context, cmd C) error
}

// CommandBus enruta los comandos hacia sus manejadores.
type CommandBus interface {
	Dispatch(ctx context.Context, cmd Command) error
}

Implementación del comando de creación de pedido y su manejador con transacción mediante pgx:

// create_order_command.go
package order

import (
	"context"
	"fmt"

	"github.com/jackc/pgx/v5"
	"github.com/jackc/pgx/v5/pgxpool"
)

// CreateOrderCommand — comando para crear un pedido.
type CreateOrderCommand struct {
	UserID    string
	Items     []OrderItem
	Currency  string
}

func (c CreateOrderCommand) CommandName() string { return "order.create" }

// OrderItem — línea de pedido.
type OrderItem struct {
	ProductID string
	Quantity  int
	Price     float64
}

// CreateOrderHandler procesa CreateOrderCommand.
type CreateOrderHandler struct {
	db        *pgxpool.Pool
	events    EventPublisher
}

func NewCreateOrderHandler(db *pgxpool.Pool, events EventPublisher) *CreateOrderHandler {
	return &CreateOrderHandler{db: db, events: events}
}

func (h *CreateOrderHandler) Handle(ctx context.Context, cmd CreateOrderCommand) error {
	// Validación
	if cmd.UserID == "" {
		return fmt.Errorf("userID is required")
	}
	if len(cmd.Items) == 0 {
		return fmt.Errorf("order must contain at least one item")
	}

	// Transacción en PostgreSQL
	tx, err := h.db.Begin(ctx)
	if err != nil {
		return fmt.Errorf("begin tx: %w", err)
	}
	defer tx.Rollback(ctx)

	var orderID string
	err = tx.QueryRow(ctx,
		`INSERT INTO orders (user_id, currency, status, created_at)
		 VALUES ($1, $2, 'pending', NOW()) RETURNING id`,
		cmd.UserID, cmd.Currency,
	).Scan(&orderID)
	if err != nil {
		return fmt.Errorf("insert order: %w", err)
	}

	for _, item := range cmd.Items {
		_, err = tx.Exec(ctx,
			`INSERT INTO order_items (order_id, product_id, quantity, price)
			 VALUES ($1, $2, $3, $4)`,
			orderID, item.ProductID, item.Quantity, item.Price,
		)
		if err != nil {
			return fmt.Errorf("insert item: %w", err)
		}
	}

	if err = tx.Commit(ctx); err != nil {
		return fmt.Errorf("commit: %w", err)
	}

	// Publicamos el evento de dominio tras el commit exitoso
	h.events.Publish(ctx, OrderCreatedEvent{
		OrderID:  orderID,
		UserID:   cmd.UserID,
		Items:    cmd.Items,
		Currency: cmd.Currency,
	})

	return nil
}

Nótese que el evento se publica después del commit exitoso de la transacción. Este es un principio fundamental: evitar situaciones en que el evento ya fue publicado pero los datos en PostgreSQL aún no fueron escritos.

Construcción del Read Model en Redis

Definamos la interfaz del Query Handler e implementemos la proyección del pedido en Redis:

// query.go
package cqrs

import "context"

// Query — interfaz marcadora para todas las consultas.
type Query interface {
	QueryName() string
}

// QueryHandler procesa una consulta y devuelve un resultado.
type QueryHandler[Q Query, R any] interface {
	Handle(ctx context.Context, query Q) (R, error)
}
// get_order_query.go
package order

import (
	"context"
	"encoding/json"
	"fmt"

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

// GetOrderQuery — consulta para obtener un pedido por ID.
type GetOrderQuery struct {
	OrderID string
}

func (q GetOrderQuery) QueryName() string { return "order.get" }

// OrderReadModel — proyección desnormalizada del pedido.
type OrderReadModel struct {
	ID       string      `json:"id"`
	UserID   string      `json:"user_id"`
	Currency string      `json:"currency"`
	Status   string      `json:"status"`
	Items    []OrderItem `json:"items"`
}

// GetOrderHandler lee el pedido desde Redis.
type GetOrderHandler struct {
	rdb *redis.Client
}

func NewGetOrderHandler(rdb *redis.Client) *GetOrderHandler {
	return &GetOrderHandler{rdb: rdb}
}

func (h *GetOrderHandler) Handle(ctx context.Context, q GetOrderQuery) (*OrderReadModel, error) {
	key := fmt.Sprintf("order:%s", q.OrderID)
	data, err := h.rdb.Get(ctx, key).Bytes()
	if err == redis.Nil {
		return nil, fmt.Errorf("order %s not found", q.OrderID)
	}
	if err != nil {
		return nil, fmt.Errorf("redis get: %w", err)
	}

	var model OrderReadModel
	if err = json.Unmarshal(data, &model); err != nil {
		return nil, fmt.Errorf("unmarshal: %w", err)
	}
	return &model, nil
}

// ProjectionUpdater actualiza Redis al recibir un evento.
type ProjectionUpdater struct {
	rdb *redis.Client
}

func (u *ProjectionUpdater) OnOrderCreated(ctx context.Context, event OrderCreatedEvent) error {
	model := OrderReadModel{
		ID:       event.OrderID,
		UserID:   event.UserID,
		Currency: event.Currency,
		Status:   "pending",
		Items:    event.Items,
	}
	data, err := json.Marshal(model)
	if err != nil {
		return err
	}
	key := fmt.Sprintf("order:%s", event.OrderID)
	return u.rdb.Set(ctx, key, data, 0).Err()
}

En Redis se utiliza una clave de tipo string order:{id} con valor en formato JSON. Para proyecciones más complejas (lista de pedidos de un usuario), son adecuadas las estructuras ZSET (ordenación por fecha) o HASH.

Sincronización y consistencia eventual

La pregunta principal en CQRS: ¿qué ocurre si Redis y PostgreSQL se dessincronizan? Es una situación normal en sistemas eventually consistent, pero debe gestionarse de forma explícita.

Escenarios de desincronización y estrategias para manejarlos:

  • Caída de Redis tras el commit en PostgreSQL. Utilice el Outbox Pattern: escriba el evento en una tabla outbox dentro de la misma transacción que los datos principales. Un worker independiente lee la tabla y publica los eventos.
  • Entrega duplicada de eventos. Haga que el Projection Updater sea idempotente: verifique la versión o el timestamp antes de sobrescribir.
  • Reconstrucción de la proyección. Implemente replay: lea los eventos desde PostgreSQL (o el event log) y recree la proyección en Redis desde cero.
  • Datos desactualizados en el cliente. Para operaciones críticas (por ejemplo, cuando el cliente consulta datos inmediatamente después de un comando exitoso), use la estrategia read-your-writes: leer temporalmente desde PostgreSQL en lugar de Redis.

La consistencia eventual no es un bug, sino un compromiso consciente entre rendimiento y consistencia estricta. Es importante documentar este contrato tanto para el equipo como para los clientes de la API.

Despliegue en Docker con Docker Compose

Entorno multicontenedor para desarrollo local y CI:

# docker-compose.yml
version: "3.9"

services:
  app:
    build: .
    ports:
      - "8080:8080"
    environment:
      DATABASE_URL: postgres://user:password@postgres:5432/orders?sslmode=disable
      REDIS_URL: redis:6379
    depends_on:
      postgres:
        condition: service_healthy
      redis:
        condition: service_healthy

  postgres:
    image: postgres:16-alpine
    environment:
      POSTGRES_USER: user
      POSTGRES_PASSWORD: password
      POSTGRES_DB: orders
    volumes:
      - pg_data:/var/lib/postgresql/data
      - ./migrations:/docker-entrypoint-initdb.d
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U user -d orders"]
      interval: 5s
      timeout: 5s
      retries: 5

  redis:
    image: redis:7-alpine
    volumes:
      - redis_data:/data
    healthcheck:
      test: ["CMD", "redis-cli", "ping"]
      interval: 5s
      timeout: 3s
      retries: 5

volumes:
  pg_data:
  redis_data:

Dockerfile para el servicio Go con compilación multietapa:

# Dockerfile
FROM golang:1.22-alpine AS builder
WORKDIR /app
COPY go.mod go.sum ./
RUN go mod download
COPY . .
RUN CGO_ENABLED=0 GOOS=linux go build -o /app/server ./cmd/server

FROM alpine:3.19
RUN apk add --no-cache ca-certificates
COPY --from=builder /app/server /server
EXPOSE 8080
ENTRYPOINT ["/server"]

Pruebas de componentes CQRS

Las pruebas se dividen en dos niveles: pruebas unitarias para la lógica de dominio y pruebas de integración para la interacción con PostgreSQL y Redis.

Prueba unitaria del Command Handler con mock de EventPublisher:

// create_order_handler_test.go
package order_test

import (
	"context"
	"testing"

	"github.com/stretchr/testify/assert"
	"github.com/stretchr/testify/mock"
)

type MockEventPublisher struct {
	mock.Mock
}

func (m *MockEventPublisher) Publish(ctx context.Context, event interface{}) {
	m.Called(ctx, event)
}

func TestCreateOrderHandler_EmptyItems_ReturnsError(t *testing.T) {
	publisher := &MockEventPublisher{}
	// En la prueba unitaria usamos nil para db — la validación ocurre antes de acceder a la BD
	handler := NewCreateOrderHandler(nil, publisher)

	err := handler.Handle(context.Background(), CreateOrderCommand{
		UserID:   "user-1",
		Items:    []OrderItem{},
		Currency: "USD",
	})

	assert.EqualError(t, err, "order must contain at least one item")
	publisher.AssertNotCalled(t, "Publish")
}

Prueba de integración con testcontainers-go:

// integration_test.go
package order_test

import (
	"context"
	"testing"

	"github.com/stretchr/testify/require"
	"github.com/testcontainers/testcontainers-go/modules/postgres"
	"github.com/testcontainers/testcontainers-go/modules/redis"
)

func TestCreateOrderHandler_Integration(t *testing.T) {
	ctx := context.Background()

	// Inicio del contenedor PostgreSQL
	pgContainer, err := postgres.RunContainer(ctx,
		testcontainers.WithImage("postgres:16-alpine"),
		postgres.WithDatabase("testdb"),
		postgres.WithUsername("test"),
		postgres.WithPassword("test"),
	)
	require.NoError(t, err)
	t.Cleanup(func() { pgContainer.Terminate(ctx) })

	// Inicio del contenedor Redis
	redisContainer, err := redis.RunContainer(ctx,
		testcontainers.WithImage("redis:7-alpine"),
	)
	require.NoError(t, err)
	t.Cleanup(func() { redisContainer.Terminate(ctx) })

	// Inicialización de dependencias y ejecución del comando
	// ... (conexión, migraciones, creación del handler)

	cmd := CreateOrderCommand{
		UserID:   "user-42",
		Items:    []OrderItem{{ProductID: "prod-1", Quantity: 2, Price: 9.99}},
		Currency: "USD",
	}
	err = handler.Handle(ctx, cmd)
	require.NoError(t, err)

	// Verificamos la existencia de la proyección en Redis
	// ...
}

El uso de testcontainers-go permite ejecutar PostgreSQL y Redis reales en contenedores aislados directamente desde las pruebas de Go, sin necesidad de levantar la infraestructura manualmente.

Rendimiento y puntos de atención

Al implementar CQRS en servicios Go de alta carga, es importante tener en cuenta lo siguiente:

  • Evite actualizar Redis de forma síncrona en el Command Handler bajo bloqueo de transacción. Publique el evento en una goroutine separada o mediante una cola.
  • Pool de conexiones pgxpool: configure MaxConns según la carga. Los valores por defecto suelen ser insuficientes para servicios de alta concurrencia.
  • Escrituras concurrentes en Redis: utilice SET NX o scripts Lua para operaciones atómicas cuando múltiples workers actualicen la misma proyección en paralelo.
  • TTL para las claves de Redis: establezca un TTL razonable para evitar la acumulación de datos obsoletos en objetos poco consultados.
  • Observabilidad: registre la latencia de comandos y consultas por separado. Un Command Handler lento no debe afectar el p99 del Query Handler.
  • No sobrecomplejice: CQRS se justifica ante alta carga o un modelo de dominio complejo. Para servicios CRUD con 100 RPS es una arquitectura innecesariamente compleja.

Conclusión

El patrón CQRS en Go con PostgreSQL como almacén de escritura y Redis como almacén de lectura proporciona una mejora real de rendimiento en servicios de alta carga cuando se implementa correctamente. La separación clara de las interfaces CommandHandler y QueryHandler, las proyecciones idempotentes y la sincronización confiable mediante eventos son la base de una arquitectura sólida. Docker y testcontainers-go garantizan un entorno reproducible para el desarrollo y las pruebas. La regla principal: adopte CQRS donde el problema de separación de carga realmente existe, y no como una tendencia arquitectónica por seguir la moda.

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