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

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

@potapov_me

Платформа

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

Контент

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

Компания

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

Аккаунт

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

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

·ИП Потапов К.С.·Политика конфиденциальности·
Сделано с ❤️ в России
  1. Миграции и версионирование схем
schema_migration

Миграции и версионирование схем

Эволюция форматов сообщений: версионирование, совместимость, миграции без downtime, schema registry.

Миграции и версионирование схем

Эволюция форматов сообщений: версионирование, совместимость, миграции без downtime, schema registry.

#Проблемы эволюции схем

Когда вы изменяете формат сообщений, возникают проблемы:

# Версия 1 class OrderV1(BaseModel): id: int amount: float # Версия 2 (добавили поле) class OrderV2(BaseModel): id: int amount: float discount: float # Новое поле # Что произойдёт? # - Старые потребители не поймут новые сообщения # - Новые потребители не обработают старые сообщения # - Сообщения в очереди могут быть разных версий

#Типы совместимости

#Backward Compatibility (Обратная совместимость)

Новые потребители читают старые сообщения.

# Новое поле optional с default class OrderV2(BaseModel): id: int amount: float discount: float = 0.0 # Default для старых сообщений

Правила:

  • ✅ Добавление optional полей с default
  • ✅ Добавление новых endpoints
  • ❌ Удаление полей
  • ❌ Изменение типа существующего поля
  • ❌ Делать required поле, которое было optional

#Forward Compatibility (Прямая совместимость)

Старые потребители читают новые сообщения.

# Игнорирование неизвестных полей class OrderV1(BaseModel): id: int amount: float class Config: extra = "ignore" # Игнорировать новые поля

Правила:

  • ✅ Добавление новых полей
  • ✅ Игнорирование неизвестных полей (extra="ignore")
  • ❌ Удаление полей, которые используют старые потребители
  • ❌ Изменение семантики существующих полей

#Full Compatibility (Полная совместимость)

И backward, и forward одновременно. Требуется для zero-downtime деплоя.

#Стратегии версионирования

#Версия в имени топика

# Публикация await broker.publish(order_v2, "orders.v2") # Потребление @broker.subscriber("orders.v1") async def handle_v1(order: OrderV1): ... @broker.subscriber("orders.v2") async def handle_v2(order: OrderV2): ...

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

  • Явная версия
  • Легко маршрутизировать
  • Можно деплоить новые потребители отдельно

Недостатки:

  • Множество топиков
  • Нужно управлять lifecycle старых версий

#Версия в сообщении

class BaseMessage(BaseModel): version: str = "1.0" class OrderV1(BaseMessage): version: str = "1.0" id: int amount: float class OrderV2(BaseMessage): version: str = "2.0" id: int amount: float discount: float # Потребитель с маршрутизацией @broker.subscriber("orders") async def handle_order(message: dict): version = message.get("version", "1.0") if version == "1.0": order = OrderV1(**message) elif version == "2.0": order = OrderV2(**message) else: raise ValueError(f"Unknown version: {version}")

#Версия в заголовках

# Публикация await broker.publish( order, "orders", headers={"schema_version": "2.0"} ) # Потребитель @broker.subscriber("orders") async def handle_order(message: dict, msg: RabbitMessage): version = msg.headers.get("schema_version", "1.0") if version == "2.0": order = OrderV2(**message) else: order = OrderV1(**message)

#Миграции без downtime

#Стратегия: Dual Write

async def publish_order(order: dict): # Публикация в обе версии await broker.publish(order, "orders.v1") # Старые потребители await broker.publish({**order, "version": "2.0"}, "orders.v2") # Новые

Этапы:

  1. Deploy: публикация в v1 и v2 одновременно
  2. Миграция потребителей с v1 на v2
  3. Когда все потребители на v2 — отключить публикацию в v1

#Стратегия: Expand and Contract

Фаза 1: Expand (расширение)

# Добавляем новое поле как optional class Order(BaseModel): id: int amount: float discount: float = 0.0 # Новое поле
  • Deploy новых потребителей (понимают v2)
  • Старые потребители работают (игнорируют discount)

Фаза 2: Migration (миграция)

  • Все потребители обновлены до v2
  • Данные мигрированы (discount заполнен)

Фаза 3: Contract (сжатие)

# Удаляем старое поле (если нужно) class OrderV3(BaseModel): id: int discount: float # amount удалён

#Стратегия: Parallel Deployment

Топик: orders

┌─────────────┐
│  Publisher  │
│   (v1+v2)   │
└──────┬──────┘
       │
       ├──> orders.v1 ──> Consumer v1 (старый)
       │
       └──> orders.v2 ──> Consumer v2 (новый)

Этапы:

  1. Создать orders.v2
  2. Deploy Consumer v2
  3. Переключить publisher на v2
  4. Дождаться обработки всех сообщений из v1
  5. Удалить Consumer v1 и orders.v1

#Schema Registry

#Confluent Schema Registry (для Kafka)

from schema_registry.client import SchemaRegistryClient, Schema client = SchemaRegistryClient(url="http://schema-registry:8081") # Регистрация схемы schema = Schema({ "type": "record", "name": "Order", "fields": [ {"name": "id", "type": "int"}, {"name": "amount", "type": "float"}, {"name": "discount", "type": "float", "default": 0.0} ] }, schema_type="AVRO") subject = "orders-value" schema_id = client.register(subject, schema) # Проверка совместимости compatibility = client.test_compatibility(subject, schema) # True если backward compatible

#Проверка совместимости

from schema_registry.client import Compatibility # Установить уровень совместимости client.set_compatibility(subject, Compatibility.BACKWARD) # Перед публикацией проверить if not client.test_compatibility(subject, new_schema): raise ValueError("Schema not compatible!") # Опубликовать с валидацией await broker.publish( order, "orders", headers={"schema_id": schema_id} )

#Пример: полная миграция

from pydantic import BaseModel from typing import Union, Literal # === Версия 1 (старая) === class OrderV1(BaseModel): version: Literal["1.0"] = "1.0" id: int amount: float user_id: int # === Версия 2 (новая) === class OrderV2(BaseModel): version: Literal["2.0"] = "2.0" id: int amount: float discount: float = 0.0 user_id: int currency: str = "USD" # === Union для потребителей === Order = Union[OrderV1, OrderV2] # === Publisher с миграцией === async def publish_order(order_data: dict): # Dual write: публикация в обе версии await broker.publish(order_data, "orders.v1") await broker.publish({**order_data, "version": "2.0"}, "orders.v2") # === Потребитель с маршрутизацией === @broker.subscriber("orders.v1") async def handle_order_v1(order: OrderV1): logger.info(f"V1 order: {order.id}") # Обработка v1 @broker.subscriber("orders.v2") async def handle_order_v2(order: OrderV2): logger.info(f"V2 order: {order.id}, discount: {order.discount}") # Обработка v2 # === Миграция: после обновления всех потребителей === @broker.subscriber("orders") async def handle_order_unified(order: Order): if order.version == "1.0": # Конвертация v1 → v2 order = OrderV2( version="2.0", id=order.id, amount=order.amount, discount=0.0, user_id=order.user_id ) # Единая обработка await process_order(order)

#Best Practices

#1. Всегда добавляйте version поле

class BaseMessage(BaseModel): version: str timestamp: str correlation_id: str

#2. Используйте optional поля с default

# ✅ Хорошо discount: float = 0.0 # ❌ Плохо discount: float # Breaks backward compatibility

#3. Игнорируйте неизвестные поля

class Config: extra = "ignore" # Для forward compatibility

#4. Документируйте изменения

## Changelog ### v2.0 (2026-03-30) - Добавлено поле `discount` (optional, default=0.0) - Добавлено поле `currency` (optional, default="USD") - Backward compatible с v1.0 ### v1.0 (2026-01-01) - Initial schema

#5. Тестируйте совместимость

def test_backward_compatibility(): # Старые данные должны валидироваться новой схемой old_data = {"id": 1, "amount": 100, "user_id": 123} order = OrderV2(**old_data) # Не должно выбросить assert order.discount == 0.0 def test_forward_compatibility(): # Новые данные должны игнорироваться старой схемой new_data = {"id": 1, "amount": 100, "user_id": 123, "discount": 10.0} order = OrderV1(**new_data) # discount игнорируется assert not hasattr(order, 'discount')

#Заключение курса

Поздравляем! Вы прошли полный курс по FastStream:

  1. ✅ Основы messaging и архитектура FastStream
  2. ✅ Работа с 4 брокерами: RabbitMQ, Kafka, Redis, NATS
  3. ✅ Публикация и подписка на сообщения
  4. ✅ Валидация данных с Pydantic
  5. ✅ Dependency Injection
  6. ✅ Тестирование приложений
  7. ✅ Middleware и обработка ошибок
  8. ✅ RPC и Request-Reply паттерны
  9. ✅ Масштабирование и производительность
  10. ✅ Production: мониторинг, логирование, деплой
  11. ✅ Миграции и версионирование схем

Теперь вы готовы строить надёжные, масштабируемые микросервисные приложения с помощью FastStream!