Когда сервис получает SIGTERM, Uvicorn даёт активным запросам шанс завершиться — но только в рамках --timeout-graceful-shutdown (в наших примерах 25–30 секунд). Если долгая цепочка вызовов не укладывается, asyncio.CancelledError прерывает её на полуслове. Операция успела начаться, но не дошла до конца — и при перезапуске сервис попробует её повторить.
Это безопасно только если операция идемпотентна: повторный вызов с теми же данными даёт тот же результат, не создавая дубликатов. Именно об этом — статья.
Что такое идемпотентность и зачем она нужна
Представьте API платежей: клиент отправил запрос на списание, но ответ не получил — соединение упало. Что делать? Попробовать ещё раз? А если первый запрос всё-таки дошёл?
Если платёжный сервис не идемпотентен, повторный запрос создаст второе списание. Клиент заплатит дважды.
Идемпотентность решает это: если вы передаёте уникальный ключ операции (Idempotency-Key), сервер запоминает результат первого вызова и при повторном просто возвращает тот же ответ — без повторного действия.
При graceful shutdown та же проблема: SIGTERM прерывает запрос в произвольный момент. Новый под получит тот же запрос снова. Без идемпотентности — дубликат.
HTTP POST с Idempotency-Key в FastAPI
Клиент отправляет заголовок Idempotency-Key — один уникальный UUID на операцию. Даже если он повторит запрос десять раз, сервер вернёт сохранённый результат без повторного действия.
from fastapi import APIRouter, Header
from uuid import UUID
router = APIRouter()
@router.post("/payments", status_code=201)
async def charge_order(
order_id: UUID,
amount: int,
idempotency_key: str = Header(..., alias="Idempotency-Key"),
service: PaymentService = Depends(get_payment_service),
) -> PaymentReceipt:
return await service.charge(order_id, amount, idempotency_key)
PaymentService.charge устроен так: сначала ищет запись по idempotency_key в базе. Если нашёл — возвращает сохранённый ответ, не делая повторного списания. Если нет — выполняет списание и сохраняет результат атомарно в одной транзакции.
Схема работы:
Ключ идемпотентности лежит в базе, а не в памяти пода, поэтому переживает перезапуск. Платим лишним поиском и лишней записью на каждый запрос — зато повтор после обрыва не списывает деньги второй раз.
Kafka-listener с защитой от повторной обработки
При работе с Kafka в ручном режиме (manual commit) offset подтверждается явно после обработки события. Риск: если CancelledError прилетит после side-effect, но до commit_offsets — при перезапуске сервис получит то же событие снова.
Защита — таблица processed_event: запись о том, что событие уже обработано, сохраняется в той же транзакции, что и само действие. При повторе обнаруживается конфликт по уникальному ключу — и обработка пропускается.
from aiokafka import AIOKafkaConsumer, TopicPartition
from sqlalchemy.ext.asyncio import AsyncSession
async def handle_order_confirmed(
event: OrderConfirmedEvent,
consumer: AIOKafkaConsumer,
tp: TopicPartition,
offset: int,
session: AsyncSession,
) -> None:
async with session.begin():
already = await session.execute(
select(ProcessedEvent).where(
ProcessedEvent.event_id == event.event_id,
ProcessedEvent.context == "billing",
)
)
if already.scalar_one_or_none():
await consumer.commit({tp: offset + 1})
return
session.add(ProcessedEvent(event_id=event.event_id, context="billing"))
session.add(OutboxEvent(
topic="payments.charge-requested",
partition_key=str(event.order_id),
payload={"order_id": str(event.order_id), "amount": str(event.total_amount),
"idempotency_key": str(event.event_id)},
))
await consumer.commit({tp: offset + 1})
Сам вызов платёжного провайдера в транзакции не делают — это сеть внутри транзакции, и SIGTERM между ответом провайдера и commit оставил бы деньги списанными без следа. Обработчик лишь кладёт задание в outbox, а worker, который пойдёт к провайдеру, передаст тот же idempotency_key в httpx-запрос. Это двойная защита: на уровне Kafka и на уровне downstream-вызова.
Если транзакция откатилась из-за CancelledError — offset тоже не подтверждается, и при перезапуске событие придёт снова. Повторная вставка в processed_event упадёт по уникальному ограничению — обработка безопасно пропустится.
Outbox-relay с двух-фазным статусом
Outbox-паттерн используется для надёжной публикации событий в Kafka: события сначала пишутся в таблицу БД, а отдельный relay-процесс их оттуда читает и отправляет.
Проблема: если SIGTERM прервёт relay между отправкой в Kafka и пометкой строки как PUBLISHED, при перезапуске та же строка уйдёт повторно.
Решение — промежуточный статус PUBLISHING:
from sqlalchemy import select, update
from datetime import datetime, timezone
async def relay_batch(session: AsyncSession, producer: AIOKafkaProducer) -> int:
now = datetime.now(timezone.utc)
picked = (
select(OutboxEvent.id)
.where(OutboxEvent.status == "PENDING")
.order_by(OutboxEvent.id)
.limit(50)
.with_for_update(skip_locked=True)
)
result = await session.execute(
update(OutboxEvent)
.where(OutboxEvent.id.in_(picked))
.values(status="PUBLISHING", locked_at=now)
.returning(OutboxEvent.id, OutboxEvent.payload, OutboxEvent.topic)
)
rows = result.fetchall()
await session.commit() # статус PUBLISHING виден другим relay и восстановителю
if not rows:
return 0
for row_id, payload, topic in rows:
await producer.send_and_wait(topic, value=payload)
await session.execute(
update(OutboxEvent)
.where(OutboxEvent.id == row_id)
.values(status="PUBLISHED", published_at=datetime.now(timezone.utc))
)
await session.commit()
return len(rows)
Захват пачки фиксируется отдельной транзакцией до отправки — иначе промежуточный статус никто не увидит, а при падении relay строки просто вернутся в PENDING вместе с откатом. У UPDATE в PostgreSQL нет LIMIT, поэтому пачку выбирает подзапрос с FOR UPDATE SKIP LOCKED.
Если SIGTERM случается между send_and_wait и UPDATE — строка зависает в статусе PUBLISHING. Отдельная фоновая задача через несколько минут возвращает такие строки обратно в PENDING. Kafka получит повторную отправку — но consumer-side processed_event защитит от дублирования.
Relay-цикл проверяет флаг готовности, а не крутится бесконечно:
async def outbox_loop(app_state: AppState) -> None:
while app_state.is_ready:
sent = await relay_batch(session, producer)
if sent == 0:
await asyncio.sleep(1.0)
Частая ошибка: httpx-retry без Idempotency-Key
# проблемный вариант
async def charge(order_id: UUID, amount: int) -> dict:
async with httpx.AsyncClient() as client:
for attempt in range(3):
try:
resp = await client.post(
f"{PAYMENT_URL}/charge",
json={"order_id": str(order_id), "amount": amount},
timeout=10.0,
)
resp.raise_for_status()
return resp.json()
except httpx.TransportError:
if attempt == 2:
raise
await asyncio.sleep(0.5)
При SIGTERM во время первой попытки: запрос ушёл, ответ не получен. Retry создаёт второй запрос. Провайдер обрабатывает оба — двойное списание.
Правильно — Idempotency-Key генерируется один раз до цикла и передаётся во всех попытках:
async def charge(order_id: UUID, amount: int, idempotency_key: str) -> dict:
async with httpx.AsyncClient() as client:
for attempt in range(3):
try:
resp = await client.post(
f"{PAYMENT_URL}/charge",
json={"order_id": str(order_id), "amount": amount},
headers={"Idempotency-Key": idempotency_key},
timeout=10.0,
)
resp.raise_for_status()
return resp.json()
except httpx.TransportError:
if attempt == 2:
raise
await asyncio.sleep(0.5)
CancelledError внутри транзакции
asyncio.CancelledError при завершении может прилететь в любом await. Если это случится внутри async with session.begin() после side-effect, но до commit — SQLAlchemy выполнит откат. Это правильно: offset тоже не подтверждается, повтор безопасен.
Если код перехватывает CancelledError — нужно либо пробросить его дальше, либо явно откатить транзакцию:
async def handle_customer_merge(event: CustomerMergeEvent, session: AsyncSession) -> None:
try:
async with session.begin():
session.add(ProcessedEvent(event_id=event.event_id, context="crm"))
await crm_service.merge(event.source_id, event.target_id)
except asyncio.CancelledError:
# откат выполнен context manager-ом; пробрасываем дальше
raise
Главное правило: except asyncio.CancelledError: pass внутри транзакции — это всегда ошибка.
Глубже: иногда ключ не нужен: PUT с идентификатором от клиентарасширенное
Прежде чем строить таблицу ключей, стоит проверить, нельзя ли обойтись без неё. Многие операции становятся идемпотентными сами, если идентификатор придумывает клиент:
@router.put("/orders/{order_id}", status_code=201)
async def put_order(order_id: UUID, body: CreateOrderRequest, session: SessionDep) -> OrderResponse:
stmt = insert(Order).values(id=order_id, **body.model_dump()).on_conflict_do_nothing(index_elements=["id"])
await session.execute(stmt)
await session.commit()
return await orders.get(session, order_id)
Повтор это тот же PUT с теми же данными: вставка упирается в первичный ключ, ON CONFLICT DO NOTHING молчит, клиент получает тот же заказ. Ни таблицы ключей, ни сохранённых ответов. Работает, пока операция сводится к одной строке с естественным идентификатором; как только у неё есть побочный эффект снаружи (списание, письмо), ключ возвращается, потому что вторую попытку надо остановить до эффекта, а не после.
Глубже: таблицы растут: срок жизни и уборкарасширенное
Две таблицы, которые появились выше, ключи идемпотентности и обработанные события, растут с каждой операцией и никогда не уменьшаются сами. Через год это самые большие таблицы в базе, и первым это заметит не диск, а замедлившаяся вставка.
Срок жизни определяет клиент повтора. Ключ идемпотентности нужен, пока клиент может повторить запрос: для платёжных провайдеров это обычно сутки, свои ключи держат 24–72 часа. Обработанные события нужны, пока событие может прийти снова, то есть не меньше срока хранения топика; при семи днях retention таблицу чистят по восьмому дню.
Уборку делают пачками, а не одним DELETE: DELETE FROM idempotency_key WHERE id IN (SELECT id FROM idempotency_key WHERE created_at < now() - interval '3 days' LIMIT 5000) в цикле фоновой задачи, пока удаляется что-то; индекс по created_at обязателен. На десятках миллионов строк в сутки дешевле секционировать таблицу по дням и отбрасывать старые секции командой DROP TABLE, которая стоит миллисекунды.
И обратная сторона у провайдера: его ключ тоже живёт ограниченно, и повтор через двое суток с тем же Idempotency-Key может быть принят как новый платёж. Поэтому повторы к чужому сервису ограничивают по времени, а не только по числу попыток.
Коротко
- Идемпотентность — повторный вызов с теми же данными даёт тот же результат без дубликатов. При graceful shutdown это обязательное свойство для любой операции, которую SIGTERM может прервать.
- HTTP POST —
Idempotency-Keyв заголовке; FastAPI читает его черезHeader(...), сервис сохраняет результат атомарно в одной транзакции. - Kafka-listener —
processed_eventи side-effect коммитятся в одной транзакции; offset подтверждается только после commit. При откате — повтор безопасен. - Outbox-relay — двух-фазный статус
PENDING → PUBLISHING → PUBLISHED; зависшие строки возвращаются вPENDINGфоновой задачей. - httpx-retry —
Idempotency-Keyгенерируется один раз до retry-цикла и передаётся во всех попытках. CancelledErrorвнутри блокаsession.begin()нужно пробрасывать: проглотить его внутри блока значит выйти из него штатно и зафиксировать половину работы.- Операция с идентификатором от клиента идемпотентна сама по себе через
PUTиON CONFLICT DO NOTHING; таблицы ключей и обработанных событий живут ограниченный срок и чистятся пачками или секциями по дням.
Что пощупать
Повтор запроса, которого клиент не дождался, в сервисе заказов практикума remodov/marketplace-system-python безопасен благодаря ключу идемпотентности: ключ и заказ пишутся в одной транзакции, проигравший гонку откатывает свой заказ исключением внутри uow.begin() и читает чужой. Тест с восемью одновременными запросами через asyncio.gather показывает, что все восемь клиентов получают один и тот же заказ.
Код: create_order.py, idempotency_repository.py, test_idempotency.py.
Сделаем сами
Ветка step-09-idempotency.