Очередь на Postgres: почему Kafka не нужна

Это история про архитектуру одного сервиса. Рассказанная до конца, а не только до того места, где обычно останавливаются статьи про архитектуру. Намеренно упрощена бизнес-логика, чтобы не размывать основной смысл.
Глава 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.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);
...
});
Проблема 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 или другого компонента. А не абстрактное «так принято в микросервисах».
KioskNews shows a cleaned-up reading view extracted from the publisher’s page — the original always lives on their site, not ours.