В постах про Кафку я упоминала, что консьюмер может повторно получить и обработать одно и то же сообщение (и это не баг), поэтому обработку лучше делать идемпотентной.
Идемпотентность означает, что многократное повторение операции даст тот же эффект, что и однократное.
Предлагаю обсудить чуть детальнее, как сделать обработку идемпотентной.
➡️ Сначала я бы посмотрела на характер сообщений: возможно, сообщения уже обладают "встроенной идемпотентностью", потому что содержат актуальное состояние, а не действие. Например, текущее количество бонусных баллов у пользователя:
{
"userId": "1",
"bonusBalance": 5000
}Если консьюмер обработает такое сообщение дважды и дважды выполнит
SET bonus_balance = 5000, то итоговое состояние останется тем же: 5000 бонусов.Можно провести аналогию с HTTP-методом PUT: повтор одного и того же запроса должен приводить к тому же целевому состоянию ресурса.
✔️ Для снапшотов важен порядок: нельзя после нового состояния применять старое. Поэтому к таким сообщениям лучше добавлять поле version и обновлять состояние только если версия из сообщения новее текущей версии в БД:
{
"userId": "1",
"bonusBalance": 5000,
"version": 2
}➡️ Другое дело, если сообщение несет информацию о каком-то событии, например, начисление бонусов:
{
"userId": "1",
"operationType": "CREDIT",
"amount": 1000
}Смысл события: начислить пользователю 1000 бонусов. Если обработать его дважды, то пользователь получит 2000 бонусов вместо 1000.
Для таких сообщений сделать обработку идемпотентной нам поможет дедупликация, то есть фильтр дублей. Для её реализации обычно используют уникальный идентификатор
messageId/eventId. По смыслу это похоже на Idempotency-Key в REST-запросах.Уникальный идентификатор может быть бизнесовым (например, идентификатор платежа) или техническим (специально сгенерирован на продюсере или до него для дедупликации).
Однажды в качестве ключа идемпотентности мы взяли комбинацию трех бизнес-полей. Аналитики пообщались с бизнесом, со смежниками, всё выглядело логично. А потом уже в процессе эксплуатации выяснилось, что в одном специфическом случае такая же комбинация полей повторяется у связанной сущности. Для нашей логики это были разные сущности, и ключ перестал подходить. Попросили ребят на передающей стороне генерировать для нас технический айдишник.
✔️ Это важный момент на этапе проектирования интеграции: что именно мы защищаем от дублей? Если важно не обработать дважды конкретное сообщение, подойдет технический messageId от продюсера. Если важно не создать второй платеж, то нужен paymentId из системы-источника.
Теперь, когда у нас есть уникальный идентификатор, мы можем фильтровать дубли сообщений.
Например, поле operationId является ключом уникальности.
🔴Если обработка сводится к вставке записи в одну таблицу БД, можно сделать дедупликацию так:
✅ добавить индекс уникальности на поле
operation_id, содержащее ключ идемпотентности✅ при добавлении записи если такой
operation_id уже есть, то ничего не делать (INSERT INTO ... ON CONFLICT DO NOTHING)🔴Если кроме вставки в таблицу с
operation_id, нужны изменения в других таблицах, то делаем всё в одной транзакции.Общая схема:
1. Консьюмер получил сообщение из Кафки
2. Открываем транзакцию
3. Пытаемся вставить запись с
operation_id4. Если вставка прошла, значит это первое получение сообщения, выполняем бизнес-логику. Если получили конфликт уникальности — это дубль, бизнес-логику пропускаем.
5. Коммитим транзакцию БД
6. После этого коммитим оффсет в Кафку
⚡️ Важно!
Если бизнес-логика долгая, например, включает в себя большие расчеты, отправку запросов в другие сервисы и тому подобное, то не нужно ее выполнять тут же до коммита оффсета. При стандартных настройках если консьюмер долго не запрашивает новые сообщения, его выгонят из консьюмер группы, и партицию могут отдать другому консьюмеру (история про захлебнувшихся консьюмеров).
В случае долгой обработки лучше просто сохранить сообщение, закоммитить оффсет, а саму обработку выполнить отдельным процессом. Такой подход называют inbox pattern.
