ESPN DeportesEl clásico de Manchester quedó definido por el VAR; silbidos para Real Madrid; Chelsea fue superado por HullInquirerDILG chief Remulla visits QC jail ahead of looming Romualdez transferThe Jerusalem PostA good deal in bad neighborhoods: Why the Gulf’s best bet remains in Jerusalem - opinionESPNPower Rankings: Everything we learned from the top 25 in Week 2RTP DesportoEuroVolley 2026. Portugal soma quarta derrota consecutiva frente à UcrâniaBBC NewsI had 11 years of chemotherapy for a cancer I didn't have20 MinutenUrsache unklar: Fischsterben im MühlebachComplete SportsBlackburn Give Injury Update On Super Eagles StarESPN CricinfoFleming to link up with T20I squad in preparation for Test coaching stintBBC عربيجماعة أنصار الله تعلن استهداف قاعدة ثانية في السعودية، ومحمد بن سلمان يلتقي قائد القيادة المركزية الأمريكيةIl Fatto QuotidianoKimi Antonelli ora ha in mano il titolo di Formula 1: quel vantaggio su Russell e il sogno di una passerella trionfaleABC News4 people hospitalized after a crane collapsed at a Miami construction site: Officials
The Daily Newsstand · Free, Always
Monday, September 14, 2026

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

Translate

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

Глава 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), чтобы Kafka не теряла сообщения → Debezium, чтобы outbox не грузил приложение → прямое чтение WAL, код эндпоинта /replication выглядит как код прямого вызова.

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

View the original on Хабр

KioskNews shows a cleaned-up reading view extracted from the publisher’s page — the original always lives on their site, not ours.