Ответ
Задача: Разработать отказоустойчивую систему обработки событий с гарантированной доставкой и идемпотентностью при высокой нагрузке (~10 000 событий в секунду).
Проблемы:
- Потеря событий при сбое после обработки бизнес-логики, но до подтверждения.
- Дублирующая обработка при повторной доставке сообщения из брокера (at-least-once delivery).
- Согласованность между изменением состояния в базе данных и фактом обработки события.
Архитектурное решение (паттерны):
-
Transactional Outbox: Для атомарности сохранения в БД и отправки события.
- Событие сначала сохраняется в той же транзакции, что и бизнес-сущность, в таблицу
outbox. - Отдельный процесс (например, CDC или scheduled task) читает
outboxи публикует события в Kafka.
- Событие сначала сохраняется в той же транзакции, что и бизнес-сущность, в таблицу
-
Идемпотентные обработчики: Каждый обработчик в потребителе проверяет, не обрабатывалось ли событие ранее.
Пример кода обработчика в потребителе Kafka:
@Component
public class EventProcessor {
@Autowired
private EventLogService eventLogService; // Сервис для записи в таблицу обработанных событий
@Autowired
private BusinessService businessService;
@Autowired
private TransactionTemplate transactionTemplate;
@KafkaListener(topics = "order-events")
public void handle(OrderEvent event) {
// 1. Проверка идемпотентности
if (eventLogService.isEventProcessed(event.getId())) {
log.info("Event {} already processed, skipping", event.getId());
return;
}
// 2. Обработка в транзакции
transactionTemplate.execute(status -> {
// Бизнес-логика
businessService.processOrder(event.getOrderId(), event.getData());
// Фиксация факта обработки события в БД
eventLogService.markEventAsProcessed(event.getId());
return null;
});
// 3. После успешной транзакции offset в Kafka будет закоммичен
}
}
Итоговые характеристики системы:
- Гарантированная доставка: Событие не потеряется благодаря Outbox.
- Идемпотентность: Повторная доставка того же события не приведет к дублированию бизнес-эффектов.
- Согласованность: Состояние БД и факт обработки события изменяются атомарно в одной транзакции.