Архитектура

Реализация CQRS на Go с PostgreSQL и Redis: разделение команд и запросов в высоконагруженных сервисах

Ruslan Ismailov Опубликовано 14 мин чтения
Р

Что такое CQRS и зачем он нужен в 2026 году

CQRS (Command Query Responsibility Segregation) — архитектурный паттерн, при котором операции записи (команды) и чтения (запросы) разделяются на уровне модели, хранилища и потока обработки. Его не следует путать с простым разделением слоёв (контроллер / сервис / репозиторий): там по-прежнему одна модель данных обслуживает и чтение, и запись.

В 2026 году нагрузка на публичные API продолжает расти: соотношение read/write в типичных потребительских сервисах достигает 100:1. Единая PostgreSQL-таблица с индексами для OLTP и сложными JOIN-запросами для дашбордов перестаёт справляться. CQRS позволяет масштабировать read side и write side независимо, выбирать оптимальное хранилище для каждой задачи и упрощать доменную модель.

Ключевые преимущества CQRS в микросервисной архитектуре:

  • Write side оптимизирован под транзакционную запись и консистентность.
  • Read side оптимизирован под скорость чтения и денормализацию.
  • Изменение read model не затрагивает доменную логику.
  • Горизонтальное масштабирование read side без нагрузки на основную БД.

Архитектура решения

Рассмотрим сервис управления заказами. Архитектура состоит из трёх слоёв:

  1. Write side — PostgreSQL. Хранит агрегаты (заказы, позиции) в нормализованной форме. Все команды проходят через доменную валидацию и сохраняются в транзакции.
  2. Read side — Redis. Хранит денормализованные проекции: готовые JSON-объекты или хеши, которые отдаются клиенту без дополнительных JOIN.
  3. Синхронизация — доменные события. После успешного коммита транзакции в PostgreSQL Command Handler публикует событие, которое обновляет Redis-проекцию. В более сложных случаях используется Outbox Pattern или Kafka.

Схема потока данных: HTTP Request → Command Handler → PostgreSQL → Domain Event → Projection Updater → Redis → Query Handler → HTTP Response.

Реализация Command Handler на Go

Определим базовые интерфейсы. Каждая команда — это value object без методов, несущий входные данные. Command Handler принимает команду и возвращает ошибку.

// command.go
package cqrs

import "context"

// Command — маркерный интерфейс для всех команд.
type Command interface {
	CommandName() string
}

// CommandHandler обрабатывает конкретную команду.
type CommandHandler[C Command] interface {
	Handle(ctx context.Context, cmd C) error
}

// CommandBus маршрутизирует команды к обработчикам.
type CommandBus interface {
	Dispatch(ctx context.Context, cmd Command) error
}

Реализация команды создания заказа и её обработчика с транзакцией через pgx:

// create_order_command.go
package order

import (
	"context"
	"fmt"

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

// CreateOrderCommand — команда создания заказа.
type CreateOrderCommand struct {
	UserID    string
	Items     []OrderItem
	Currency  string
}

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

// OrderItem — позиция заказа.
type OrderItem struct {
	ProductID string
	Quantity  int
	Price     float64
}

// CreateOrderHandler обрабатывает 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 {
	// Валидация
	if cmd.UserID == "" {
		return fmt.Errorf("userID is required")
	}
	if len(cmd.Items) == 0 {
		return fmt.Errorf("order must contain at least one item")
	}

	// Транзакция в 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)
	}

	// Публикуем доменное событие после успешного коммита
	h.events.Publish(ctx, OrderCreatedEvent{
		OrderID:  orderID,
		UserID:   cmd.UserID,
		Items:    cmd.Items,
		Currency: cmd.Currency,
	})

	return nil
}

Обратите внимание: событие публикуется после успешного коммита транзакции. Это фундаментальный принцип — не допускать ситуации, когда событие опубликовано, а данные в PostgreSQL не записаны.

Построение Read Model в Redis

Определим интерфейс Query Handler и реализуем проекцию заказа в Redis:

// query.go
package cqrs

import "context"

// Query — маркерный интерфейс для всех запросов.
type Query interface {
	QueryName() string
}

// QueryHandler обрабатывает запрос и возвращает результат.
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 — запрос на получение заказа по ID.
type GetOrderQuery struct {
	OrderID string
}

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

// OrderReadModel — денормализованная проекция заказа.
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 читает заказ из 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 обновляет Redis при получении события.
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()
}

В Redis используется строковый ключ order:{id} со значением в формате JSON. Для более сложных проекций (список заказов пользователя) подойдут структуры ZSET (сортировка по дате) или HASH.

Синхронизация и eventual consistency

Главный вопрос при CQRS: что происходит, если Redis и PostgreSQL рассинхронизированы? Это нормальная ситуация для eventually consistent систем, но её нужно обрабатывать явно.

Сценарии рассинхронизации и стратегии их обработки:

  • Падение Redis после коммита в PostgreSQL. Используйте Outbox Pattern: записывайте событие в таблицу outbox внутри той же транзакции, что и основные данные. Отдельный воркер читает таблицу и публикует события.
  • Повторная доставка события. Сделайте Projection Updater идемпотентным: проверяйте версию или timestamp перед перезаписью.
  • Восстановление проекции. Реализуйте replay: читайте события из PostgreSQL (или event log) и пересоздавайте Redis-проекцию с нуля.
  • Staleness на клиенте. Для критичных операций (например, после успешной команды клиент сразу запрашивает данные) используйте стратегию read-your-writes: временно читать из PostgreSQL, а не из Redis.

Eventual consistency — это не баг, а сознательный компромисс между производительностью и строгой согласованностью. Важно документировать этот контракт для команды и для клиентов API.

Деплой в Docker с Docker Compose

Многоконтейнерная среда для локальной разработки и 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 для Go-сервиса с многоэтапной сборкой:

# 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"]

Тестирование CQRS-компонентов

Тестирование разделяется на два уровня: юнит-тесты для доменной логики и интеграционные тесты для взаимодействия с PostgreSQL и Redis.

Юнит-тест Command Handler с моком 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{}
	// В юнит-тесте используем nil для db — валидация происходит до обращения к БД
	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")
}

Интеграционный тест с 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()

	// Запуск 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) })

	// Запуск Redis-контейнера
	redisContainer, err := redis.RunContainer(ctx,
		testcontainers.WithImage("redis:7-alpine"),
	)
	require.NoError(t, err)
	t.Cleanup(func() { redisContainer.Terminate(ctx) })

	// Инициализация зависимостей и выполнение команды
	// ... (подключение, миграции, создание 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)

	// Проверяем наличие проекции в Redis
	// ...
}

Использование testcontainers-go позволяет запускать реальные PostgreSQL и Redis в изолированных контейнерах прямо из Go-тестов, без необходимости поднимать инфраструктуру вручную.

Производительность и подводные камни

Реализуя CQRS в высоконагруженных Go-сервисах, важно учитывать следующее:

  • Избегайте синхронного обновления Redis в Command Handler под блокировкой транзакции. Публикуйте событие в отдельной горутине или через очередь.
  • Пул соединений pgxpool: настройте MaxConns исходя из нагрузки. Дефолтные значения часто занижены для highload-сервисов.
  • Конкурентные записи в Redis: используйте SET NX или Lua-скрипты для атомарных операций при параллельном обновлении одной проекции из нескольких воркеров.
  • TTL для Redis-ключей: устанавливайте разумный TTL, чтобы избежать накопления устаревших данных при редко запрашиваемых объектах.
  • Observability: логируйте latency команд и запросов отдельно. Медленный Command Handler не должен влиять на p99 Query Handler.
  • Не переусложняйте: CQRS оправдан при высокой нагрузке или сложной доменной модели. Для CRUD-сервисов с 100 RPS это избыточная архитектура.

Итог

Паттерн CQRS на Go с PostgreSQL в роли write-хранилища и Redis в роли read-хранилища даёт реальный прирост производительности в высоконагруженных сервисах при правильной реализации. Чёткое разделение интерфейсов CommandHandler и QueryHandler, идемпотентные проекции и надёжная синхронизация через события — основа устойчивой архитектуры. Docker и testcontainers-go обеспечивают воспроизводимую среду для разработки и тестирования. Главное правило: внедряйте CQRS там, где проблема разделения нагрузки действительно существует, а не как архитектурный тренд ради тренда.

Технологии

Теги

Руслан Исмаилов

Senior Web / Backend разработчик. Senior web/backend разработчик с 9-летним опытом. Стек: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, микросервисы, CI/CD. Подробнее обо мне →