Идемпотентность фоновых задач и консьюмеров Kafka¶
Заголовок Idempotency-Key привлекает внимание, потому что у него есть имя и спецификация. Но та же проблема возникает у каждого worker очереди, причём чаще: очереди по своей природе доставляют как минимум один раз. Worker отправляет письмо со счётом, погибает до подтверждения, задача приходит снова. Таймаут видимости истекает, пока медленный worker ещё работает, и второй забирает ту же задачу. Партиция Kafka перераспределяется, последний незафиксированный пакет повторяется. Очередь во всех случаях работает правильно. Без дедупликации каждый случай даёт повторный эффект. Нужен тот же ключ идемпотентности, что в HTTP, с одним отличием в его составе.
Числа получены скриптом статьи с Redis 7 в контейнере и внутрипроцессной очередью at-least-once. Версии: idempotency-kit 0.3.0, redis-py 8.1.0, Python 3.13.
Сбой до подтверждения¶
Очередь, не теряющая задачи, хранит их доступными для повторной доставки до подтверждения worker. Порядок работы: выполнить эффект, затем подтвердить. Между ними процесс может погибнуть. Очередь правильно доставит задачу снова, но worker неправильно повторит эффект:
plain handler deliveries=2 mails sent=2
handler under an idempotency key = job id deliveries=2 mails sent=1
В обоих вариантах две доставки. Обычный обработчик отправил счёт дважды. Обработчик с ключом по идентификатору задачи — один раз: повтор нашёл сохранённую запись первого выполнения и вернул результат, не вызывая отправку письма. Это весь механизм из статьи об HTTP, применённый к задаче; хранилище — тот же Redis.
async def handle(job: dict) -> None:
await coordinator.coordinate("mail.invoice", job["job_id"], 3600, adapter, mailer.send_invoice, job["order_id"])
В показанном сценарии гарантия доставки at-least-once и сохранённый результат под ключом дают однократное выполнение задачи. Это та же идея композиции, что в статье об exactly-once для outbox и inbox, только с очередью вместо Kafka.
Два worker, одна задача¶
Более неприятен конкурентный повтор. Worker медленный, таймаут видимости истёк, второй получил ту же задачу, теперь оба исполняют её одновременно. Кеш результата не помогает: никто ещё не закончил. Нужна предварительная резервация:
in_flight='wait' {'A': 'em_1', 'B': 'em_1'} mails sent=1
in_flight='raise' {'A': 'em_1', 'B': 'in progress, requeue'} mails sent=1
Worker A зарезервировал ключ и начал работу. B нашёл резервацию. В режиме wait он дождался результата A и вернул ту же квитанцию: оба подтвердили задачу, очередь чиста. В режиме raise он получил ошибку «задача выполняется». Это подходит очереди, предпочитающей отложенный повтор простаивающему worker: B возвращает задачу, A завершает её, следующая доставка получает сохранённую квитанцию. В обоих режимах отправлено одно письмо. Обычный обработчик отправил бы два, без признаков дублирования в журналах worker.
Из чего составлять ключ¶
Здесь задачи отличаются от HTTP. Для запроса ключ создаёт клиент, и он означает «этот запрос». Для задачи ключ выбирает worker. Очевидный идентификатор задания отличает экземпляр задания, но не обязательно нужный эффект:
key = job id: two jobs for order-9 -> mails sent=2 (the key must be the effect's identity, not the delivery's)
key = 'order-9:invoice': the same two jobs -> mails sent=2 (a colon is not allowed in a key; the record fails validation and the action runs unprotected)
key = 'order-9.invoice': the same two jobs -> mails sent=1
Для одного заказа поставлены две задачи: отправитель повторил постановку либо два участка кода независимо решили отправить счёт. По идентификаторам задач это два ключа и два письма. По эффекту «счёт для заказа 9» — один ключ и одно письмо, сколько бы задач его ни несло. Ключ должен обозначать заказ и действие, а не сообщение, попросившее его выполнить.
Средняя строка — моя ошибка при написании эксперимента, которую я оставил как вполне вероятную ошибку читателя. Ключ содержал двоеточие, зарезервированное библиотекой как разделитель ключа хранения, и не прошёл валидацию. Координатор обрабатывает такую проблему как ошибку хранилища: считает, журналирует и запускает действие без защиты. Две задачи отправили два письма. Правило документировано, ошибка записана в журнал, но именно это нужно проверять тестом: несохраняемый ключ оставляет задачу без идемпотентности, не выбрасывая исключения.
Две доставки могут представлять одну бизнес-операцию. Для дедупликации нужен её идентификатор; внешним эффектам по-прежнему нужна собственная защита от повторов при сбоях.
Потребители: inbox и ключ решают разные задачи¶
У Kafka consumer встречаются все три сценария выше, а ещё есть inbox. Они защищают разные части обработки; потребителю с внешними действиями нужны оба механизма.
Строка inbox фиксируется в одной транзакции с эффектом в БД. Повторное сообщение конфликтует с ней, обработчик не запускается. Это защищает строку счёта, запись ledger, изменение статуса — всё, что записано через ту же session. Для этих изменений однократность обеспечивается общей фиксацией строки и эффекта.
Письмо в эту транзакцию не входит. Если обработчик пишет ledger и отправляет письмо, защищена запись, но не отправка. Сбой между отправкой и commit откатит строку, а повтор отправит письмо снова. Внешнему эффекту нужен собственный ключ: идентификатор сообщения или, лучше, идентичность действия, выведенная из него, с такой же обёрткой, как у задачи выше:
async def on_order_created(event: InboxEvent, repo: InboxEventRepository) -> None:
await repo.session.execute(ledger.insert().values(order_id=event.payload["order_id"])) # the inbox protects this
await coordinator.coordinate("mail.invoice", f"{event.payload['order_id']}.invoice", 86400, # the key protects this
adapter, mailer.send_invoice, event.payload["order_id"])
Одна транзакция для строки, одна резервация для внешнего действия: в этом эксперименте сообщение, доставленное трижды, создаёт одну строку и одно письмо. Для окна между самой внешней отправкой и сохранением её результата нужна также поддержка ключа принимающей стороной; резервация Redis не превращает отправку в транзакцию БД.
Срок хранения¶
Ключ должен жить дольше максимально поздней повторной доставки. У очереди с dead-letter это могут быть дни, а пакетная задача после сбоя может возобновиться с контрольной точки через сутки. Обычно выбирают двадцать четыре часа, для ручного повторного запуска после инцидентов — неделю. Минимальная минута здесь несущественна. Важно пережить срок возможного повторного получения задачи: после истечения записи повтор для механизма ключей снова становится новым действием.
Всё выше — idempotency-kit, тот же координатор и те же резервации, что в HTTP. Между двумя статьями изменилось лишь то, кто выбирает ключ и что тот обозначает.
Суть — в третьей таблице: два, два, один.