В КУРСЕ?

Разбираемся в теме

Kafka доставила событие повторно: как не выполнить одно действие дважды в Python-сервисе

Асинхронный сервис на Python может успешно изменить данные и затем аварийно завершиться до сохранения позиции чтения в Kafka. После перезапуска событие окажется обработано повторно. Для такой ситуации недостаточно общего обещания надёжной очереди. Рассмотрим обычного потребителя, который фиксирует позицию после обработки, и отделим повторное получение сообщения от повторного бизнес-результата во внешней базе данных.

Получить сообщение и завершить его обработку не одно и то же

В Kafka позиция потребителя позволяет определить, откуда продолжать чтение. Представим событие с идентификатором Е1: добавить три единицы в учебный счётчик резервов. Сервис записал изменение в базу, но завершился до фиксации следующей позиции чтения. При восстановлении он может снова получить Е1. Если обработчик без проверки повторит прибавление, счётчик увеличится на шесть вместо трёх. Это не обязательно означает, что исходное событие было отправлено дважды. Повтор может возникнуть на стороне обработки после сбоя между двумя самостоятельными действиями.

Идемпотентность связывают с устойчивым идентификатором

Обработчик можно спроектировать так, чтобы повтор одного события не создавал дополнительного результата. В учебной модели база хранит уникальный идентификатор обработанного события вместе с соответствующим изменением счётчика в одной транзакции. Тогда оба изменения сохраняются вместе либо не сохраняются. При повторном Е1 уже зарегистрированный идентификатор не позволяет снова применить прибавление. Простого отдельного чтения и последующей записи недостаточно для защиты от всех гонок: нужны согласованные ограничения и атомарность. Также идентификатор должен оставаться тем же при повторе, иначе система увидит новое событие вместо повторного.

Гарантия заканчивается там, где заканчивается согласованная операция

Транзакционные возможности Kafka позволяют согласовывать чтение и запись в её потоках при соответствующей настройке. Но они не делают любое действие во внешней системе автоматически однократным. Запись в отдельную базу, отправка письма и обращение к стороннему сервису имеют собственные границы. В нашем примере защита касается только счётчика и журнала событий, сохранённых совместно. Если после этого приложение выполняет ещё одно независимое действие, его повторяемость надо анализировать отдельно. Асинхронный код и увеличение числа обработчиков сами по себе эту задачу не решают.

Попробуйте на практике

Проследите два сбоя на бумаге. Настоящий кластер Kafka, база данных, программный код и доступ к чужим системам не нужны.

  1. Нарисуйте исходное состояние: счётчик резервов равен нулю, журнал обработанных событий пуст. Запишите Е1 с действием прибавить три и расположите фиксацию позиции чтения после обработки.
  2. Сначала рассмотрите незащищённый вариант: прибавление сохранилось, процесс завершился до фиксации позиции. Повторите событие и получите шесть, объяснив происхождение лишнего результата.
  3. Теперь объедините запись Е1 и увеличение счётчика в одну условную транзакцию с уникальностью идентификатора. После успешного сохранения и такого же сбоя повтор Е1 должен оставить счётчик равным трём.
  4. Рассмотрите сбой до сохранения транзакции: журнал и счётчик остаются прежними, поэтому повтор может применить изменение один раз. Отдельно отметьте, что отправка письма вне этой транзакции её защитой не охвачена.

Как проверить результат. В незащищённой модели получилось шесть, в согласованной три. Сбой до сохранения не оставил половины операции. Фиксация позиции не перепутана с подтверждением внешнего результата, а гарантия не распространена на независимые действия.

Частые вопросы

Поможет ли генерировать новый идентификатор при каждой попытке?

Для распознавания повтора это помешает: одинаковое бизнес-событие будет выглядеть новым. Правило идентификации должно соответствовать событию и сохраняться при его повторной обработке.

Достаточно ли включить идемпотентного производителя?

Это защищает определённый участок публикации в Kafka, но не заменяет анализ внешних действий потребителя. Повторное изменение отдельной базы требует собственного согласованного решения.

Самостоятельный разбор темы. Содержание конкретной обучающей программы здесь не представлено.

Зарегистрируйтесь, чтобы уточнить возможность доступа к этому материалу

Зарегистрироваться
← К списку материалов