Реализация CQRS на Go с PostgreSQL и Redis: разделение команд и запросов в высоконагруженных сервисах
Что такое 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 без нагрузки на основную БД.
Архитектура решения
Рассмотрим сервис управления заказами. Архитектура состоит из трёх слоёв:
- Write side — PostgreSQL. Хранит агрегаты (заказы, позиции) в нормализованной форме. Все команды проходят через доменную валидацию и сохраняются в транзакции.
- Read side — Redis. Хранит денормализованные проекции: готовые JSON-объекты или хеши, которые отдаются клиенту без дополнительных JOIN.
- Синхронизация — доменные события. После успешного коммита транзакции в 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. Подробнее обо мне →