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

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

@potapov_me

Платформа

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

Контент

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

Компания

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

Аккаунт

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

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

·ИП Потапов К.С.·Политика конфиденциальности·
Сделано с ❤️ в России
  1. Event-driven архитектура
event_driven

Event-driven архитектура

Kafka, RabbitMQ, событийная архитектура, паттерны

Открыть лабораториюv1.0Запускается локально из публичного репозитория

Event-driven архитектура с FastAPI

Событийная архитектура для масштабируемых приложений.

#Что такое event-driven?

Event-driven архитектура — сервисы общаются через события (events). Производитель публикует событие, потребители реагируют.

Order Service → [OrderCreated] → Kafka → Payment Service
                                      → Email Service
                                      → Notification Service

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

  • Слабая связанность — сервисы не знают друг о друге
  • Масштабируемость — каждый сервис масштабируется независимо
  • Отказоустойчивость — события сохраняются в очереди

#Apache Kafka

#Установка

docker run -d --name kafka -p 9092:9092 confluentinc/cp-kafka

#Producer

from aiokafka import AIOKafkaProducer import json producer = AIOKafkaProducer( bootstrap_servers='localhost:9092', value_serializer=lambda v: json.dumps(v).encode() ) await producer.start() await producer.send( 'orders', {'order_id': 123, 'user_id': 1, 'total': 99.99} ) await producer.stop()

#Consumer

from aiokafka import AIOKafkaConsumer import json consumer = AIOKafkaConsumer( 'orders', bootstrap_servers='localhost:9092', value_deserializer=lambda m: json.loads(m.decode()), group_id='payment-service' ) await consumer.start() async for message in consumer: order = message.value await process_order(order) await consumer.stop()

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

#Публикация событий

from fastapi import FastAPI from aiokafka import AIOKafkaProducer from contextlib import asynccontextmanager producer = AIOKafkaProducer(bootstrap_servers='kafka:9092') @asynccontextmanager async def lifespan(app): await producer.start() yield await producer.stop() app = FastAPI(lifespan=lifespan) @app.post('/orders') async def create_order(order: OrderCreate): # Создание заказа order = await orders_service.create(order) # Публикация события await producer.send('orders', { 'event': 'OrderCreated', 'order_id': order.id, 'user_id': order.user_id }) return order

#Подписка на события

from fastapi import FastAPI from aiokafka import AIOKafkaConsumer from contextlib import asynccontextmanager import asyncio consumer = AIOKafkaConsumer( 'orders', bootstrap_servers='kafka:9092', group_id='email-service' ) @asynccontextmanager async def lifespan(app): await consumer.start() task = asyncio.create_task(consume_events()) yield await consumer.stop() task.cancel() app = FastAPI(lifespan=lifespan) async def consume_events(): async for message in consumer: event = message.value if event['event'] == 'OrderCreated': await send_order_email(event)

#RabbitMQ

#Установка

docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:management

#Producer (aio-pika)

import aio_pika async def publish_order(order: dict): connection = await aio_pika.connect_robust("amqp://guest:guest@rabbitmq/") async with connection: channel = await connection.channel() exchange = await channel.declare_exchange('orders', aio_pika.ExchangeType.DIRECT) queue = await channel.declare_queue('orders_queue') await exchange.publish( aio_pika.Message(body=json.dumps(order).encode()), routing_key='order.created' )

#Consumer

async def consume(): connection = await aio_pika.connect_robust("amqp://guest:guest@rabbitmq/") async with connection: channel = await connection.channel() exchange = await channel.declare_exchange('orders', aio_pika.ExchangeType.DIRECT) queue = await channel.declare_queue('orders_queue') await queue.bind(exchange, routing_key='order.created') async with queue.iterator() as queue_iter: async for message in queue_iter: async with message.process(): order = json.loads(message.body) await process_order(order)

#Event Sourcing

from datetime import datetime, timezone class EventStore: def __init__(self): self.events = [] def append(self, event: dict): self.events.append({ 'id': len(self.events), 'timestamp': datetime.now(timezone.utc), **event }) def get_events(self, aggregate_id: str): return [e for e in self.events if e.get('aggregate_id') == aggregate_id] # Восстановление состояния из событий def rebuild_state(events: list): state = {} for event in events: if event['type'] == 'OrderCreated': state = {'status': 'created', **event['data']} elif event['type'] == 'OrderPaid': state['status'] = 'paid' elif event['type'] == 'OrderShipped': state['status'] = 'shipped' return state

#CQRS (Command Query Responsibility Segregation)

# Write model (commands) class OrderCommand: def create(self, data: dict): # Валидация, бизнес-логика event = {'type': 'OrderCreated', 'data': data} event_store.append(event) # Read model (queries) class OrderQuery: def get_order(self, order_id: str): # Оптимизированный read model return read_db.query(f"SELECT * FROM orders WHERE id = {order_id}")

#Паттерны

#Event Carried State Transfer

# Событие содержит все необходимые данные { 'event': 'OrderCreated', 'order_id': 123, 'user_id': 1, 'total': 99.99, 'items': [...], 'shipping_address': {...} }

#Dead Letter Queue

consumer = AIOKafkaConsumer( 'orders', bootstrap_servers='kafka:9092', group_id='payment-service' ) async for message in consumer: try: await process_order(message.value) await consumer.commit() except Exception as e: # Отправка в DLQ await producer.send('orders-dlq', { 'error': str(e), 'original': message.value })

#Исполняемая лаборатория

Лаборатория «EventBus, idempotency key, key-based partitioning» доступна в публичном GitLab. Она проверяет in-memory EventBus с pub/sub, идемпотентность через ключ, key-based partitioning для консьюмеров и корректную обработку повторных событий.

Склонируйте репозиторий и запустите тесты из его корня:

git clone --branch v1.0 --depth 1 https://gitlab.potapov.me/courses/fastapi_pro.git cd fastapi_pro uv sync --group test uv run --group test pytest event_driven/tests

#Чеклист event-driven

  • Kafka/RabbitMQ broker
  • Producer для публикации
  • Consumer для подписки
  • Обработка ошибок (DLQ)
  • Идемпотентность потребителей
  • Мониторинг очередей

Далее: GraphQL