Мини-проект: CQRS-система с Command → Read Model
Мини-проект: CQRS-система на Go
Соберём полный CQRS-поток: Command → запись → Read Model. Никаких брокеров - чистое разделение команд и запросов с оптимизированными read-проекциями.
Если нужны брокеры - это трек Message Brokers: там RabbitMQ + Kafka + Outbox паттерн.
Архитектура
HTTP POST /api/lessons/:id/complete
│
▼
CompleteLesson (Command Handler)
│
├─ INSERT completions (write model)
└─ UPDATE user_progress (read model - в одной транзакции)
HTTP GET /api/users/:id/progress
│
▼
GetUserProgress (Query Handler)
│
└─ SELECT user_progress (read model - денормализован, без JOIN)
Ключевое: write model - таблица completions (иммутабельный журнал). Read model - user_progress (денормализованный, быстрый SELECT).
Структура проекта
cmd/
api/main.go
internal/
domain/
lesson.go -- LessonCompletion struct, интерфейсы
usecase/
complete_lesson.go
get_progress.go
repository/
completion_repo.go
progress_repo.go
handler/
lesson_handler.go
config/
wire.go -- DI через Wire
Доменные типы
// internal/domain/lesson.go
type LessonCompletion struct {
ID uuid.UUID
UserID uuid.UUID
LessonID string
TrackID string
CompletedAt time.Time
}
type UserProgress struct {
UserID uuid.UUID
CompletedLessons int
LastActivityAt time.Time
ProgressByTrack map[string]int
}
// ErrProgressNotFound - доменная ошибка «прогресса нет».
//
// Нужна, чтобы handler мог отличить «такого пользователя нет» от «база
// недоступна»: без неё обе ситуации выглядят как `err != nil` и получают
// один код ответа.
var ErrProgressNotFound = errors.New("progress not found")
// Записывающие методы принимают транзакцию: две записи должны попасть
// в одну команду COMMIT, а транзакцией управляет use case.
type CompletionRepository interface {
SaveTx(ctx context.Context, tx *sql.Tx, c *LessonCompletion) error
FindByUser(ctx context.Context, userID uuid.UUID) ([]*LessonCompletion, error)
}
type ProgressRepository interface {
UpdateTx(ctx context.Context, tx *sql.Tx, userID uuid.UUID, trackID string) error
// Get возвращает ErrProgressNotFound, если строки нет.
Get(ctx context.Context, userID uuid.UUID) (*UserProgress, error)
}
Здесь так сделано намеренно - чтобы в одном примере было видно, где именно
начинается и заканчивается транзакция. Способ спрятать её от портов
показан в уроке Репозитории и Unit of Work:
uow.Do(ctx, func(repos) error { ... }), где границу держит адаптер,
а use case о *sql.Tx не знает.
Выбор между двумя вариантами - не вопрос «правильности». Прямая передача проще и честно показывает механику; Unit of Work чище по зависимостям и оправдан, когда репозиториев больше двух или когда СУБД реально может поменяться.
Command: CompleteLesson
// internal/usecase/complete_lesson.go
type CompleteLessonUseCase struct {
completions CompletionRepository
progress ProgressRepository
db *sql.DB
}
type CompleteLessonInput struct {
UserID uuid.UUID
LessonID string
TrackID string
}
func (uc *CompleteLessonUseCase) Execute(ctx context.Context, in CompleteLessonInput) error {
tx, err := uc.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin tx: %w", err)
}
defer tx.Rollback()
completion := &domain.LessonCompletion{
ID: uuid.New(),
UserID: in.UserID,
LessonID: in.LessonID,
TrackID: in.TrackID,
CompletedAt: time.Now(),
}
// Write model: записать факт прохождения
if err := uc.completions.SaveTx(ctx, tx, completion); err != nil {
return fmt.Errorf("save completion: %w", err)
}
// Read model: обновить денормализованную проекцию
if err := uc.progress.UpdateTx(ctx, tx, in.UserID, in.TrackID); err != nil {
return fmt.Errorf("update progress: %w", err)
}
return tx.Commit()
}
INSERT в write model и UPDATE read model идут одной транзакцией, поэтому
проекция не может разъехаться с журналом: либо обе записи применились, либо
ни одна.
Query: GetUserProgress
// internal/usecase/get_progress.go
type GetProgressUseCase struct {
progress ProgressRepository
}
func (uc *GetProgressUseCase) Execute(ctx context.Context, userID uuid.UUID) (*domain.UserProgress, error) {
p, err := uc.progress.Get(ctx, userID)
if err != nil {
// %w сохраняет цепочку, поэтому errors.Is в handler увидит
// domain.ErrProgressNotFound сквозь эту обёртку.
return nil, fmt.Errorf("get progress: %w", err)
}
return p, nil
}
Query не трогает completions - читает только из денормализованного user_progress. Никаких JOIN, никаких агрегаций в реальном времени.
Схема БД
-- Write model: иммутабельный журнал
CREATE TABLE completions (
id UUID PRIMARY KEY,
user_id UUID NOT NULL,
lesson_id TEXT NOT NULL,
track_id TEXT NOT NULL,
completed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
UNIQUE (user_id, lesson_id)
);
-- Read model: денормализованная проекция
CREATE TABLE user_progress (
user_id UUID PRIMARY KEY,
completed_lessons INT NOT NULL DEFAULT 0,
last_activity_at TIMESTAMPTZ,
progress_by_track JSONB NOT NULL DEFAULT '{}'
);
UNIQUE (user_id, lesson_id) в completions - встроенная защита от повторного прохождения.
user_progress.user_id- первичный ключ, а PostgreSQL создаёт под PK уникальный индекс автоматически. Второй индекс на тот же столбец - точная копия, которая только замедляет запись.completions.user_id- ведущий столбец составногоUNIQUE (user_id, lesson_id). Запросы видаWHERE user_id = $1этот индекс используют, как разобрано в уроке Индексы: составной индекс работает по префиксу столбцов.
Единственный сценарий, в котором отдельный узкий индекс оправдан, - когда
запросов по user_id очень много, а составной индекс заметно шире, поэтому
хуже помещается в кэш. Это решается замером EXPLAIN (ANALYZE, BUFFERS)
на реальных данных, а не добавлением индекса на всякий случай: каждый
индекс - это цена на каждом INSERT в completions, а пишем мы туда
в той же транзакции, что и в read model.
HTTP Handler
// internal/handler/lesson_handler.go
func (h *LessonHandler) CompleteLesson(w http.ResponseWriter, r *http.Request) {
userID := middleware.UserIDFromContext(r.Context())
lessonID := chi.URLParam(r, "lessonID")
trackID := chi.URLParam(r, "trackID")
err := h.completeLesson.Execute(r.Context(), usecase.CompleteLessonInput{
UserID: userID,
LessonID: lessonID,
TrackID: trackID,
})
if err != nil {
http.Error(w, "internal error", http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusNoContent)
}
func (h *LessonHandler) GetProgress(w http.ResponseWriter, r *http.Request) {
userID := middleware.UserIDFromContext(r.Context())
progress, err := h.getProgress.Execute(r.Context(), userID)
// Разные причины - разные коды. Одна ветка `if err != nil` с 404
// превратила бы недоступную базу в «нет такого пользователя»:
// клиент перестал бы повторять запрос, а мониторинг увидел бы 404
// вместо ошибки сервера.
switch {
case errors.Is(err, domain.ErrProgressNotFound):
http.Error(w, "not found", http.StatusNotFound)
return
case err != nil:
http.Error(w, "internal error", http.StatusInternalServerError)
return
}
json.NewEncoder(w).Encode(progress)
}
Docker Compose для разработки
services:
postgres:
image: postgres:16-alpine
environment:
POSTGRES_DB: cqrs_demo
POSTGRES_USER: app
POSTGRES_PASSWORD: secret
ports:
- "5432:5432"
healthcheck:
test: pg_isready -U app
interval: 5s
retries: 5
api:
build: .
depends_on:
postgres:
condition: service_healthy
environment:
DATABASE_URL: postgres://app:secret@postgres:5432/cqrs_demo
ports:
- "8080:8080"
Только postgres и api - никаких брокеров. Это и есть чистый CQRS без event-driven.
Тест use case
func TestCompleteLessonUseCase(t *testing.T) {
db := testdb.Open(t) // тестовая БД из testcontainers
repo := repository.NewCompletionRepo(db)
progressRepo := repository.NewProgressRepo(db)
uc := usecase.NewCompleteLessonUseCase(repo, progressRepo, db)
userID := uuid.New()
err := uc.Execute(context.Background(), usecase.CompleteLessonInput{
UserID: userID,
LessonID: "go-01",
TrackID: "go",
})
require.NoError(t, err)
// Проверяем write model
completions, _ := repo.FindByUser(context.Background(), userID)
assert.Len(t, completions, 1)
// Проверяем read model
progress, _ := progressRepo.Get(context.Background(), userID)
assert.Equal(t, 1, progress.CompletedLessons)
}
Типичная ошибка
Обновлять read model асинхронно (через фоновую горутину или брокер), когда в этом нет необходимости. Если read model обновляется в той же транзакции - он всегда консистентен, нет eventual consistency, нет задержек. Брокер нужен, когда событие должны обработать несколько систем - тогда смотри Message Brokers.
Мини-задание
- Реализуй
CompletionRepository.SaveTxиProgressRepository.UpdateTxпринимающие*sql.Tx - Добавь UNIQUE constraint и проверь что повторный вызов возвращает специфичную ошибку (не 500)
- Напиши интеграционный тест с testcontainers: PostgreSQL без моков
- Добавь endpoint
GET /api/users/:id/completionsвозвращающий полный список пройденных уроков из write model