Загрузка платежей с проверкой блокировки записи в БД

2. Загрузка платежей с проверкой блокировки записи в БД

Условие задачи:
Дан API-эндпоинт на FastAPI для загрузки файла со списком платежей. Файл содержит до 1000 записей, каждая из которых включает ИНН получателя, номер счета, сумму и назначение платежа.
Нужно доработать обработку загрузки так, чтобы при сохранении платежей дополнительно выполнялась проверка: не заблокирована ли соответствующая запись в базе данных.
Если запись заблокирована, такая операция должна корректно обрабатываться и учитываться при формировании результата обработки файла.
Предложи дополнительные способы оптимизации

# ---
# Дано: загрузка списка платежей (до 1000 в файле): ИНН, реквизиты
# Клиент загружает списки платежей, исходим из того что мы знаем кто это делает. Момент с авторизацией и идентификацией клиента специально опущен

from fastapi import FastAPI, UploadFile, HTTPException
from pydantic import BaseModel
from typing import List

app = FastAPI()


class PaymentData(BaseModel):
    inn: str  # ИНН получателя
    account_number: str  # Номер (ID) счета
    amount: float  # Сумма платежа
    purpose: str  # Назначение платежа


class PaymentResponse(BaseModel):
    status: str
    message: str
    processed_count: int
    errors_count: int


def parse_file(file):
    ...


def save_payments(data: list[PaymentData]):
    ...


@app.post("/api/payments/upload", response_model=PaymentResponse)
async def upload_payments(customer_code, file: UploadFile):
    try:
        payments_data: list[PaymentData] = parse_file(file)
        processed_count, errors_count = save_payments(payments_data)
    except Exception as e:
        raise HTTPException(status_code=500, detail=f"Ошибка при обработке файла: {str(e)}")

    return PaymentResponse(
        status="success",
        message="Файл успешно обработан",
        processed_count=processed_count,
        errors_count=errors_count
    )
Спойлеры к решению
Подсказки
  • Текущий PaymentResponse слишком бедный: он не показывает, какие строки файла не обработались и почему.
  • Сумму платежа лучше хранить как Decimal, а не float.
  • UploadFile в FastAPI читается асинхронно, поэтому parse_file(file) лучше сделать async.
  • Проверку блокировки лучше делать на уровне БД.
  • Если блокировка бизнесовая, например поле is_blocked, нужно проверять это поле.
  • Если речь именно про database lock, можно использовать SELECT ... FOR UPDATE NOWAIT или SELECT ... FOR UPDATE SKIP LOCKED.
  • При обработке файла до 1000 записей не стоит делать по одному запросу на каждую проверку, лучше группировать данные и делать batch-запросы.
Решение

1. Что не так в текущем варианте #

Исходная схема:

payments_data: list[PaymentData] = parse_file(file)
processed_count, errors_count = save_payments(payments_data)

Проблемы:

1. parse_file(file) синхронный, хотя UploadFile обычно читается асинхронно.
2. save_payments(data) ничего не говорит о причинах ошибок.
3. amount имеет тип float, что плохо для денег.
4. Нет построчного результата обработки.
5. Нет проверки существования счёта.
6. Нет проверки блокировки записи.
7. Нет защиты от дублей в файле.
8. Нет транзакционной логики.
9. Нет разделения успешных и неуспешных строк.
10. Любая ошибка превращается в HTTP 500, хотя часть ошибок может быть нормальным результатом обработки файла.

Для денег лучше использовать:

Decimal

а не:

float

Потому что float может давать ошибки округления.

2. Улучшенные модели ответа #

Лучше возвращать не только количество ошибок, но и список ошибок по строкам файла.

from decimal import Decimal
from pydantic import BaseModel, Field


class PaymentData(BaseModel):
    row_number: int
    inn: str
    account_number: str
    amount: Decimal = Field(gt=0)
    purpose: str


class PaymentRowError(BaseModel):
    row_number: int
    inn: str | None = None
    account_number: str | None = None
    code: str
    message: str


class PaymentResponse(BaseModel):
    status: str
    message: str
    processed_count: int
    errors_count: int
    errors: list[PaymentRowError] = []

Пример ошибки по строке:

{
  "row_number": 15,
  "inn": "7700000000",
  "account_number": "40702810000000000001",
  "code": "ACCOUNT_LOCKED",
  "message": "Счёт заблокирован"
}

Так клиент сможет понять, какие именно записи не были сохранены.

3. Вариант 1 — если блокировка бизнесовая #

Например, в таблице счетов есть поле:

is_blocked BOOLEAN NOT NULL DEFAULT FALSE

Тогда перед сохранением платежа нужно проверить, что счёт существует и не заблокирован.

Пример таблиц:

CREATE TABLE accounts (
    id BIGSERIAL PRIMARY KEY,
    customer_code TEXT NOT NULL,
    inn TEXT NOT NULL,
    account_number TEXT NOT NULL,
    is_blocked BOOLEAN NOT NULL DEFAULT FALSE,

    UNIQUE (customer_code, inn, account_number)
);

CREATE TABLE payments (
    id BIGSERIAL PRIMARY KEY,
    account_id BIGINT NOT NULL REFERENCES accounts(id),
    amount NUMERIC(18, 2) NOT NULL,
    purpose TEXT NOT NULL,
    status TEXT NOT NULL,
    created_at TIMESTAMP NOT NULL DEFAULT now()
);

Пример обработки через asyncpg:

from decimal import Decimal

import asyncpg
from fastapi import FastAPI, UploadFile, HTTPException
from pydantic import BaseModel, Field


app = FastAPI()


class PaymentData(BaseModel):
    row_number: int
    inn: str
    account_number: str
    amount: Decimal = Field(gt=0)
    purpose: str


class PaymentRowError(BaseModel):
    row_number: int
    inn: str | None = None
    account_number: str | None = None
    code: str
    message: str


class PaymentResponse(BaseModel):
    status: str
    message: str
    processed_count: int
    errors_count: int
    errors: list[PaymentRowError] = []


async def parse_file(file: UploadFile) -> list[PaymentData]:
    content = await file.read()

    # Здесь должна быть реальная логика парсинга CSV/XLSX/JSON.
    # Важно: row_number лучше проставлять сразу при парсинге.
    raise NotImplementedError


async def save_payments(
    conn: asyncpg.Connection,
    customer_code: str,
    payments: list[PaymentData],
) -> tuple[int, list[PaymentRowError]]:
    processed_count = 0
    errors: list[PaymentRowError] = []

    async with conn.transaction():
        for payment in payments:
            account = await conn.fetchrow(
                """
                SELECT id, is_blocked
                FROM accounts
                WHERE customer_code = $1
                  AND inn = $2
                  AND account_number = $3
                """,
                customer_code,
                payment.inn,
                payment.account_number,
            )

            if account is None:
                errors.append(
                    PaymentRowError(
                        row_number=payment.row_number,
                        inn=payment.inn,
                        account_number=payment.account_number,
                        code="ACCOUNT_NOT_FOUND",
                        message="Счёт не найден",
                    )
                )
                continue

            if account["is_blocked"]:
                errors.append(
                    PaymentRowError(
                        row_number=payment.row_number,
                        inn=payment.inn,
                        account_number=payment.account_number,
                        code="ACCOUNT_LOCKED",
                        message="Счёт заблокирован",
                    )
                )
                continue

            await conn.execute(
                """
                INSERT INTO payments (
                    account_id,
                    amount,
                    purpose,
                    status
                )
                VALUES ($1, $2, $3, $4)
                """,
                account["id"],
                payment.amount,
                payment.purpose,
                "created",
            )

            processed_count += 1

    return processed_count, errors

Endpoint:

@app.post("/api/payments/upload", response_model=PaymentResponse)
async def upload_payments(customer_code: str, file: UploadFile):
    try:
        payments_data = await parse_file(file)

        async with app.state.pool.acquire() as conn:
            processed_count, errors = await save_payments(
                conn=conn,
                customer_code=customer_code,
                payments=payments_data,
            )

    except HTTPException:
        raise
    except Exception as e:
        raise HTTPException(
            status_code=500,
            detail=f"Ошибка при обработке файла: {str(e)}",
        )

    return PaymentResponse(
        status="success" if not errors else "partial_success",
        message="Файл обработан",
        processed_count=processed_count,
        errors_count=len(errors),
        errors=errors,
    )

4. Вариант 2 — если речь именно про DB-lock #

Если нужно проверить, что строка не заблокирована другой транзакцией, можно использовать:

SELECT ...
FOR UPDATE NOWAIT

NOWAIT означает: если строка уже заблокирована другой транзакцией, не ждать, а сразу выбросить ошибку.

Пример:

from asyncpg.exceptions import LockNotAvailableError


async def save_payments_with_db_lock_check(
    conn: asyncpg.Connection,
    customer_code: str,
    payments: list[PaymentData],
) -> tuple[int, list[PaymentRowError]]:
    processed_count = 0
    errors: list[PaymentRowError] = []

    async with conn.transaction():
        for payment in payments:
            try:
                account = await conn.fetchrow(
                    """
                    SELECT id
                    FROM accounts
                    WHERE customer_code = $1
                      AND inn = $2
                      AND account_number = $3
                    FOR UPDATE NOWAIT
                    """,
                    customer_code,
                    payment.inn,
                    payment.account_number,
                )

            except LockNotAvailableError:
                errors.append(
                    PaymentRowError(
                        row_number=payment.row_number,
                        inn=payment.inn,
                        account_number=payment.account_number,
                        code="ROW_LOCKED",
                        message="Запись сейчас заблокирована другой операцией",
                    )
                )
                continue

            if account is None:
                errors.append(
                    PaymentRowError(
                        row_number=payment.row_number,
                        inn=payment.inn,
                        account_number=payment.account_number,
                        code="ACCOUNT_NOT_FOUND",
                        message="Счёт не найден",
                    )
                )
                continue

            await conn.execute(
                """
                INSERT INTO payments (
                    account_id,
                    amount,
                    purpose,
                    status
                )
                VALUES ($1, $2, $3, $4)
                """,
                account["id"],
                payment.amount,
                payment.purpose,
                "created",
            )

            processed_count += 1

    return processed_count, errors

В этом варианте заблокированная строка не ломает обработку всего файла. Она просто попадает в результат как ошибка конкретной строки.

Например:

{
  "status": "partial_success",
  "message": "Файл обработан",
  "processed_count": 998,
  "errors_count": 2,
  "errors": [
    {
      "row_number": 12,
      "inn": "7700000000",
      "account_number": "40702810000000000001",
      "code": "ROW_LOCKED",
      "message": "Запись сейчас заблокирована другой операцией"
    },
    {
      "row_number": 45,
      "inn": "7800000000",
      "account_number": "40702810000000000002",
      "code": "ACCOUNT_NOT_FOUND",
      "message": "Счёт не найден"
    }
  ]
}

5. Более оптимальный batch-подход #

Для файла до 1000 строк можно не ходить в БД по каждой строке отдельно. Лучше сначала собрать все уникальные пары:

inn + account_number

Потом одним запросом получить все подходящие счета.

Идея:

keys = {
    (payment.inn, payment.account_number)
    for payment in payments
}

Дальше одним SQL-запросом получить счета:

SELECT id, inn, account_number, is_blocked
FROM accounts
WHERE customer_code = $1
  AND (inn, account_number) = ANY($2)

Но с asyncpg удобнее передать данные через unnest.

Пример:

async def get_accounts_map(
    conn: asyncpg.Connection,
    customer_code: str,
    payments: list[PaymentData],
) -> dict[tuple[str, str], asyncpg.Record]:
    inns = [payment.inn for payment in payments]
    account_numbers = [payment.account_number for payment in payments]

    rows = await conn.fetch(
        """
        WITH input_data AS (
            SELECT *
            FROM unnest($2::text[], $3::text[]) AS t(inn, account_number)
        )
        SELECT DISTINCT
            a.id,
            a.inn,
            a.account_number,
            a.is_blocked
        FROM accounts a
        JOIN input_data i
          ON i.inn = a.inn
         AND i.account_number = a.account_number
        WHERE a.customer_code = $1
        """,
        customer_code,
        inns,
        account_numbers,
    )

    return {
        (row["inn"], row["account_number"]): row
        for row in rows
    }

После этого можно обрабатывать список платежей уже без отдельного SELECT на каждую строку:

async def save_payments_batch(
    conn: asyncpg.Connection,
    customer_code: str,
    payments: list[PaymentData],
) -> tuple[int, list[PaymentRowError]]:
    processed_count = 0
    errors: list[PaymentRowError] = []
    rows_to_insert = []

    async with conn.transaction():
        accounts_map = await get_accounts_map(
            conn=conn,
            customer_code=customer_code,
            payments=payments,
        )

        for payment in payments:
            key = (payment.inn, payment.account_number)
            account = accounts_map.get(key)

            if account is None:
                errors.append(
                    PaymentRowError(
                        row_number=payment.row_number,
                        inn=payment.inn,
                        account_number=payment.account_number,
                        code="ACCOUNT_NOT_FOUND",
                        message="Счёт не найден",
                    )
                )
                continue

            if account["is_blocked"]:
                errors.append(
                    PaymentRowError(
                        row_number=payment.row_number,
                        inn=payment.inn,
                        account_number=payment.account_number,
                        code="ACCOUNT_LOCKED",
                        message="Счёт заблокирован",
                    )
                )
                continue

            rows_to_insert.append(
                (
                    account["id"],
                    payment.amount,
                    payment.purpose,
                    "created",
                )
            )

        if rows_to_insert:
            await conn.executemany(
                """
                INSERT INTO payments (
                    account_id,
                    amount,
                    purpose,
                    status
                )
                VALUES ($1, $2, $3, $4)
                """,
                rows_to_insert,
            )

        processed_count = len(rows_to_insert)

    return processed_count, errors

Такой подход лучше, чем делать 1000 отдельных SELECT.

6. Если нужно batch-поведение с DB-lock #

Если задача именно про строки, заблокированные другой транзакцией, можно использовать FOR UPDATE SKIP LOCKED.

Идея:

SELECT ...
FROM accounts
WHERE ...
FOR UPDATE SKIP LOCKED

Тогда БД вернёт только те строки, которые удалось заблокировать. Заблокированные другой транзакцией строки будут пропущены.

После этого логика такая:

1. Собрали все счета из файла.
2. Попытались получить их через FOR UPDATE SKIP LOCKED.
3. То, что вернулось, можно обрабатывать.
4. То, что не вернулось, но существует в БД, считаем временно заблокированным.
5. То, чего нет в БД, считаем ACCOUNT_NOT_FOUND.

Но тут важно отличить:

не найдено в БД

от:

есть в БД, но заблокировано другой транзакцией

Для этого можно сделать два запроса:

1. Обычный SELECT — узнать, какие счета вообще существуют.
2. SELECT ... FOR UPDATE SKIP LOCKED — узнать, какие удалось заблокировать.

7. Что лучше вернуть клиенту #

Не стоит возвращать просто:

{
  "status": "success",
  "processed_count": 900,
  "errors_count": 100
}

Лучше вернуть детальную информацию:

{
  "status": "partial_success",
  "message": "Файл обработан частично",
  "processed_count": 900,
  "errors_count": 100,
  "errors": [
    {
      "row_number": 10,
      "inn": "7700000000",
      "account_number": "40702810000000000001",
      "code": "ACCOUNT_LOCKED",
      "message": "Счёт заблокирован"
    }
  ]
}

Статусы можно сделать такими:

success         — все строки обработаны
partial_success — часть строк обработана, часть отклонена
failed          — файл вообще не удалось обработать

8. Дополнительные способы оптимизации #

1. Читать и парсить файл потоково, а не загружать огромный файл целиком.
2. Ограничить размер файла на уровне API.
3. Проверять content-type и расширение файла.
4. Валидировать строки файла до похода в БД.
5. Удалять дубли внутри файла до сохранения.
6. Делать batch-проверку счетов.
7. Делать batch-insert платежей.
8. Использовать пул соединений к БД.
9. Добавить индексы на customer_code, inn, account_number.
10. Не держать транзакцию открытой во время парсинга файла.
11. Использовать Decimal для суммы.
12. Для долгой обработки файла можно сохранять upload job и обрабатывать файл в фоне через очередь.
13. Добавить idempotency-key для повторной загрузки файла.
14. Хранить hash файла, чтобы не обработать один и тот же файл дважды.
15. Добавить таблицу payment_imports для истории загрузок.

Индексы:

CREATE UNIQUE INDEX idx_accounts_customer_inn_account
ON accounts (customer_code, inn, account_number);

CREATE INDEX idx_payments_account_id
ON payments (account_id);

Для истории загрузок:

CREATE TABLE payment_imports (
    id BIGSERIAL PRIMARY KEY,
    customer_code TEXT NOT NULL,
    file_name TEXT NOT NULL,
    file_hash TEXT NOT NULL,
    status TEXT NOT NULL,
    processed_count INTEGER NOT NULL DEFAULT 0,
    errors_count INTEGER NOT NULL DEFAULT 0,
    created_at TIMESTAMP NOT NULL DEFAULT now(),

    UNIQUE (customer_code, file_hash)
);

Это позволит защититься от повторной обработки одного и того же файла.

9. Итоговая архитектурная идея #

1. Endpoint принимает файл.
2. Файл парсится и валидируется.
3. Создаётся запись об импорте файла.
4. Счета из файла проверяются batch-запросом.
5. Заблокированные записи не сохраняются как платежи.
6. По каждой проблемной строке формируется ошибка.
7. Валидные платежи сохраняются пачкой.
8. Клиент получает summary + список ошибок.

Главное: заблокированная запись не должна приводить к падению всей загрузки. Это нормальный частичный результат обработки файла, который должен попасть в errors с понятным кодом, например ACCOUNT_LOCKED или ROW_LOCKED.