Мини-проект: 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
}
type CompletionRepository interface {
Save(ctx context.Context, c *LessonCompletion) error
FindByUser(ctx context.Context, userID uuid.UUID) ([]*LessonCompletion, error)
}
type ProgressRepository interface {
Update(ctx context.Context, userID uuid.UUID, trackID string) error
Get(ctx context.Context, userID uuid.UUID) (*UserProgress, error)
}
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 в одной транзакции - read model всегда консистентен с write 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 {
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 '{}'
);
CREATE INDEX ON completions (user_id);
CREATE INDEX ON user_progress (user_id);
UNIQUE (user_id, lesson_id) в completions - встроенная защита от повторного прохождения.
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 {
http.Error(w, "not found", http.StatusNotFound)
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