Проблема данных в распределённых системах
В монолите одна БД и 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:
- Create Order
- Charge Customer (компенсация: Refund)
- Notify Restaurant (компенсация: Cancel Notification)
- Assign Courier (компенсация: Unassign)
- Track Delivery
Если ресторан отклоняет заказ на шаге 3, происходит компенсация: Refund на шаге 2.
LinkedIn -- CQRS для ленты
LinkedIn строит персонализированную ленту через CQRS:
- Write: публикация поста (нормализованная модель)
- Read: лента пользователя (денормализованная, pre-computed для каждого пользователя)
- Событие PostPublished запускает пересчёт лент подписчиков