HardТеория7 min

Dual-write problem и решения

Почему dual-write ломается, решения: Transactional outbox, CDC, Event sourcing, 2PC. ASCII failure scenarios

Суть проблемы

Сервис хочет записать изменение в две системы одновременно:

  • БД + очередь сообщений (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'],
        ];
    }
}
## Decision tree
Нужна ли атомарность 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