Python
Yesterday

FastStream: пишем асинхронные сервисы на Python без тонны обвязки

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

Именно вокруг такого подхода строятся событийные приложения. А FastStream помогает описывать их на Python без лишней инфраструктурной обвязки.

Что такое FastStream

FastStream — это Python-фреймворк для создания приложений, которые обмениваются сообщениями через брокеры: RabbitMQ, Kafka, NATS, Redis и MQTT.

Проще говоря, он помогает писать consumers и publishers примерно так же, как FastAPI помогает писать HTTP API. Ты объявляешь обычную асинхронную функцию, описываешь входные данные через аннотации типов или Pydantic-модель, а затем связываешь её с очередью или топиком при помощи декоратора:

@broker.subscriber("orders.created")
async def handle_order(event: OrderCreated) -> None:
    ...

FastStream берёт на себя подключение к брокеру, получение и декодирование сообщений, валидацию данных, вызов обработчиков, внедрение зависимостей и корректное завершение работы приложения.

При этом это не попытка спрятать RabbitMQ, Kafka и NATS за одной универсальной кнопкой. У каждого брокера остаются собственные возможности, настройки и модель доставки. FastStream лишь даёт им похожий способ описания приложения и убирает большую часть повторяющегося инфраструктурного кода.

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

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

Зачем он вообще нужен

Представь интернет-магазин.

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

Создан заказ №123.

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

Для передачи таких сообщений используют брокеры — например, RabbitMQ или Kafka. Брокер можно представить как почтовое отделение внутри системы: одно приложение отправляет туда сообщение, а другое забирает и обрабатывает его.

Функцию, которая получает такие сообщения, обычно называют обработчиком или consumer. По сути, это обычный Python-код:

async def handle_order(message):
    await reserve_products(message)
    await send_notification(message)

Но одной функции недостаточно.

Нужно подключиться к брокеру, указать нужную очередь, получить сообщение, превратить JSON в Python-объект, проверить данные, обработать ошибку и сообщить брокеру, что всё прошло успешно. Также нужно корректно закрыть соединение при остановке приложения.

В небольшом примере это выглядит несложно. В настоящем проекте одинаковый технический код быстро начинает повторяться в каждом обработчике.

FastStream берёт эту работу на себя.

Разработчику остаётся описать, откуда приходит сообщение и какие данные ожидает функция:

@broker.subscriber("orders.created")
async def handle_order(event: OrderCreated) -> None:
    await process_order(event)

Здесь orders.created — имя очереди с сообщениями о новых заказах, а OrderCreated — модель ожидаемых данных.

FastStream сам получает сообщение, преобразует его в объект OrderCreated, проверяет поля и вызывает функцию. Если данные имеют неправильный формат, ошибка обнаружится ещё до запуска основной логики.

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

Это особенно полезно, когда обработчиков становится не два или три, а несколько десятков.

Схема работы приложения на FastStream

Устанавливаем FastStream

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

Для RabbitMQ:

pip install "faststream[rabbit]"

Для Kafka через AIOKafka:

pip install "faststream[kafka]"

Для Kafka через Confluent:

pip install "faststream[confluent]"

Аналогично доступны дополнительные зависимости nats, redis и mqtt. Для встроенной командной строки понадобится отдельный extra:

pip install "faststream[cli]"

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

Первый сервис с RabbitMQ

Представим небольшой сервис, который получает событие о создании заказа и отправляет команду сервису уведомлений.

from decimal import Decimal
from uuid import UUID

from faststream import FastStream
from faststream.rabbit import RabbitBroker
from pydantic import BaseModel, PositiveInt


broker = RabbitBroker(
    "amqp://guest:guest@localhost:5672/",
)

app = FastStream(broker)


class OrderItem(BaseModel):
    product_id: UUID
    quantity: PositiveInt
    price: Decimal


class OrderCreated(BaseModel):
    order_id: UUID
    user_id: UUID
    items: list[OrderItem]


class SendNotification(BaseModel):
    user_id: UUID
    text: str


@broker.subscriber("orders.created")
@broker.publisher("notifications.send")
async def handle_order(
    event: OrderCreated,
) -> SendNotification:
    total = sum(
        item.price * item.quantity
        for item in event.items
    )

    return SendNotification(
        user_id=event.user_id,
        text=(
            f"Заказ {event.order_id} создан. "
            f"Сумма: {total}"
        ),
    )

Здесь есть два основных декоратора.

@broker.subscriber("orders.created") регистрирует функцию как обработчик очереди. Когда в очередь приходит сообщение, FastStream декодирует его и пытается собрать объект OrderCreated.

@broker.publisher("notifications.send") публикует возвращённое функцией значение в другую очередь.

Связка из subscriber и publisher позволяет описать простую цепочку обработки почти без транспортного кода. Декораторы работают с аннотациями типов и моделями Pydantic, поэтому входные данные проходят валидацию до выполнения основной логики обработчика.

Запустить приложение можно встроенной командой:

faststream run app:app

Здесь app до двоеточия — имя Python-модуля, а второе app — объект FastStream.

Во время разработки доступна автоматическая перезагрузка:

faststream run app:app --reload

CLI также поддерживает запуск нескольких процессов через --workers, генерацию AsyncAPI и отправку тестовых сообщений в брокер.

Публикация сообщений без декоратора

Не каждая публикация является результатом обработки другого сообщения. Иногда команду нужно отправить из фоновой задачи, startup-хука или обычного метода сервиса.

Для этого используется broker.publish():

from uuid import UUID


async def request_order_recalculation(
    order_id: UUID,
) -> None:
    await broker.publish(
        {
            "order_id": str(order_id),
            "reason": "price_changed",
        },
        queue="orders.recalculate",
    )

Название параметра назначения зависит от транспорта. У RabbitMQ это может быть queue, у Kafka — topic, у NATS — subject, у Redis — channel или параметры stream.

И вот здесь видно важное свойство FastStream: внешне брокеры похожи, но их собственная терминология и возможности не стираются.

Я бы не стал прятать broker.publish() в абстракцию вида UniversalMessageBus, если проект не собирается действительно менять брокер.

Обычно такая абстракция быстро становится либо слишком примитивной, либо начинает повторять API самого FastStream.

Гораздо полезнее вынести публикацию конкретных событий в отдельный класс:

from uuid import UUID


class OrderEventPublisher:
    def __init__(self, message_broker: RabbitBroker) -> None:
        self._broker = message_broker

    async def order_cancelled(
        self,
        order_id: UUID,
        user_id: UUID,
    ) -> None:
        await self._broker.publish(
            {
                "order_id": str(order_id),
                "user_id": str(user_id),
            },
            queue="orders.cancelled",
        )

Так бизнес-код не знает название очереди, но транспорт при этом не маскируется под несуществующий универсальный интерфейс.

Зависимости

В обработчики редко приходит только сообщение. Обычно нужны репозиторий, клиент другого сервиса, настройки или объект текущей транзакции.

У FastStream есть собственная система dependency injection:

from collections.abc import AsyncIterator
from typing import Annotated

from faststream import Depends


class OrderRepository:
    async def save(self, order: OrderCreated) -> None:
        print(f"Saving order {order.order_id}")


async def get_order_repository() -> AsyncIterator[OrderRepository]:
    repository = OrderRepository()

    try:
        yield repository
    finally:
        pass


OrderRepositoryDep = Annotated[
    OrderRepository,
    Depends(get_order_repository),
]


@broker.subscriber("orders.created")
async def save_order(
    event: OrderCreated,
    repository: OrderRepositoryDep,
) -> None:
    await repository.save(event)

Зависимости могут быть вложенными. Их также можно назначать конкретному subscriber, целому router или брокеру, чтобы одна и та же проверка применялась ко всем обработчикам.

Но я бы не переносил в зависимости всю бизнес-логику.

Хорошая зависимость создаёт и освобождает ресурс: сессию базы данных, клиент API, репозиторий или контекст трассировки. Если внутри get_order_service() уже выполняется половина обработки заказа, разобраться в потоке выполнения будет трудно.

Подтверждение сообщений и повторная обработка

Самая неприятная ошибка при работе с брокерами выглядит примерно так:

@broker.subscriber("payments")
async def process_payment(event: PaymentEvent) -> None:
    await charge_card(event)
    await save_payment(event)

Карта уже списана, но база данных временно недоступна. Обработчик завершается ошибкой, сообщение приходит повторно, и карта списывается ещё раз.

FastStream управляет подтверждением сообщений и позволяет выбирать политику через AckPolicy. Например, NACK_ON_ERROR предназначен для повторной доставки при необработанной ошибке:

from faststream import AckPolicy


@broker.subscriber(
    "orders.created",
    ack_policy=AckPolicy.NACK_ON_ERROR,
)
async def process_order(event: OrderCreated) -> None:
    await order_service.process(event)

Для полного ручного управления предусмотрена политика MANUAL.

При этом ack, nack и reject физически работают по-разному в разных брокерах. В RabbitMQ это отдельные протокольные операции. В Kafka подтверждение связано с фиксацией offset, а nack может приводить к возврату позиции consumer. В Redis Streams отрицательное подтверждение не эквивалентно RabbitMQ nack. FastStream унифицирует намерение, но не может отменить различия транспортов.

Повторная доставка — не замена идемпотентности.

Consumer должен быть готов получить одно событие несколько раз. Для этого можно хранить идентификаторы обработанных сообщений, использовать уникальные ограничения в базе и проектировать операции так, чтобы повторный вызов не создавал второй результат.

Для публикации событий после изменения базы пригодится transactional outbox. Иначе легко получить ещё одну классическую проблему: данные в PostgreSQL сохранились, а публикация сообщения не произошла.

FastStream хорошо обрабатывает сообщения. Гарантировать атомарность между произвольной базой данных и внешним брокером он за приложение не может.

Middleware

Когда в каждом обработчике появляются одинаковые try, логирование времени и установка trace ID, пора выносить техническую логику в middleware.

Middleware FastStream может выполнять код до и после обработки или публикации сообщения. Через него удобно реализовать логирование, метрики, трассировку, обработку исключений и собственную стратегию повторных попыток.

Условный middleware измерения времени может выглядеть так:

from time import monotonic

from faststream import BaseMiddleware


class TimingMiddleware(BaseMiddleware):
    async def consume_scope(self, call_next, message):
        started_at = monotonic()

        try:
            return await call_next(message)
        finally:
            duration = monotonic() - started_at
            print(
                "Message processed in "
                f"{duration:.3f} seconds"
            )

Затем он передаётся брокеру:

broker = RabbitBroker(
    "amqp://guest:guest@localhost:5672/",
    middlewares=[TimingMiddleware],
)

Не стоит писать middleware на каждый чих. Проверка существования заказа — бизнес-правило и должна остаться в сервисе. А вот trace ID, метрики и техническая обработка исключений действительно относятся к общей инфраструктуре.

AsyncAPI вместо документации в голове

В HTTP-проекте мы привыкли открывать Swagger и смотреть доступные endpoint. С брокерами часто всё иначе.

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

FastStream генерирует AsyncAPI-схему на основе зарегистрированных subscriber и publisher. Документацию можно экспортировать в JSON или YAML либо запустить как отдельную HTML-страницу.

Для локального просмотра используется команда:

faststream docs serve app:app

По умолчанию страница документации будет доступна на порту 8000. Актуальная версия FastStream также поддерживает интерфейс Try It Out для отправки тестовых сообщений через страницу AsyncAPI.

Документация не заменит нормальное описание семантики события.

Схема покажет, что поле status является строкой. Но она не объяснит, можно ли получить completed после cancelled, считается ли событие фактом или командой и разрешено ли добавлять новые значения без обновления consumer.

Поэтому AsyncAPI стоит использовать вместе с текстовым описанием контракта, а не вместо него.

Тестирование без настоящего RabbitMQ

Одна из самых приятных возможностей FastStream — тестовый брокер.

Он перехватывает публикации и направляет сообщения зарегистрированным обработчикам в памяти. Благодаря этому большинство тестов можно запускать без Docker Compose и настоящей очереди. При необходимости тот же тестовый клиент умеет работать и с реальным брокером через параметр with_real=True.

Допустим, у нас есть обработчик:

@broker.subscriber("orders.created")
async def handle_order(event: OrderCreated) -> None:
    await order_service.process(event)

Тест выглядит так:

from uuid import uuid4

import pytest
from faststream.rabbit import TestRabbitBroker


@pytest.mark.asyncio
async def test_handle_order() -> None:
    message = {
        "order_id": str(uuid4()),
        "user_id": str(uuid4()),
        "items": [
            {
                "product_id": str(uuid4()),
                "quantity": 2,
                "price": "19.90",
            },
        ],
    }

    async with TestRabbitBroker(broker) as test_broker:
        await test_broker.publish(
            message,
            queue="orders.created",
        )

        handle_order.mock.assert_called_once()

В тестовом режиме FastStream добавляет обработчикам mock-объекты, через которые можно проверять количество вызовов и переданные аргументы.

Однако весь проект тестировать только таким способом не нужно.

Сам обработчик лучше оставлять тонким:

@broker.subscriber("orders.created")
async def handle_order(
    event: OrderCreated,
    service: OrderServiceDep,
) -> None:
    await service.process(event)

После этого бизнес-логику OrderService можно тестировать обычными unit-тестами. TestRabbitBroker нужен для проверки связывания очереди, декодирования, валидации, зависимостей и публикации результата.

Один-два интеграционных теста стоит запустить с настоящим RabbitMQ. In-memory режим не проверит сетевые ошибки, настройки exchange, права пользователя, durable-очереди и поведение брокера после перезапуска.

Интеграция FastStream с FastAPI

Теперь к самому интересному сценарию.

Допустим, HTTP API принимает запрос на создание заказа. Сам заказ обрабатывается асинхронно через RabbitMQ. При этом хочется держать HTTP endpoint и consumer в одном приложении.

FastStream предоставляет специальные router-классы для FastAPI:

from faststream.rabbit.fastapi import RabbitRouter

Такой router может одновременно содержать HTTP-маршруты и обработчики сообщений. Он подключается к приложению через обычный app.include_router(). В режиме FastAPI-интеграции обработчики используют систему зависимостей FastAPI: нужно импортировать Depends из fastapi, а не из faststream.

Соберём небольшой пример:

from uuid import UUID, uuid4

from fastapi import Depends, FastAPI, status
from faststream.rabbit.fastapi import RabbitRouter
from pydantic import BaseModel, PositiveInt


router = RabbitRouter(
    "amqp://guest:guest@localhost:5672/",
    schema_url="/asyncapi",
    include_in_schema=True,
)


class CreateOrderRequest(BaseModel):
    user_id: UUID
    product_id: UUID
    quantity: PositiveInt


class OrderAccepted(BaseModel):
    order_id: UUID


class OrderCreated(BaseModel):
    order_id: UUID
    user_id: UUID
    product_id: UUID
    quantity: PositiveInt


class OrderService:
    async def process(
        self,
        event: OrderCreated,
    ) -> None:
        print(f"Processing order {event.order_id}")


def get_order_service() -> OrderService:
    return OrderService()


@router.post(
    "/orders",
    response_model=OrderAccepted,
    status_code=status.HTTP_202_ACCEPTED,
)
async def create_order(
    request: CreateOrderRequest,
) -> OrderAccepted:
    order_id = uuid4()

    event = OrderCreated(
        order_id=order_id,
        user_id=request.user_id,
        product_id=request.product_id,
        quantity=request.quantity,
    )

    await router.broker.publish(
        event.model_dump(mode="json"),
        queue="orders.created",
    )

    return OrderAccepted(order_id=order_id)


@router.subscriber("orders.created")
async def process_order(
    event: OrderCreated,
    service: OrderService = Depends(get_order_service),
) -> None:
    await service.process(event)


app = FastAPI(
    title="Order API",
)

app.include_router(router)

HTTP endpoint возвращает 202 Accepted, потому что запрос принят, но обработка заказа ещё не завершена.

Consumer получает событие из orders.created. В нём можно использовать fastapi.Depends, включая уже существующие зависимости приложения. Сам router управляет подключением брокера в рамках жизненного цикла FastAPI. Для старых версий FastAPI до 0.112.2 документация FastStream отдельно указывает необходимость вручную передать router.lifespan_context; в современных версиях достаточно подключения router.

AsyncAPI будет доступен по пути:

/asyncapi

Путь задаётся через schema_url. Его также можно включить в основную OpenAPI-схему приложения с помощью include_in_schema=True.

FastAPI Depends или FastStream Depends?

Здесь легко ошибиться.

В отдельном FastStream-приложении используется:

from faststream import Depends

В FastAPI-интеграции:

from fastapi import Depends

Когда FastStream работает как plugin для FastAPI, он встраивается в механизм зависимостей FastAPI. Обычный faststream.Context в таком режиме также заменяется специальными аннотациями из модуля интеграции соответствующего брокера.

Я бы не смешивал оба варианта в одном модуле. Иначе по импорту Depends становится невозможно понять, какая система его обрабатывает.

Публикация из HTTP endpoint

Внутри HTTP-маршрута брокер доступен через:

router.broker

Поэтому можно отправить сообщение напрямую:

@router.post("/reports")
async def create_report(request: CreateReportRequest):
    await router.broker.publish(
        request.model_dump(mode="json"),
        queue="reports.generate",
    )

    return {"status": "accepted"}

Официальная интеграция поддерживает такой сценарий и позволяет использовать broker как внутри router, так и через зависимость FastAPI.

В большом проекте я всё же вынес бы публикацию в сервис:

class ReportPublisher:
    def __init__(self, broker: RabbitBroker) -> None:
        self._broker = broker

    async def generate(
        self,
        request: CreateReportRequest,
    ) -> None:
        await self._broker.publish(
            request.model_dump(mode="json"),
            queue="reports.generate",
        )

Endpoint тогда отвечает за HTTP, publisher — за контракт сообщения, а consumer — за выполнение задачи.

Держать HTTP API и consumers вместе или разделить?

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

Объединённое приложение удобно, когда проект небольшой, HTTP и consumer используют одинаковые зависимости, разворачиваются одной командой и имеют похожие требования к ресурсам.

Например, административный сервис может принимать настройки через REST и тут же обрабатывать несколько служебных очередей. Делить его на два deployment только ради архитектурной красоты нет смысла.

Разделение полезнее, когда нагрузка отличается.

HTTP API может требовать десять быстрых экземпляров, а consumer — два процесса с большим объёмом памяти. Или обработчик сообщений периодически падает из-за внешнего сервиса, но HTTP API должен продолжать отвечать. В таком случае лучше создать отдельные точки запуска:

src/
├── api/
│   └── app.py
├── worker/
│   └── app.py
├── application/
├── domain/
└── infrastructure/

При этом бизнес-сервисы, модели событий и репозитории остаются общими. Разделяются только процессы и жизненные циклы.

Микросервисы для этого не обязательны. Один репозиторий, одна кодовая база и два deployment часто дают нужную независимость без сетевого зоопарка между внутренними модулями.

Что FastStream не решает

FastStream делает работу с брокером приятнее, но не проектирует систему вместо разработчика.

Он не решит автоматически:

  • какие события считать публичным контрактом;
  • сколько раз consumer может получить одно сообщение;
  • как публиковать событие атомарно вместе с транзакцией базы;
  • когда сообщение нужно повторить, а когда отправить в dead-letter queue;
  • как менять формат событий без поломки старых consumer;
  • как масштабировать обработчики с учётом partition Kafka;
  • как наблюдать цепочку из пяти асинхронных сервисов.

Здесь по-прежнему нужны идемпотентность, версионирование контрактов, outbox, метрики, tracing и нормальная стратегия обработки ошибок.

Фреймворк убирает рутину. Архитектуру он не отменяет.

Когда FastStream действительно полезен

FastStream особенно хорошо подходит проектам, где уже любят FastAPI и Pydantic.

Порог входа получается небольшим: те же асинхронные функции, аннотации типов, модели запросов и dependency injection. Разработчику не приходится сначала изучать большую внутреннюю платформу, чтобы написать один consumer.

При этом проект не ограничивается игрушечными сценариями. Можно работать с нативными возможностями брокера, вручную управлять подтверждениями, писать middleware, создавать routers и тестировать обработчики как в памяти, так и через реальный транспорт.

Мне больше всего нравится, что код перестаёт быть набором callbacks вокруг клиента RabbitMQ.

Обработчик снова выглядит как функция:

@broker.subscriber("orders.created")
async def handle_order(
    event: OrderCreated,
    service: OrderServiceDep,
) -> None:
    await service.process(event)

Видно, какое событие приходит, какой сервис используется и что происходит дальше.

Остальное тоже никуда не исчезло. Соединение с брокером, декодирование, подтверждения, тестовый транспорт и документация всё ещё существуют. Просто теперь они не размазаны по каждому consumer.

Именно в этом FastStream хорош.

Он не делает брокеры простыми. Он делает работу с ними аккуратнее.