Skip to content

Repository files navigation

Async Payments Processing Service

Микросервис асинхронной обработки платежей. API принимает запрос на оплату, надёжно (через transactional outbox) кладёт событие в RabbitMQ, воркер эмулирует обращение к платёжному шлюзу и уведомляет клиента через webhook.

Реализованы: идемпотентность, transactional outbox, ретраи с экспоненциальной задержкой, dead-letter queue, аутентификация по X-API-Key, миграции, Docker.


Стек

FastAPI · Pydantic v2 · SQLAlchemy 2.0 (async) · PostgreSQL · RabbitMQ (FastStream) · Alembic · Docker / docker-compose.

Архитектура

   client
     │  POST /api/v1/payments  (X-API-Key, Idempotency-Key)
     ▼
 ┌────────┐   1 транзакция    ┌──────────────────────────┐
 │  API   │ ───────────────►  │ payments + outbox (PG)   │
 └────────┘                   └──────────────────────────┘
     ▲                                   │  outbox relay (polling)
     │ GET /payments/{id}                ▼
     │                          exchange "payments" ──► queue payments.process
     │                                                        │
     │                                                   ┌──────────┐
     │                                                   │ consumer │  эмуляция шлюза (2–5с)
     │                                                   └──────────┘
     │                       успех │                         │ ошибка
     │                             ▼                         ▼
     │                    status=succeeded         payments.retry.{0,1}  (TTL backoff)
     │                       + webhook                       │ (после N попыток)
     │                                                       ▼
     └───────────────────────────────────────────────  payments.dlq
                                                      status=failed + webhook

Поток:

  1. API в одной транзакции создаёт Payment(pending) и запись в outbox.
  2. Фоновый outbox-relay (внутри api) публикует событие в обменник payments.
  3. consumer читает payments.process, эмулирует шлюз (2–5 c, 90% успех / 10% ошибка), обновляет статус и шлёт webhook.
  4. При ошибке сообщение уходит в очередь ретрая с TTL (экспоненциальный backoff) и возвращается в обработку; после исчерпания попыток — в payments.dlq, платёж помечается failed.

Структура проекта

app/
  api/payments.py       # HTTP-эндпоинты (создание/получение)
  services/payment_service.py  # идемпотентное создание + outbox в 1 транзакции
  repositories/         # СЛОЙ ДОСТУПА К ДАННЫМ (Repository + Unit of Work)
    payment.py          #   PaymentRepository (Protocol) + реализация на SQLAlchemy
    outbox.py           #   OutboxRepository (Protocol) + реализация
    unit_of_work.py     #   UnitOfWork: репозитории в одной транзакции (commit/rollback)
  outbox/relay.py       # publisher: outbox -> RabbitMQ
  worker/consumer.py    # обработка платежа, ретраи, DLQ
  worker/webhook.py     # доставка webhook с ретраями
  broker.py             # топология RabbitMQ (exchange/queues/retry/DLQ)
  models.py schemas.py config.py database.py security.py enums.py
  main.py               # FastAPI (сервис api)
  worker_app.py         # FastStream (сервис consumer)
alembic/                # миграции
docker-compose.yml Dockerfile requirements.txt

Слои. Бизнес-логика (services, worker, outbox/relay) обращается к данным только через абстракции repositories/ и UnitOfWork, а не к AsyncSession напрямую. Это изолирует ORM/SQLAlchemy от бизнес-кода: репозиторий можно подменить моком в тестах, а хранилище/ORM — заменить, не трогая сервисы. Транзакционная граница (включая запись агрегата и outbox-события одним commit) принадлежит UnitOfWork.

Быстрый старт

docker compose up --build

Поднимется: postgres, rabbitmq (UI http://localhost:15672, guest/guest), migrate (накатит миграции и выйдет), api (http://localhost:8000, Swagger /docs), consumer.

API

Все эндпоинты требуют заголовок X-API-Key (по умолчанию secret-api-key-change-me).

Создать платёж

curl -X POST http://localhost:8000/api/v1/payments \
  -H "X-API-Key: secret-api-key-change-me" \
  -H "Idempotency-Key: order-42" \
  -H "Content-Type: application/json" \
  -d '{
    "amount": "199.90",
    "currency": "RUB",
    "description": "Подписка Pro",
    "metadata": {"order_id": 42},
    "webhook_url": "https://webhook.site/<your-uuid>"
  }'

Ответ 201 Created — платёж в статусе pending. Повтор того же запроса с тем же Idempotency-Key вернёт 200 OK и тот же платёж (нового не создаётся).

Получить платёж

curl http://localhost:8000/api/v1/payments/<id> \
  -H "X-API-Key: secret-api-key-change-me"

Через несколько секунд статус станет succeeded (или failed после ретраев). Если указан webhook_url, на него придёт уведомление о результате.

Как реализованы требования

  • Transactional Outbox. Payment и OutboxEvent пишутся в одной транзакции (services/payment_service.py). Отдельный relay (outbox/relay.py) с SELECT … FOR UPDATE SKIP LOCKED публикует события в брокер и помечает их published. Гарантия at-least-once: даже при сбое публикации событие переотправится, дубликаты обезвреживает идемпотентность потребителя.
  • Идемпотентность. (1) API — уникальный idempotency_key; повторный запрос возвращает существующий платёж (защита от гонок — через IntegrityError). (2) Consumer — перед обработкой проверяет, что платёж ещё pending; повторная доставка сообщения не переобрабатывает платёж.
  • Retry + экспоненциальный backoff. На каждый «тир» ретрая — отдельная очередь payments.retry.{i} с x-message-ttl = base·2^i и x-dead-letter-exchange обратно в payments.process. Отдельная очередь на тир исключает head-of-line blocking. По умолчанию MAX_ATTEMPTS=3, RETRY_BASE_DELAY=5 → задержки 5с и 10с.
  • Dead Letter Queue. После исчерпания попыток сообщение публикуется в payments.dlq с заголовком x-death-reason, платёж → failed.
  • Webhook + ретраи. worker/webhook.py шлёт POST с экспоненциальным backoff (по умолчанию 3 попытки). Ретраи webhook не переобрабатывают платёж.
  • Аутентификация. X-API-Key на всех эндпоинтах (security.py).

Конфигурация (env)

Переменная Значение по умолчанию Назначение
DATABASE_URL postgresql+asyncpg://…/payments подключение к PostgreSQL
RABBITMQ_URL amqp://guest:guest@…:5672/ подключение к RabbitMQ
API_KEY secret-api-key-change-me ключ для X-API-Key
SUCCESS_RATE 0.9 вероятность успешной обработки
MAX_ATTEMPTS 3 попыток обработки до DLQ
RETRY_BASE_DELAY 5 база backoff (сек): base·2^i
WEBHOOK_MAX_ATTEMPTS 3 попыток доставки webhook

Демонстрация ретраев и DLQ

Чтобы гарантированно увидеть ретраи и попадание в DLQ, поднимите сервис с SUCCESS_RATE=0 (каждая обработка «падает»):

SUCCESS_RATE=0 docker compose up --build

Создайте платёж — в логах consumer будут видны ретраи (#1, #2) и финальное -> DLQ; в RabbitMQ UI появится сообщение в payments.dlq, платёж получит статус failed. Удобно ловить webhook на https://webhook.site.

Замечания и компромиссы

  • Один обработчик-consumer, как и просили: обработка + статус + webhook в одном месте.
  • Outbox-relay запущен внутри api как фоновая задача (можно вынести в отдельный сервис; FOR UPDATE SKIP LOCKED уже позволяет масштабировать его горизонтально).
  • Backoff реализован на нативном x-message-ttl RabbitMQ (без плагина delayed-message), отдельная очередь на каждый тир исключает HOL-blocking.
  • Тесты и CI не входили в объём задания, но код спроектирован тестируемым: бизнес-логика зависит от абстракций repositories / UnitOfWork (Protocol), поэтому в юнит-тестах доступ к данным подменяется фейковым репозиторием без БД; брокер/сессия инъектируются.

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages