Это история про архитектуру одного сервиса. Рассказанная до конца, а не только до того места, где обычно останавливаются статьи про архитектуру. Намеренно упрощена бизнес-логика, чтобы не размывать основной смысл.
Глава 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-сервис никогда не узнает, что событие произошло.
И это не экзотика на 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.Outboxbuilder.Services.AddOutbox(options =>{ options.PollingInterval = TimeSpan.FromMilliseconds(100); // Для тестов options.BatchSize = 50; options.MaxAttempts = 3;}).WithEfCore<AppDbContext>().AddMessageOutbox();// 2. Регистрируем наш адаптер, который просто дергает Kafkabuilder.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); ...});
Проблема dual write решена. Но взамен мы получили новую.
Цена атомарности
|
Цена |
Суть |
|---|---|
|
Write amplification |
|
|
Vacuum |
Постоянно обновляемая outbox-таблица генерирует мёртвые кортежи |
|
Фоновый диспетчер |
Ещё один процесс, конкурирующий за соединения к Postgres |
|
Задержка |
Интервал опроса. Не миллисекунды. Сотни миллисекунд. |
И главное — вся эта нагрузка живёт внутри того же процесса и той же базы, что обслуживает бизнес-логику. Частый поллинг создаёт постоянный фоновый I/O даже тогда, когда сообщений нет.
Глава 4. Debezium. Выносим боль за пределы сервиса
Всю эту нагрузку не обязательно держать в том же процессе и постоянно дергать базу. Postgres и так пишет каждое изменение в WAL (Write-Ahead Log). Debezium — это Kafka Connect коннектор, который читает WAL с помощью логической репликации и публикует результат в Kafka-топик. Приложение делает только INSERT. Outbox не нужен. Всю работу по надёжной доставке берёт на себя отдельный процесс.
Опрос сообщества 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" }}
Важные настройки:
-
slot.drop.on.stop: true— без этого при удалении коннектора слот останется висеть и будет копить WAL, пока не кончится диск. Подробнее о том, как настраивать репликацию и администрировать слоты. -
poll.interval.ms: 100— по умолчанию 500 мс. Нужно для тестов.
Мы решили проблему нагрузки на приложение. Мы не решили проблему количества движущихся частей: теперь у нас 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. Отдельного «отправить в очередь» просто не требуется:
Ноль 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);});
Глава 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.msKafka Connect (по умолчанию 500 мс). Уменьшение до 100 мс: 426 → 507 итер/сек. Узкое место — таймер опроса. -
ZeroAlloc.Outbox: Задержки дает не интервал опроса, а последовательная отправка внутри батча (
BatchSize = 50). Уменьшение интервала с 1с до 100мс: 28 → 56 итер/сек (×2, а не ×10). Дальнейшее уменьшение бессмысленно — нужно делать батч на клиенте, но для этого надо писать свой Outbox.
Глава 8. «А что если нагрузка вырастет?»
Практически любой разговор о нужности Кафки сводится к этому аргументу.
Но насколько реально можно поиметь проблемы:
-
Вряд ли у вас будет так много данных https://topicpartition.io/definitions/small-data, ваш бизнес и кодовая база не растут так быстро.
-
По данным aiven.io 80% Kafka кластеров не превышают 1 МБ/с
-
Отчет RedPanda за 2023-2024 год показывает, что у 56% компаний трафик данных ≤1 МБ/с
-
Postgres достаточно быстр чтобы на современных дисках держать огромные нагрузки
Наш собственный замер: до 5 600 сообщений в секунду и 1,3 МБ/с через прямую WAL-репликацию — на одном CPU, без единой оптимизации.
Подавляющее большинство систем, которые сегодня платят за Kafka, физически не приближаются к границам, за которые логическая репликация не выходит.
Когда Kafka действительно нужна
-
Много разнородных downstream-потребителей не под вашим контролем
-
Десятки тысяч сообщений/сек и растёт
-
Бизнесу нужна долгая история с произвольной перемоткой \ пререпроигрыванием истории, и вы можете это сделать за разумное время
-
Много медленных консьюмеров и их число меняется (не совместимо с предыдущим)
Когда логической репликации достаточно
-
Вы владеете и продюсером, и консьюмером
-
Нагрузка до десятков тысяч сообщений/сек
-
Хотите удешевить инфраструктуру
-
Критична консистентность и задержки
Финал: мы вернулись туда, откуда начали
Мы сделали полный круг. Прямой вызов → Kafka, чтобы разнести сервисы → outbox (через ZeroAlloc.Outbox), чтобы Kafka не теряла сообщения → Debezium, чтобы outbox не грузил приложение → прямое чтение WAL, код эндпоинта /replication выглядит как код прямого вызова.
Мораль: Используйте Postgres, пока не столкнулись с проблемами масштабирования. Когда столкнётесь — у вас будет конкретная измеренная метрика, чтобы обосновать добавление Kafka или другого компонента. А не абстрактное «так принято в микросервисах».
ссылка на оригинал статьи https://habr.com/ru/articles/1082062/