Суть проблемы
Сервис хочет записать изменение в две системы одновременно:
- БД + очередь сообщений (Kafka/RabbitMQ)
- БД + поисковый индекс (Elasticsearch)
- БД_A + БД_B (кросс-сервисная)
- БД + cache (Redis)
Задача -- сохранить консистентность: либо обе записи удались, либо ни одна.
В распределённой системе без специальных средств это невозможно сделать атомарно. Каждая комбинация "DB commit + внешний вызов" имеет failure window, где одна система зафиксировала изменение, а вторая -- нет.
Что может пойти не так
Код:
db.insert(order) -- Step 1
db.commit()
kafka.send(event) -- Step 2
Возможные точки отказа:
t0: ─── insert + commit ──────┬── send ────────
│
│
v
FAILURE WINDOW
Сценарий A: краш между commit и send
-> DB имеет order, Kafka не знает -> другие сервисы рассогласованы
Сценарий B: сеть до Kafka timeout
-> DB имеет order, Kafka неизвестно: дошло или нет
-> retry? если дошло -> дубль. Если нет -> потеря
Сценарий C: send успешен, но ответ потерян
-> ACK не вернулся, код считает "не отправилось"
-> retry -> дубль события
Сценарий D: Kafka положил, а наш процесс упал до return
-> клиент увидит ошибку HTTP, повторит -> дубль
Поменяем порядок:
kafka.send(event) -- Step 1
db.insert(order) -- Step 2
db.commit()
Теперь:
Сценарий E: send успешен, а commit упал
-> событие ушло в другие сервисы о несуществующем заказе
-> eventual consistency нарушена в обратную сторону
Сценарий F: send таймаут; мы считаем "не отправилось"
-> rollback DB. Но на самом деле send успел.
-> событие без order есть
Нет порядка, который бы решал проблему. Нужны другие подходы.
Решение 1: Transactional Outbox
Главная идея -- превратить dual-write в single-write, записывая всё в одну БД.
┌─────────────────────────────┐
│ BEGIN TRANSACTION │
│ INSERT INTO orders ... │
│ INSERT INTO outbox ... │
│ COMMIT │ <-- атомарно
└──────────┬──────────────────┘
│
v
┌──────────────┐
│ outbox table │
└──────┬───────┘
│
v
[Worker polls outbox]
│
v
[Publish to Kafka]
│
v
UPDATE outbox SET sent_at = NOW()
Гарантия: at-least-once. Событие будет опубликовано, возможно повторно (если worker упал между publish и UPDATE). Consumers должны быть идемпотентны (inbox pattern).
Подробная реализация -- в 3.outbox.md.
Плюсы
- Атомарность через ACID
- Простая отладка: вся активность в одной БД
- Контроль над payload события
Минусы
- Дополнительная таблица и worker
- Дополнительная latency (polling)
- Бизнес-логика и инфраструктура смешаны в БД
Решение 2: Change Data Capture (CDC)
CDC читает бинарный лог БД (WAL в Postgres, binlog в MySQL) и стримит изменения в Kafka. Не требует кода в приложении -- всё решается инфраструктурой.
┌─────────────────────────────┐
│ BEGIN TRANSACTION │
│ INSERT INTO orders ... │
│ COMMIT │
└──────────┬──────────────────┘
│
│ (WAL)
v
┌──────────────┐
│ Debezium │ (reads WAL, decoded)
└──────┬───────┘
│
v
Kafka: dbserver.public.orders
│
v
Downstream consumers
Debezium
Kafka Connect плагин, поддерживает Postgres (logical replication), MySQL (binlog row), MongoDB (oplog), SQL Server.
Пример топика:
{
"op": "c", // create (c/u/d)
"before": null,
"after": {"id": 123, "user_id": 45, "total": 9900},
"source": {"db": "orders", "table": "orders", "lsn": 12345678}
}
Плюсы
- Никакого кода в приложении -- "невидимое" для сервиса
- Низкая latency (10-100 ms)
- Автоматически покрывает все CRUD операции, включая админ-правки
Минусы
- Schema события = схема БД (а не бизнес-событие)
- Изменение БД ломает downstream -- схемы связаны
- Операционная сложность (Kafka Connect кластер, Debezium, schema registry)
- Нельзя обогащать события бизнес-логикой
Когда CDC
- Репликация БД -> OLAP (ClickHouse, BigQuery)
- Кеш -> БД (Redis updated from PG changes)
- Поисковый индекс (PG -> Elasticsearch)
- Нет возможности менять код (legacy)
Для бизнес-событий лучше outbox: вы контролируете что публикуется.
Решение 3: Event Sourcing
Радикальный подход: не хранить состояние, хранить только события.
state = fold(reduce_fn, initial_state, events)
Order aggregate:
[OrderCreated, OrderPaid, OrderShipped]
Current state computed by replaying events.
При изменении -- добавить новое событие в event store. Событие одновременно:
- Единица записи
- Факт для других сервисов (publish из event store)
- Источник для materialized views
Нет dual-write, потому что нет двух sources of truth -- только event store.
Подробно разобрано в 17.arch-patterns/4.cqrs-es.md.
Плюсы
- Нет dual-write problem вообще
- Полная история изменений (audit, time-travel)
- Replay → любые projections
Минусы
- Радикально меняет модель -- не для каждого проекта
- Сложность: reduce, snapshots, versioning событий
- Сложнее read model (нужен projection worker)
Когда event sourcing
- Финансы, учёт (audit important)
- Сложная домен-модель с богатой историей
- Команда готова к паттерну
Иначе -- overkill. Outbox достаточен для "нам просто нужны события при записи".
Решение 4: Two-Phase Commit (2PC)
Исторический ответ -- распределённая транзакция через координатор.
Coordinator DB1 Queue
│ │ │
│─── prepare ──►│ │
│ │── locks ─────│
│◄── ready ─────│ │
│─── prepare ───────────────► │
│◄── ready ────────────────── │
│ │
│─── commit ───►│ │
│─── commit ─────────────────► │
2PC существует в Java EE (JTA/XA), не применяется в PHP/Go реально.
Плюсы
- ACID-like атомарность через системы
- Стандарт XA
Минусы
- Блокирующий: ресурсы заблокированы на время prepare
- Координатор = SPOF
- Плохая availability: любая нода недоступна -> blocked transaction
- Latency: 2 round trips минимум
- Не все системы поддерживают (Stripe не умеет 2PC)
- Практически не используется в современных микросервисах
Когда 2PC
Почти никогда. Исторически -- banking, XA-совместимые RDBMS. В новых системах -- saga + outbox.
Сравнение решений
| Свойство | Outbox | CDC | Event Sourcing | 2PC |
|---|---|---|---|---|
| Consistency | Eventual | Eventual | Eventual | Strong |
| Latency events | 100-500 ms | 10-100 ms | 0 (write is event) | 2× RTT |
| Код в приложении | + | - | ++ | + |
| Инфраструктура | +1 worker | Kafka Connect | Event store | XA |
| Schema событий | Бизнес | DB schema | Бизнес | - |
| Application-controlled | + | - | + | + |
| Подходит для новых проектов | ✓ лучший дефолт | ✓ репликация | редко | нет |
Когда stale data -- OK
Не все dual-write проблемы нужно решать строго. Иногда eventual consistency нормальна:
- Cache: Redis отстал на секунду -- ок. Если критично -- cache-aside pattern с TTL (см.
02.design-approaches/4.caching-strategies.md). - Search index: заказ появился в поиске через 5 секунд после create -- ок.
- Analytics / dashboards: статистика с лагом в минуту -- всегда ок.
- Notifications: email на минуту позже -- ок.
Для этих случаев достаточно периодического reconciliation job'а:
-- Weekly: find DB rows missing from search index
SELECT o.id FROM orders o
LEFT JOIN search_index s ON s.id = o.id
WHERE s.id IS NULL AND o.created_at < NOW() - INTERVAL '1 hour';
Пошлём эти в индексатор и забудем.
Idempotent retry как "наивное" решение
Часто слышу: "просто сделай retry с идемпотентностью":
db.commit()
retry(fn() => kafka.send(event)) // с idempotency key
Проблема: что если процесс падает между commit и первым send? Retry некому запустить. Потеряли событие.
Работает только с outbox / CDC / event sourcing -- retry должен читать "незавершённое действие" из persistent store. Это и есть outbox.
CDC consumer: псевдокод
<?php
declare(strict_types=1);
/**
* Consume CDC events from Debezium topic.
* Map row-level change to domain events or sync-downstream.
*/
final class DebeziumOrderConsumer
{
public function __construct(
private readonly SearchIndex $index,
private readonly InboxRepo $inbox,
) {}
public function handle(array $envelope): void
{
// Idempotency key = source LSN or (topic, partition, offset)
$eventId = $envelope['source']['lsn'] ?? null;
if ($eventId === null || !$this->inbox->markProcessed($eventId)) {
return; // already handled
}
switch ($envelope['op']) {
case 'c': // create
case 'u': // update
$this->index->upsert(
id: $envelope['after']['id'],
doc: $this->toDoc($envelope['after']),
);
break;
case 'd': // delete
$this->index->remove($envelope['before']['id']);
break;
case 'r': // snapshot read
$this->index->upsert(
id: $envelope['after']['id'],
doc: $this->toDoc($envelope['after']),
);
break;
}
}
private function toDoc(array $row): array
{
return [
'id' => $row['id'],
'userId' => $row['user_id'],
'total' => $row['total_cents'] / 100,
'createdAt' => $row['created_at'],
];
}
}
Нужна ли атомарность DB + очередь?
│
├── Нет (eventual через reconciliation) → skip patterns
│
├── Да, уже есть Kafka Connect в стеке
│ └── CDC (Debezium)
│
├── Да, нужны бизнес-события
│ └── Transactional Outbox
│
├── Да, и это центральная модель (финансы, audit)
│ └── Event Sourcing
│
└── Легаси Java EE с XA-драйверами
└── 2PC (только если уже выбран)
Pitfalls
- Rely на "порядок операций" -- "сначала commit, потом send" не работает. Это самое частое заблуждение.
- Retry без persistent store -- потерянные retries при краше.
- Outbox без мониторинга lag -- тихий рост backlog.
- CDC без schema registry -- изменение колонки сломало все downstream.
- Event sourcing без snapshots -- replay миллиона событий при каждом запросе.
- 2PC в микросервисах 2025 -- не надо.
- Игнорировать eventual -- "просто добавим retry" в 20 местах без outbox.
- Double publish -- outbox + CDC на одной таблице -> события в двух топиках.
Выводы
- Dual-write в распределённой системе атомарно невозможен
- Transactional Outbox -- лучший default для большинства проектов
- CDC (Debezium) -- когда нужна репликация БД в другие системы без кода
- Event Sourcing -- радикальное решение для доменов с богатой историей
- 2PC -- практически не используется в современных стеках
- Иногда eventual consistency -- это нормально; reconciliation job проще событий
- Ни один retry не спасает без persistent store для незавершённых действий
- Выбор зависит от контекста: бизнес-события vs репликация БД vs audit