Перейти к основному контенту
Tech Path Finder
КурсыИнтервьюКод-ревьюБлог
Tech Path Finder

Персонализированный путеводитель в IT. Квизы, мок-интервью, код ревью и аналитика прогресса.

@potapov_me

Платформа

  • Курсы
  • Прогресс
  • Мок-интервью
  • Код ревью
  • Живое ревью с ИИ
  • Тренажёр переговоров
  • Закладки

Контент

  • Блог
  • Главная
  • Обратная связь

Компания

  • О проекте
  • Тарифы
  • Условия использования
  • Конфиденциальность
  • Согласие на обработку данных
  • Cookie
  • Реквизиты

Аккаунт

  • Войти
  • Зарегистрироваться
  • Профиль

© 2026 Tech Path Finder. Все права защищены.

·ИП Потапов К.С.·Политика конфиденциальности·
Сделано с ❤️ в России
  1. Основы работы с NATS
nats_basics

Основы работы с NATS

Высокопроизводительная коммуникация: subjects, queue groups, wildcards. JetStream для персистентности.

Основы работы с NATS

Высокопроизводительная коммуникация: subjects, queue groups, wildcards. JetStream для персистентности.

#Что такое NATS?

NATS — это высокопроизводительный брокер сообщений с простой моделью subjects, разработанный для:

  • Микросервисной коммуникации
  • Service mesh и service discovery
  • IoT и edge computing
  • Real-time приложений

Ключевые особенности:

  • Производительность: ~500K msg/s (быстрее RabbitMQ и Redis)
  • Простота: модель subjects без сложных exchanges
  • Wildcards: подписка по паттернам (*, >)
  • Queue groups: балансировка нагрузки между потребителями
  • JetStream: персистентность и streaming (опционально)

#NATS Core vs JetStream

ХарактеристикаNATS CoreJetStream
ПерсистентностьНет (at-most-once)Да (at-least-once)
ХранениеСообщения не хранятсяЛог сообщений на диске
DeliveryFire-and-forgetГарантированная доставка
ReplayНетДа (чтение с любого offset)
СценарийReal-time, service meshEvent sourcing, очереди

#Архитектура NATS

┌─────────────┐
│  Publisher  │
└──────┬──────┘
       │ PUBLISH orders.new {...}
       ▼
┌─────────────────┐
│     NATS        │
│    Server       │
└──────┬──────────┘
       │
       ├────────────────────┐
       │                    │
       ▼                    ▼
┌─────────────┐      ┌─────────────┐
│ Subscriber  │      │ Subscriber  │
│ Queue Group │      │ Queue Group │
│  (workers)  │      │ (analytics) │
└─────────────┘      └─────────────┘

#Ключевые концепции

#Subject (Субъект)

Тема сообщения в формате dot-separated:

orders.new
orders.paid
orders.shipped
user.created
user.deleted

Иерархия:

  • orders.new — новый заказ
  • orders.paid — оплаченный заказ
  • orders.* — все события заказов (wildcard)

#Wildcards (Подстановочные знаки)

  • * — ровно один токен
  • > — ноль или более токенов (должен быть в конце)
# Подписка на все события заказов
orders.*       → orders.new, orders.paid, orders.shipped
orders.>       → orders.new, orders.paid, orders.shipped, orders.shipped.international

# Подписка на все события
>              → все сообщения

# Подписка на все user-события
user.*         → user.created, user.deleted
user.>         → user.created, user.deleted, user.profile.updated

#Queue Group (Группа очереди)

Балансировка нагрузки между потребителями:

Queue Group: "order-workers"

┌─────────────┐     ┌─────────────┐     ┌─────────────┐
│  Worker 1   │     │  Worker 2   │     │  Worker 3   │
│  orders.new │     │  orders.new │     │  orders.new │
└─────────────┘     └─────────────┘     └─────────────┘
       ▲                   ▲                   ▲
       └───────────────────┼───────────────────┘
                           │
                    NATS Server
              (round-robin балансировка)

Правила:

  • Сообщение получает только один потребитель из queue group
  • Между разными queue groups сообщение доставляется всем
  • Аналог consumer groups в Kafka, но проще

#Работа с NATS в FastStream

#Конфигурация

from faststream import FastStream from faststream.nats import NatsBroker broker = NatsBroker("nats://localhost:4222") app = FastStream(broker)

#Publisher

@broker.publisher("orders.new") async def create_order(order: dict): await broker.publish(order, "orders.new")

#Subscriber

from pydantic import BaseModel class Order(BaseModel): id: int amount: float @broker.subscriber("orders.new") async def handle_order(order: Order): print(f"New order: {order.id}")

#Queue Group

@broker.subscriber( "orders.new", queue="order-workers" # Queue group для балансировки ) async def handle_order(order: Order): # Только один воркер из группы получит сообщение ...

#Wildcard Subscription

# Подписка на все события заказов @broker.subscriber("orders.*") async def handle_order_event(event: dict): ... # Подписка на все вложенные события @broker.subscriber("orders.>") async def handle_all_orders(event: dict): ...

#Request-Reply в NATS

NATS имеет встроенную поддержку RPC:

# Server @broker.subscriber("calculate") async def handle_calc(data: dict) -> dict: return {"result": data["a"] + data["b"]} # Client response = await broker.request( {"a": 2, "b": 3}, "calculate", timeout=5.0 ) print(response) # {"result": 5}

Преимущества перед ручным RPC:

  • Автоматический reply subject
  • Встроенный таймаут
  • Балансировка между серверами (queue group)

#JetStream: персистентность

JetStream добавляет персистентность к NATS:

#Конфигурация JetStream

from faststream.nats import NatsBroker broker = NatsBroker( "nats://localhost:4222", jetstream=True # Включить JetStream )

#Stream (Поток)

Stream — это лог хранения сообщений:

from faststream.nats import JetStreamConfig # Создание stream await broker.add_stream( name="orders", subjects=["orders.*"], # Subjects для stream storage="file" # file или memory )

#Consumer

@broker.subscriber( "orders.new", stream="orders", # JetStream stream durable="order-worker", # Имя потребителя (персистентное) ack=True, # Подтверждение обработки manual_ack=True # Ручной ack после обработки ) async def handle_order(order: Order, message): try: await db.save(order) await message.ack() # Подтвердить except Exception: await message.nack() # Вернуть в stream

#Deliver Policy

Политика доставки для новых потребителей:

@broker.subscriber( "orders.new", stream="orders", deliver_policy="all" # all, last, new, start_time, start_seq ) async def handle_order(order: Order): ...

Опции:

  • all — читать с начала stream
  • last — только последнее сообщение
  • new — только новые сообщения (после подписки)
  • start_time — с указанного времени
  • start_seq — с указанной последовательности

#Продвинутые возможности

#Pull Consumer

Потребитель сам запрашивает сообщения (вместо push):

@broker.subscriber( "orders.new", stream="orders", pull=True, # Pull consumer batch_size=10, # Запрашивать по 10 сообщений batch_timeout_ms=1000 ) async def handle_orders(orders: list[Order]): # Обработать batch ...

Преимущества:

  • Контроль потока (backpressure)
  • Нет риска перегрузки потребителя
  • Подходит для batch processing

#KV Bucket (Key-Value)

NATS JetStream KV — распределённое key-value хранилище:

# Создание KV await broker.create_kv_bucket( name="config", history=10 # Хранить 10 версий ) # Запись await broker.put("config", "app.version", "1.2.3") # Чтение value = await broker.get("config", "app.version") # Подписка на изменения @broker.subscriber("config.>", stream="KV_config") async def handle_config_change(key: str, value: str): ...

Сценарии:

  • Конфигурация приложений
  • Service discovery
  • Feature flags

#Object Store

Хранение бинарных объектов:

# Создание Object Store await broker.create_object_store( name="uploads", storage="file" ) # Загрузка await broker.put_object( "uploads", "file.pdf", data=binary_data ) # Скачивание data = await broker.get_object("uploads", "file.pdf")

#Пример: Service Mesh

NATS идеален для service mesh благодаря производительности и простоте:

from faststream import FastStream from faststream.nats import NatsBroker from pydantic import BaseModel class ServiceRequest(BaseModel): service: str method: str payload: dict class ServiceResponse(BaseModel): status: int data: dict broker = NatsBroker("nats://localhost:4222") app = FastStream(broker) # Service Discovery @broker.subscriber("service.discover") async def discover_service(request: ServiceRequest) -> ServiceResponse: # Вернуть информацию о сервисе return ServiceResponse(status=200, data={"host": "user-service:8080"}) # RPC вызов сервиса @broker.subscriber("service.call", queue="rpc-workers") async def call_service(request: ServiceRequest) -> ServiceResponse: # Проксировать вызов к целевому сервису response = await http.post(f"{request.service}/{request.method}", request.payload) return ServiceResponse(status=response.status, data=response.json()) # События сервисов @broker.subscriber("events.>") async def handle_event(event: dict): # Логирование всех событий print(f"Event: {event}")

#Мониторинг NATS

#NATS CLI

# Информация о сервере nats server info # Статистика соединений nats connection ls # Мониторинг subjects nats sub ">" --count

#Monitoring Port

NATS экспортирует метрики на порту 8222:

curl http://localhost:8222/connz # Соединения curl http://localhost:8222/subsz # Подписки curl http://localhost:8222/routez # Маршруты (в кластере) curl http://localhost:8222/jz # JetStream статус

#Prometheus метрики

NATS экспортирует метрики для Prometheus:

  • nats_connection_count — активные соединения
  • nats_sent_msgs — отправленные сообщения
  • nats_received_msgs — полученные сообщения
  • nats_jetstream_storage_used — использованное место JetStream

#Сравнение брокеров

ХарактеристикаNATS CoreNATS JetStreamRabbitMQKafka
Производительность~500K msg/s~100K msg/s~10K msg/s~1M msg/s
ПерсистентностьНетДаДаДа
МодельSubjectsStreamsExchanges/QueuesTopics/Partitions
СложностьНизкаяСредняяСредняяВысокая
Идеально дляService mesh, real-timeEvent sourcing, очередиTask queues, RPCEvent streaming

Далее: Публикация сообщений (Publishing)