HardТеория3 min

Управление данными в микросервисах

CQRS, Event Sourcing, Outbox pattern, database per service, Saga и стратегии согласованности данных

Проблема данных в распределённых системах

В монолите одна БД и ACID транзакции решают всё. В микросервисах каждый сервис владеет своими данными, и нет единой транзакции. Это фундаментальный сдвиг.

Монолит: ACID транзакция
┌─────────────────────────────────────┐
│  BEGIN TRANSACTION                  │
│    INSERT INTO orders ...           │
│    UPDATE inventory SET qty = ...   │
│    INSERT INTO payments ...         │
│  COMMIT                             │
└─────────────────────────────────────┘
  Всё или ничего. Просто и надёжно.

Микросервисы: НЕТ единой транзакции
┌──────────┐  ┌──────────┐  ┌──────────┐
│  Orders  │  │Inventory │  │ Payments │
│  Service │  │ Service  │  │ Service  │
│  ┌────┐  │  │  ┌────┐  │  │  ┌────┐  │
│  │ DB │  │  │  │ DB │  │  │  │ DB │  │
│  └────┘  │  │  └────┘  │  │  └────┘  │
└──────────┘  └──────────┘  └──────────┘
  Каждый сервис -- своя транзакция.
  Как обеспечить согласованность?

Database per Service

Зачем отдельная БД

Свойство Shared DB Database per Service
Независимый деплой Нет (миграции связаны) Да
Изоляция нагрузки Нет (один нагружает -- страдают все) Да
Технологическая свобода Одна СУБД Polyglot persistence
Масштабирование Общее Независимое per service
Согласованность ACID Eventual Consistency
Сложность Низкая Высокая

Стратегия перехода от Shared DB

Этап 1: Shared DB (текущее состояние)
┌───────┐  ┌───────┐  ┌───────┐
│Svc A  │  │Svc B  │  │Svc C  │
└───┬───┘  └───┬───┘  └───┬───┘
    └──────────┼──────────┘
               │
        ┌──────┴──────┐
        │   One DB    │
        └─────────────┘

Этап 2: Отдельные схемы
┌───────┐  ┌───────┐  ┌───────┐
│Svc A  │  │Svc B  │  │Svc C  │
└───┬───┘  └───┬───┘  └───┬───┘
    │          │          │
    v          v          v
┌────────────────────────────┐
│  One DB                    │
│  ┌──────┐┌──────┐┌──────┐ │
│  │Sch A ││Sch B ││Sch C │ │
│  └──────┘└──────┘└──────┘ │
└────────────────────────────┘

Этап 3: Отдельные инстансы
┌───────┐  ┌───────┐  ┌───────┐
│Svc A  │  │Svc B  │  │Svc C  │
└───┬───┘  └───┬───┘  └───┬───┘
    │          │          │
┌───┴───┐  ┌──┴────┐  ┌──┴────┐
│ DB A  │  │ DB B  │  │ DB C  │
└───────┘  └───────┘  └───────┘

CQRS -- Command Query Responsibility Segregation

Идея CQRS

Разделить модели чтения и записи. Write model оптимизирована для бизнес-логики, Read model -- для запросов.

┌─────────────────────────────────────────────────────┐
│                     CQRS                            │
│                                                     │
│  Commands (Write)           Queries (Read)          │
│  ┌──────────────┐           ┌──────────────┐        │
│  │ CreateOrder   │           │ GetOrderList │        │
│  │ CancelOrder   │           │ GetOrderById │        │
│  │ UpdateOrder   │           │ SearchOrders │        │
│  └──────┬───────┘           └──────┬───────┘        │
│         │                          │                │
│         v                          v                │
│  ┌──────────────┐           ┌──────────────┐        │
│  │ Write Model  │  events   │ Read Model   │        │
│  │ (normalized  │──────────>│ (denormal.   │        │
│  │  domain)     │           │  projections)│        │
│  └──────┬───────┘           └──────┬───────┘        │
│         │                          │                │
│         v                          v                │
│  ┌──────────────┐           ┌──────────────┐        │
│  │ PostgreSQL   │           │ Elasticsearch│        │
│  │ (write-opt.) │           │ / Redis      │        │
│  └──────────────┘           │ (read-opt.)  │        │
│                             └──────────────┘        │
└─────────────────────────────────────────────────────┘

Когда использовать CQRS

  • Разная нагрузка -- 90% чтения, 10% записи (или наоборот)
  • Разные модели -- запись нормализована, чтение денормализовано
  • Разные хранилища -- запись в PostgreSQL, чтение из Elasticsearch
  • Масштабирование -- read replicas для чтения, master для записи

Когда НЕ использовать

  • Простые CRUD-приложения
  • Когда read и write модели одинаковы
  • Когда eventual consistency не приемлема

Пример: Order Read Model

Write DB (PostgreSQL):              Read DB (Elasticsearch):
┌────────────────┐                  ┌──────────────────────┐
│ orders         │                  │ order_projections    │
│   id           │   Projection     │   order_id           │
│   user_id      │   ──────────>    │   user_name          │
│   status       │                  │   user_email         │
│   created_at   │                  │   items_count        │
│                │                  │   total_amount       │
│ order_items    │                  │   status             │
│   order_id     │                  │   product_names      │
│   product_id   │                  │   created_at         │
│   quantity     │                  │   shipping_address   │
│   price        │                  └──────────────────────┘
└────────────────┘
  JOIN нужен для получения            Всё в одном документе
  полной информации                   Запрос без JOIN

Event Sourcing

Идея Event Sourcing

Вместо хранения текущего состояния, храним последовательность событий. Текущее состояние = replay всех событий.

Традиционный подход (State Sourcing):
┌──────────────────────┐
│ orders               │
│ id: 123              │
│ status: shipped      │  Знаем текущее состояние
│ total: $50.00        │  НЕ знаем как дошли сюда
│ updated_at: ...      │
└──────────────────────┘

Event Sourcing:
┌──────────────────────────────────────────┐
│ events (append-only log)                 │
│                                          │
│ 1. OrderCreated    {id:123, total:$50}   │
│ 2. PaymentReceived {id:123, amount:$50}  │
│ 3. ItemPacked      {id:123, box:"A"}     │
│ 4. OrderShipped    {id:123, tracking:..} │
│                                          │
│ Текущее состояние = replay событий 1-4   │
│ Полная история = аудит бесплатно         │
└──────────────────────────────────────────┘

Преимущества Event Sourcing

  • Полный аудит -- каждое изменение записано
  • Temporal queries -- состояние на любой момент времени
  • Replay -- можно пересчитать проекции (read models)
  • Debug -- точно видно, что произошло и когда
  • Event-driven -- события уже есть, можно подписаться

Проблемы Event Sourcing

  • Сложность -- кривая обучения, другое мышление
  • Хранилище -- event store может вырасти до огромных размеров
  • Эволюция событий -- как менять формат старых событий (upcasting)
  • Eventual consistency -- проекции обновляются с задержкой

Snapshot для оптимизации

Replay миллиона событий -- долго. Snapshot -- сохранённое состояние после N событий.

Events: 1, 2, 3, ... 999, 1000, 1001, 1002, 1003
                            ^
                     Snapshot @1000
                     (полное состояние)

Восстановление: загрузить snapshot @1000 + replay 1001-1003
Вместо: replay 1-1003 (все 1003 события)

Outbox Pattern

Проблема двойной записи

Как атомарно записать в БД и отправить событие? Две отдельные операции могут упасть по отдельности.

❌ Проблема двойной записи (dual write):

  1. INSERT INTO orders ... -- ✅ Успех
  2. kafka.send(OrderCreated) -- ❌ Ошибка сети!

  Результат: заказ в БД есть, событие не отправлено.
  Другие сервисы не узнают о заказе.

  Или наоборот:
  1. INSERT INTO orders ... -- ❌ Ошибка!
  2. kafka.send(OrderCreated) -- ✅ Успех

  Результат: события есть, заказа нет.

Решение: Transactional Outbox

┌──────────────────────────────────────────────┐
│  Одна транзакция (ACID):                     │
│                                              │
│  BEGIN;                                      │
│    INSERT INTO orders (id, ...)              │
│      VALUES ('ord-123', ...);                │
│                                              │
│    INSERT INTO outbox (id, topic, payload)   │
│      VALUES (gen_id(), 'orders',             │
│        '{"type":"OrderCreated",...}');        │
│  COMMIT;                                     │
│                                              │
│  Outbox таблица -- в ТОЙ ЖЕ базе данных     │
└──────────────────┬───────────────────────────┘
                   │
                   │ Relay process
                   │ (отдельный процесс)
                   v
┌──────────────────────────────────────────────┐
│  Outbox Relay (CDC или Polling):             │
│                                              │
│  1. SELECT FROM outbox WHERE sent = false    │
│  2. kafka.send(event)                        │
│  3. UPDATE outbox SET sent = true            │
│                                              │
│  Гарантия: at-least-once delivery            │
│  Consumer должен быть идемпотентным!         │
└──────────────────────────────────────────────┘

CDC (Change Data Capture) вместо Polling

┌──────────┐     ┌──────────┐     ┌──────────┐
│  Orders  │     │ Debezium │     │  Kafka   │
│  Service │     │  (CDC)   │     │          │
│          │     │          │     │          │
│ orders   │     │ Читает   │     │ orders   │
│ table    │────>│ WAL      │────>│ topic    │
│          │     │ (binlog) │     │          │
│ outbox   │     │          │     │          │
│ table    │────>│          │────>│          │
└──────────┘     └──────────┘     └──────────┘

Debezium подключается к WAL (Write-Ahead Log) PostgreSQL
или binlog MySQL и стримит изменения в Kafka.
Не нужен polling -- изменения приходят в реальном времени.

Saga Pattern -- подробно

Choreography Saga

┌──────────────────────────────────────────────────────┐
│  Choreography Saga: Create Order                     │
│                                                      │
│  Order Service                                       │
│    │                                                 │
│    ├── OrderCreated ──────> Inventory Service         │
│    │                           │                     │
│    │                           ├── InventoryReserved  │
│    │                           │    ──> Payment Svc   │
│    │                           │          │           │
│    │                           │          ├── Payment │
│    │                           │          │   Processed│
│    │                           │          │   ──> Order│
│    │                           │          │      Svc   │
│    │<── OrderConfirmed ────────┘──────────┘           │
│                                                      │
│  При ошибке Payment:                                 │
│    PaymentFailed ──> Inventory: ReleaseInventory      │
│    PaymentFailed ──> Order: CancelOrder               │
└──────────────────────────────────────────────────────┘

Orchestration Saga

┌──────────────────────────────────────────────────────┐
│  Orchestration Saga: Create Order                    │
│                                                      │
│  ┌─────────────────────┐                             │
│  │  Order Saga         │                             │
│  │  Orchestrator       │                             │
│  │                     │                             │
│  │  State Machine:     │                             │
│  │  PENDING            │                             │
│  │    -> RESERVING     │──> Reserve(Inventory)       │
│  │    -> CHARGING      │──> Charge(Payment)          │
│  │    -> CONFIRMED     │──> Confirm(Order)           │
│  │                     │                             │
│  │  Compensations:     │                             │
│  │  CHARGING_FAILED    │                             │
│  │    -> RELEASING     │──> Release(Inventory)       │
│  │    -> CANCELLED     │──> Cancel(Order)            │
│  └─────────────────────┘                             │
│                                                      │
│  Оркестратор хранит состояние саги и управляет       │
│  всеми шагами, включая компенсации.                  │
└──────────────────────────────────────────────────────┘

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

API Composition

Для запросов, требующих данные из нескольких сервисов:

GET /api/v1/order-details/123

┌──────────────────┐
│  API Composer    │
│  (BFF / Gateway) │
└──┬──────┬──────┬─┘
   │      │      │
   │      │      │  Параллельные запросы
   v      v      v
┌─────┐ ┌─────┐ ┌─────┐
│Order│ │User │ │Ship │
│ Svc │ │ Svc │ │ Svc │
└─────┘ └─────┘ └─────┘

Composer агрегирует ответы в один объект:
{
  "order": { ... },
  "user": { "name": "Ivan" },
  "shipping": { "status": "in_transit" }
}

Materialized View / Read Store

Предварительно собранная денормализованная модель:

Order Service ──OrderCreated──┐
User Service ──UserUpdated────┤     ┌──────────────┐
Payment Service ──PaymentOK───┼────>│ Read Store   │
Shipping Service ──Shipped────┘     │ (Projection) │
                                    │              │
                                    │ order_views: │
                                    │  order_id    │
                                    │  user_name   │
                                    │  total       │
                                    │  status      │
                                    │  tracking    │
                                    └──────────────┘

Реальные примеры

Uber -- Event Sourcing для Trip Service

Uber использует Event Sourcing для хранения истории поездок. Каждое изменение состояния поездки -- событие:

  • TripRequested -> DriverMatched -> DriverArrived -> TripStarted -> TripCompleted

Это позволяет:

  • Восстановить состояние поездки при сбое
  • Рассчитать стоимость на основе событий
  • Анализировать паттерны для ML-моделей

Wolt / DoorDash -- Saga для заказа еды

Процесс заказа еды -- классическая Saga:

  1. Create Order
  2. Charge Customer (компенсация: Refund)
  3. Notify Restaurant (компенсация: Cancel Notification)
  4. Assign Courier (компенсация: Unassign)
  5. Track Delivery

Если ресторан отклоняет заказ на шаге 3, происходит компенсация: Refund на шаге 2.

LinkedIn -- CQRS для ленты

LinkedIn строит персонализированную ленту через CQRS:

  • Write: публикация поста (нормализованная модель)
  • Read: лента пользователя (денормализованная, pre-computed для каждого пользователя)
  • Событие PostPublished запускает пересчёт лент подписчиков

Проверь себя

Зачем нужен Snapshot в Event Sourcing?

Какую проблему решает Outbox Pattern?

Что такое CDC (Change Data Capture) в контексте Outbox Pattern?

Что такое Event Sourcing?

В чём суть CQRS паттерна?