Есть классическая backend-задача. Cоздаем заказ, сохраняем его в БД и хотим отправить событие в Kafka:
Наивно это выглядит так:
@Transactional
public void createOrder(CreateOrderRequest request) {
Order order = orderRepository.save(new Order(request));
kafkaTemplate.send("order-created", new OrderCreatedEvent(order.getId()));
}
И тут есть ряд проблем
База данных и Kafka - это две разные системы. Транзакция в PostgreSQL не управляет Kafka
Kafka не знает, что происходит в транзакции PostgreSQL. А значит, “сохранить заказ” и “отправить событие” не являются одной атомарной операцией 🙃
А еще транзакции должны быть короткими. Пока транзакция открыта, она держит connection к БД, что может стать боттлнеком. Если внутри транзакции делать сетевой вызов в Kafka или внешний сервис, и он подвиснет, у тебя подвиснет вся транзакция. А потом connection pool забьется и новые запросы будут висеть, ждать свободный коннекшн
В общем, это плохой паттерн в транзакции делать какую-то долгую операцию, особенно внешние вызовы
И тут начинаются варианты, как это можно пофиксить:
1️⃣ Можно сначала отправить событие в Kafka, а потом сохранить в БД в транзакици
А что если потом транзакция в БД откатится?
В итоге в Kafka улетело
OrderCreated, но самого заказа в базе нет. Другие сервисы начали реагировать на событие о заказе, которого не существует.2️⃣ Сначала сохранить заказ в БД, а событие отправить после коммита
Но если приложение упадёт между коммитом и отправкой в Kafka?
Заказ в базе будет, а события не будет. Другие сервисы не узнают, что заказ создался
❗️Именно эту проблему решает Transactional Outbox Pattern
Идея простая: в одной транзакции с бизнес-данными мы сохраняем не только заказ, но и событие в отдельную таблицу -
outbox_eventsТо есть внутри одной БД-транзакции происходит сразу две записи:
@Transactional
public void createOrder(CreateOrderRequest request) {
Order order = orderRepository.save(new Order(request));
outboxRepository.save(new OutboxEvent(
OrdersEvents.ORDER_CREATED,
order.getId(),
toJson(new OrderCreatedEvent(order))
));
}
Теперь заказ и событие точно сохраняются атомарно. Либо сохранилось и то, и другое. Либо не сохранилось ничего
🤔 Окей, а дальше как?
А уже потом отдельный процесс или задача по расписанию (poller) читает
outbox_events и публикует события в Kafka
Service → orders + outbox_events → Poller → Kafka
Это может быть обычный poller, который раз в секунду берёт новые события из таблицы, отправляет их в Kafka и помечает как опубликованные.
Либо можно использовать CDC-подход, например Debezium: он читает изменения из журнала WAL SQL базы и отправляет их дальше в Kafka
PostgreSQL WAL → Debezium → Kafka
Но у Outbox есть важный нюанс: обычно он даёт at-least-once delivery. То есть событие может быть доставлено больше одного раза
Например, poller отправил событие в Kafka, но упал до того, как пометил его как опубликованное. После рестарта он отправит это же событие ещё раз
Поэтому consumer’ы должны быть идемпотентными: уметь спокойно обработать одно и то же событие повторно. Обычно для этого используют
eventId и таблицу уже обработанных событий. Писал про идемптоентность тут❓ На собеседованиях по этой теме часто спрашивают
• зачем нужен Outbox Pattern?
• почему нельзя просто отправить событие в Kafka внутри
@Transactional?• что будет, если БД закоммитилась, а Kafka-send упал?
• что такое outbox-таблица?
• почему consumer должен быть идемпотентными и что будет если НЕ?
💡 Outbox Pattern нужен, когда изменение в БД и публикация события должны быть логически связаны, но физически происходят в разных системах. Мы сохраняем событие рядом с бизнес-данными в рамках одной транзакции и потом доставляем его во внешнюю систему
Накидай огня, если пост зашел 🔥
А в комментах расскажи, какие моменты по Кафке тебе еще не понятны, выберу одну из тем и напишу пост 😎
#хардовая_польза
