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
Дополнительные зависимости устанавливаются отдельно для каждого брокера.
pip install "faststream[rabbit]"
pip install "faststream[kafka]"
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
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 хорош.
Он не делает брокеры простыми. Он делает работу с ними аккуратнее.