КДПВ
Картинка для привлечения внимания

Это история про архитектуру одного сервиса. Рассказанная до конца, а не только до того места, где обычно останавливаются статьи про архитектуру. Намеренно упрощена бизнес-логика, чтобы не размывать основной смысл.

Глава 1. В начале было просто

Сервис принимал запрос, писал строку в базу, отвечал 201. Это весь код:

app.MapPost("/direct", async (Message dto, AppDbContext db, CancellationToken ct) =>
{
    Message msg = new() { Content = dto.Content, CreatedAt = DateTime.UtcNow };
    db.Messages.Add(msg);
    await db.SaveChangesAsync(ct);
    return Results.Created($"/messages/{msg.Id}", msg);
});

Никто не пишет статей про этот код, потому что в нём нечего обсуждать. Но именно об него разбиваются все последующие решения — держите его в голове, он ещё пригодится в финале.

Глава 2. «Нам нужна очередь» — и вот почему это не глупость

Рано или поздно один сервис перестаёт быть одним сервисом. Появляется второй, которому важно узнать о том же событии — обсчитать аналитику, обновить поисковый индекс, отправить письмо. Возникает соблазн просто дёрнуть его по HTTP, и это первая ошибка, которую все совершают и все же исправляют: синхронный вызов делает вас настолько же надёжным, насколько надёжен самый хрупкий из ваших соседей.

Значит нужно асинхронно, через брокер. И почти всегда этим брокером оказывается Kafka. Это не карго-культ, а рациональный выбор:

  • 80%+ компаний из Fortune 100 её используют, клиентские библиотеки есть для всех языков (kafka.apache.org/powered-by)

  • Проверена в бою: выросла из LinkedIn, где гоняли миллиарды событий в день. Netflix, Uber, Goldman Sachs — все на ней (подробнее о том, как Kafka устроена внутри)

  • Масштабируется линейно: партиции + consumer groups, добавляй брокеров без изменения кода

  • Готовая модель доставки: pull, персистентный лог, репликация, настраиваемое хранение

  • В конце концов это модно.

В коде это выглядит так:

db.Messages.Add(msg);
await db.SaveChangesAsync(ct);                    // (1) записали в базу
await producer.ProduceAsync(KafkaConsumer.Topic, new()
{
    Timestamp = new(msg.CreatedAt),
    Key = msg.Id,
    Value = msg
}, ct);                                            // (2) отправили в Kafka

Спросите себя: что случится, если процесс упадёт между строкой (1) и строкой (2)? В базе — запись есть. В Kafka — ничего. Downstream-сервис никогда не узнает, что событие произошло.

Клиент → HTTP POST → [Сервис] → INSERT → Postgres → HTTP 201
схема 1

И это не экзотика на 0.001% инцидентов. CancellationToken по закрытию соединения клиента — таймаут, закрытие или обновление страницы, кроме того случается рестарт пода, деплой, OOM-killer, да просто исключение внутри ProduceAsync — любое из этого гарантированно создаёт дыру.

Это dual write problem: два независимых ресурса обновляются не атомарно. Нельзя обернуть INSERT в Postgres и ProduceAsync в Kafka в одну транзакцию. Они просто не знают друг о друге.

Делать запись и отправку в обратном порядке - получается еще хуже.

«Просто ретраить» не помогает:

  • Сервис падает до ретрая → событие потеряно

  • Ретрай проходит, но и оригинал прошёл → дубликат

  • Ретраим и базу тоже → дубликат заказа

Глава 3. Гарантия отправки

Может, распределённая транзакция? 2PC? К сожалению нет:

  • Kafka не поддерживает XA. RabbitMQ не поддерживает. SQS не поддерживает.

  • 2PC блокирующий: координатор падает → участники висят с локами бесконечно

  • 30–40% потери пропускной способности по сравнению с локальными транзакциями

  • Требует, чтобы все участники были доступны одновременно

Вот к чему мы на самом деле пришли: если вы пишете в Kafka напрямую из бизнес-транзакции — вы уже нарушаете гарантии. Спорить с этим бессмысленно — можно только либо принять outbox, либо жить с потерянными сообщениями.

На помощь приходит паттерн Transactional Outbox.

  • Его называют «каноническим решением» для надёжной публикации событий — формулировка из разбора у Chris Richardson, microservices.io.

  • AWS Prescriptive Guidance рекомендует именно его.

  • Confluent включает outbox как обязательный шаг в собственный курс по микросервисам.

Идея простая до гениальности: записать и результат работы, и намерение отправить сообщение в одной транзакции. Обе таблицы — в одной базе. Одна ACID-транзакция. Либо обе записи закоммичены, либо обе откачены.

Вам не нужно писать диспетчер outbox и таблицу руками. Оно уже давно реализовано в сотнях библиотек. Например ZeroAlloc.Outbox. Всё, что от вас требуется — это зарегистрировать сервисы в DI и написать крошечный класс-адаптер для отправки в Kafka.

Вот как выглядит настройка в Program.cs:

// 1. Регистрируем сам Outbox из библиотеки ZeroAlloc.Outbox
builder.Services.AddOutbox(options =>
{
    options.PollingInterval = TimeSpan.FromMilliseconds(100); // Для тестов
    options.BatchSize = 50;
    options.MaxAttempts = 3;
})
.WithEfCore<AppDbContext>()
.AddMessageOutbox();

// 2. Регистрируем наш адаптер, который просто дергает Kafka
builder.Services.AddTransient<IOutboxDispatcher<Message>, OutboxDispatcher>();

И сам адаптер — 5 строк кода, библиотека сама берет на себя фоновый опрос, батчинг, ретраи и транзакционность:

public class OutboxDispatcher(IProducer<int, Message> producer) : IOutboxDispatcher<Message>
{
    public async ValueTask DispatchAsync(Message message, CancellationToken ct)
    => await producer.ProduceAsync(KafkaConsumer.Topic, new()
    {
        Timestamp = new(message.CreatedAt),
        Key = message.Id,
        Value = message
    }, ct);
}

В коде эндпоинта мы просто инжектим IOutboxWriter<Message> и пишем в той же транзакции:

app.MapPost("/outbox", async (Message dto, AppDbContext db, IOutboxWriter<Message> outbox, ...) =>
{
    var id = await db.Database.CreateExecutionStrategy().ExecuteInTransactionAsync(async (ct) =>
    {
        Message msg = new() { Content = dto.Content, CreatedAt = DateTime.UtcNow };
        db.Messages.Add(msg);
        await db.SaveChangesAsync(ct);
        
        // ← Пишем в outbox в ТОЙ ЖЕ транзакции. ZeroAlloc.Outbox сам всё сохранит
        await outbox.WriteAsync(msg, ct: ct);      
        return msg.Id;
    }, ct => Task.FromResult(false), ct);
    ...
});
Postgres (WAL) → Debezium (Kafka Connect) → Kafka topic → консьюмер(ы)
схема 2

Проблема dual write решена. Но взамен мы получили новую.

Цена атомарности

Цена

Суть

Write amplification

INSERT + INSERT + UPDATE на каждое сообщение

Vacuum

Постоянно обновляемая outbox-таблица генерирует мёртвые кортежи

Фоновый диспетчер

Ещё один процесс, конкурирующий за соединения к Postgres

Задержка

Интервал опроса. Не миллисекунды. Сотни миллисекунд.

И главное — вся эта нагрузка живёт внутри того же процесса и той же базы, что обслуживает бизнес-логику. Частый поллинг создаёт постоянный фоновый I/O даже тогда, когда сообщений нет.

Глава 4. Debezium. Выносим боль за пределы сервиса

Всю эту нагрузку не обязательно держать в том же процессе и постоянно дергать базу. Postgres и так пишет каждое изменение в WAL (Write-Ahead Log). Debezium — это Kafka Connect коннектор, который читает WAL с помощью логической репликации и публикует результат в Kafka-топик. Приложение делает только INSERT. Outbox не нужен. Всю работу по надёжной доставке берёт на себя отдельный процесс.

Postgres (WAL) → Debezium (Kafka Connect) → Kafka topic → консьюмер(ы)
схема 3

Опрос сообщества Debezium 2026 года показал, что 91.3% респондентов уже активно используют его в проде (результаты опроса), а список пользователей включает организации разного масштаба. Выглядит так, что решению можно доверять.

Конфигурация коннектора:

{
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "table.include.list": "public.messages",
    "plugin.name": "pgoutput",
    "slot.name": "debezium_slot",
    "publication.name": "debezium_publication",
    "slot.drop.on.stop": "true",
    "snapshot.mode": "no_data",
    "poll.interval.ms": "100"
  }
}

Важные настройки:

Мы решили проблему нагрузки на приложение. Мы не решили проблему количества движущихся частей: теперь у нас Postgres, Kafka и отдельная JVM для Kafka Connect.

Глава 5. Стоп. А зачем Kafka в этой схеме?

Посмотрите ещё раз на диаграмму главы 4. Debezium читает WAL, который Postgres пишет в любом случае. Единственное, что реально делает Kafka в этой конкретной схеме — это ретранслирует то, что уже лежит в WAL, ещё через один сетевой хоп.

WAL — это последовательный, append-only лог с независимыми читателями, каждый из которых хранит свою позицию. Опишите Kafka человеку, который не знает, что это Kafka, — вы только что описали WAL.

Kafka

PostgreSQL

Topic

Таблица

Partition + фильтр

PUBLICATION

Consumer + offset

Replication slot

Broker

WAL + pgoutput

Produce

INSERT

Consume

Чтение потока репликации

Если оба потребителя события — ваш код, то Kafka в этой схеме — это плата за ретрансляцию того, что уже существует. INSERT уже попал в WAL. Отдельного «отправить в очередь» просто не требуется:

 ┌────────────────────────────────────┐ │  BEGIN TRANSACTION                 │ │    INSERT INTO messages (...);     │  ← это одновременно и запись, │  COMMIT                            │     и «отправка в очередь» └────────────────────────────────────┘ │ ▼  (автоматически, через WAL) Replication slot → подписчик получает InsertMessage
схема 4

Ноль write amplification. Ноль outbox-таблиц. Нет фоновых диспетчеров. Нет отдельной JVM для Kafka Connect.

Глава 6. Читаем логическую репликацию

Настройка на стороне БД, лучше прямо в миграции:

CREATE PUBLICATION rep_pub FOR TABLE messages;
SELECT * FROM pg_create_logical_replication_slot('rep_slot', 'pgoutput');

Для чтения используем Npgsql.Replication — часть штатного драйвера. Никаких сторонних библиотек:

await using var conn = new LogicalReplicationConnection(connectionString);
await conn.Open(stoppingToken);
var slot = new PgOutputReplicationSlot("rep_slot");
await foreach (var message in conn.StartReplication(
    slot, new PgOutputReplicationOptions("rep_pub", PgOutputProtocolVersion.V4, binary: true),
    stoppingToken))
{
    if (message is InsertMessage insertMessage)
    {
        Message msg = await ReadMessageAsync(insertMessage, stoppingToken);
        completions.Complete(msg.Id, msg);
    }
    conn.SetReplicationStatus(message.WalEnd);  // обязательно!
    await conn.SendStatusUpdate(stoppingToken);
}

Пример чтения логической репликации Postgres я уже приводил в статье Ваш кэш в Redis неэффективен, что с этим делать?

Код эндпонита:

app.MapPost("/replication", async (Message dto, AppDbContext db, ...) =>
{
    Message msg = new() { Content = dto.Content, CreatedAt = DateTime.UtcNow };
    db.Messages.Add(msg);
    await db.SaveChangesAsync(ct);
    return Results.Created($"/messages/{msg.Id}", msg);
});
 ┌────────────────────────────────────┐ │  BEGIN TRANSACTION                 │ │    INSERT INTO messages (...);     │  ← это одновременно и запись, │  COMMIT                            │     и «отправка в очередь» └────────────────────────────────────┘ │ ▼  (автоматически, через WAL) Replication slot → подписчик получает InsertMessage
схема 5

Глава 7. Бенчмарк

Приложение в Aspire с эндпонитами, как в статье. Продьюсер и консьюмер в одном процессе. Консьюмер просто сигналит вызывающему потоку, что сообщение обработано. Код по ссылке https://github.com/gandjustas/habr-aspnet-kafka.

Методика: k6, 250 виртуальных пользователей, каждый запрос ждёт реального подтверждения доставки. Железо — Intel i9-9900KF, 64 ГБ RAM, вся инфраструктура в контейнерах.

Результат:

Эндпоинт

Итер/сек

avg

p95

CPU приложения на сообщение

/replication

5 625

43.9 ms

65.1 ms

0.382 ms

/direct

5 133

48.0 ms

76.9 ms

0.353 ms

/naive (Kafka)

705

350.5 ms

410.4 ms

1.274 ms

/debezium (100мс)

477

507.8 ms

908.4 ms

1.647 ms

/outbox (100мс)

56.2

4 378.4 ms

4 804.5 ms

3.250 ms

/replication обгоняет /naive в 8 раз по throughput и по задержке. Причём в CPU-времени на сообщение разрыв ещё честнее: 0.382 мс против 1.274 мс — Kafka-путь в 3.3 раза дороже по факту потраченных тактов, при этом менее надежен.

Почему Debezium и outbox гораздо медленнее

  • Debezium: poll.interval.ms Kafka Connect (по умолчанию 500 мс). Уменьшение до 100 мс: 426 → 507 итер/сек. Узкое место — таймер опроса.

  • ZeroAlloc.Outbox: Задержки дает не интервал опроса, а последовательная отправка внутри батча (BatchSize = 50). Уменьшение интервала с 1с до 100мс: 28 → 56 итер/сек (×2, а не ×10). Дальнейшее уменьшение бессмысленно — нужно делать батч на клиенте, но для этого надо писать свой Outbox.

Глава 8. «А что если нагрузка вырастет?»

Практически любой разговор о нужности Кафки сводится к этому аргументу.

Но насколько реально можно поиметь проблемы:

Наш собственный замер: до 5 600 сообщений в секунду и 1,3 МБ/с через прямую WAL-репликацию — на одном CPU, без единой оптимизации.

Подавляющее большинство систем, которые сегодня платят за Kafka, физически не приближаются к границам, за которые логическая репликация не выходит.

Когда Kafka действительно нужна

  • Много разнородных downstream-потребителей не под вашим контролем

  • Десятки тысяч сообщений/сек и растёт

  • Бизнесу нужна долгая история с произвольной перемоткой \ пререпроигрыванием истории, и вы можете это сделать за разумное время

  • Много медленных консьюмеров и их число меняется (не совместимо с предыдущим)

Когда логической репликации достаточно

  • Вы владеете и продюсером, и консьюмером

  • Нагрузка до десятков тысяч сообщений/сек

  • Хотите удешевить инфраструктуру

  • Критична консистентность и задержки

Финал: мы вернулись туда, откуда начали

Мы сделали полный круг. Прямой вызов → Kafka, чтобы разнести сервисы → outbox (через ZeroAlloc.Outbox), чтобы сообщения не терялись → Debezium, чтобы outbox не грузил приложение → прямое чтение WAL, код эндпоинта /replication выглядит как код прямого вызова.

Мораль: Используйте Postgres, пока не столкнулись с проблемами масштабирования. Когда столкнётесь — у вас будет конкретная измеренная метрика, чтобы обосновать добавление Kafka или другого компонента. А не абстрактное «так принято в микросервисах».