Введение
Почти любая распределенная платформа рано или поздно сталкивается с необходимостью синхронизировать данные между своими системами. И чем выше требования к надежности и скорости доставки изменений, тем меньше остается простых решений. Бизнес-пользователям при работе с интерфейсом одной из платформ нужно знать о сущностях, заведённых во второй. Первая хранит их в реляционной СУБД (неважно, какой именно), а вторая – в документоориентированной MongoDB. До кучи, обе платформы разнесены по сети. Можно, конечно, предложить открыть оба интерфейса бок о бок и сверять данные глазами. Но это справедливо негативно скажется на финансовом обеспечении команды разработки, да и как-то не по-человечески так относиться к своим коллегам. Поэтому приходится зарабатывать доверие в коллективе честным трудом.
Встаёт вопрос: как в условиях распределённой структуры наладить транспортировку данных для взаимодействия независимых систем? Прикручивать ко всем продуктам кастомные интеграции дорого и непродуктивно.
Хочется представить, как мы реализовали интеграцию двух независимых платформ со своими базами данных при помощи инструментов MongoDB для мониторинга данных и брокера сообщений. В данной статье описать путь от самого простого решения на примитивных выгрузках метаданных одной пачкой до реализации на событиях Change Stream посредством брокера сообщений.
Вводные
У нас есть:
Источник данных, которым выступает MongoDB.
Потребитель, внешняя система, которую нужно держать в курсе изменений в источнике.
Представим, что в MongoDB в базе данных "Megafon" есть коллекция "default", содержащая множество документов, представляющих бизнес-объекты. Они все имеют полу-структурированное наполнение: у всех из них есть атрибуты id (идентификатор объекта), name (имя), type (тип) и createdAt (дата создания объекта). Пользователям "потребителя" необходим справочник всех объектов типа "Rule" в источнике данных, который связывает идентификатор объекта id с его именем name.
Не будем подробно останавливаться на потребителе, поскольку это не так важно. Достаточно знать, что он умеет получать и хранить данные из системы-источника.
Рассмотрим все основные функции, которые нужны для решения подобной задачи, с примерами и комментариями, объясняющими причины принятия тех или иных решений.
GET /metadata
Начнём с того, что опишем модель интересующих нас документов. Допустим, в системе-источнике пользователями заводятся объекты, содержащие идентификатор, имя и тип объекта:
from datetime import datetime from typing import Literal from uuid import UUID class Rule: id: UUID createdAt: datetime name: str type: Literal["Rule"]
Саму задачу можно решить одной простой выгрузкой. Достаточно поднять простой микросервис с GET-эндпоинтом, возвращающим содержимое базы данных в желаемом формате. Пусть он возвращает список документов в формате JSON с тремя строковыми атрибутами, содержащими идентификатор, имя и тип объекта:
{ "id": "23571ab2-2b98-4b5f-a3b1-e4ca19d1031c", "name": "Example Name", "type": "Rule" }
К примеру, обработчик GET-запроса можно реализовать так:
from pymongo.asynchronous.collection import AsyncCollection MDB_CREDS = { "username" = "mdbuser", "password" = "mdbuser", } MDB_PROPS = { "host" = "localhost", "port" = 27017, "uuid_representation" = "STANDARD", "db": "Megafon", "collection": "default", } MATCH = { "type": "Rule" } def get_messages(collection: AsyncCollection) -> tuple[dict]: messages = [] await for rule in collection.find(filter=MATCH): messages.append( { "id": str(rule["id"]), # Конвертация `uuid.UUID` к `str`. "name": rule["name"], "type": rule["type"], } ) return tuple(messages)
Так потребитель сможет запрашивать все доступные в источнике данные. Простое начало сложной фичи.

К подобной интеграции возникают справедливые вопросы: зачем гонять такое количество данных, если потребителю нужен один новый объект из тысяч уже выгруженных? Как можно оптимизировать процесс, снизив нагрузку на сеть?
Мне, пожалуйста, фильтрованного
Можно расширить содержимое сообщения, добавив в него дату создания объекта createdAt. Так потребитель сможет найти самый последний полученный документ и передать его идентификатор на вход запросу. Зная дату создания этого объекта, несложно вернуть отфильтрованный срез данных, сэкономив на трафике.
Расширим наш документ, добавив дату создания объекта createdAt:
{ "id": "23571ab2-2b98-4b5f-a3b1-e4ca19d1031c", "name": "Example Name", "createdAt": "2026-07-12T12:00:00.000", "type": "Rule" }
Доработаем обработчик, чтобы он стал выглядеть следующим образом:
def get_messages(collection: AsyncCollection, return_after: UUID | None = None) -> tuple[dict]: messages = [] last_sent = None if return_after: await for rule in collection.find(filter=MATCH | {"id": return_after}): last_sent = rule break _filter = {**MATCH} if last_sent: _filter |= { "id": {"$ne": last_sent["id"]}, "createdAt": {"$gte": last_sent["createdAt"]}, } await for rule in collection.find(filter=_filter): messages.append( { "id": str(rule["id"]), # Конвертация `uuid.UUID` к `str`. "name": rule["name"], "createdAt": rule["createdAt"], "type": rule["type"], } ) return tuple(messages) )
И бизнес сыт, и сеть цела. Но вот незадача: объектам в источнике можно редактировать наименование! Как отразить изменение наименований уже выгруженных сущностей? Фильтрация по дате создания объекта тут ничем не поможет.
Потоки изменений
Для таких задач в MongoDB реализованы потоки изменений (Change Streams) – инструмент по отслеживанию изменений данных в реальном времени. Используя потоки изменений, наше приложение может следить за состоянием базы данных, моментально сообщая обо всех корректировках потребителю.
Потоки изменений работают на трёх уровнях:
весь деплоймент;
база данных;
коллекция.
Также потоки изменений используют средства фреймворка агрегаций. Мы можем отфильтровать те события, которые попадают на обработку в наше приложение. Это полезно, если в коллекции содержатся другие объекты, не интересующие потребителя: добавив шаг фильтрации $match, можно уменьшить поток событий, обрабатываемый приложением.
Примерно так будет выглядеть конвейер агрегации, возвращающий события типа "insert", "update" и "delete" для документов типа "Rule":
WATCH_PIPELINE = [ { "$match": { "operationType": {"$in": ["insert", "update", "delete"]}, "fullDocument.type": "Rule", # Если передать методу `watch` аргумент `full_document="updateLookup"`, можно фильтроваться по содержимому документа через атрибут `fullDocument`. }, }, ]
В качестве примера достаточно отслеживать события на уровне коллекции. Реализуем асинхронный генератор событий:
from pymongo.asynchronous.collection import AsyncCollection async def watch_collection( pipeline: dict, collection: AsyncCollection, resume_token=None ): async with await coll.watch(pipeline, full_document="updateLookup", resume_after=resume_token) as stream: async for change in stream: yield change, stream.resume_token
Возникает закономерный вопрос: как передать эти данные потребителю? GET-запрос для таких задач, очевидно, не пригоден. Как вариант, мы можем использовать API потребителя, отправляя ему сообщения по мере создания новых или редактирования существующих объектов. А что делать в том случае, если потребитель или наше приложение оказались недоступны?
А если потребитель недоступен?
Можно реализовать очередь сообщений. Это будет простым решением без усложнения архитектуры платформы. Все сообщения, генерируемые базой данных, помещаются в очередь, и отправка документа потребителю не триггерится до тех пор, пока он не объявится "в сети".
Пример реализации основного цикла приложения на FIFO-очереди, реализованной при помощи asyncio:
from asyncio import create_task, gather, Queue from logging import getLogger logger = getLogger() EVENT_QUEUES_DICT: dict[str, Queue] = {} # Так можно отслеживать несколько коллекций, оставляя доступ к очереди извне: к примеру, для мониторинга состояния очередей. def get_listener_label(db: str, collection: str) -> str: """Возвращает ключ очереди событий в словаре `EVENT_QUEUES_DICT` для указанной коллекции `collection` базы данных `db`.""" return f"{db}.{coll}" def get_queue_length(db: str, collection: str) -> int: """Возвращает размер очереди событий для указанной коллекции `collection` базы данных `db`.""" listener_label = get_listener_label(db, collection) try: events_queue = EVENT_QUEUES_DICT[listener_label] except KeyError as ex: logger.error(f"Очереди событий для коллекции '{db}.{collection}' не существует.") raise ex return events_queue.qsize() def process_event(event: dict) -> dict | None: """Обрабатывает событие, возвращая сформированное сообщение.""" full_document = event.get("fullDocument") if full_document is None: return None return { "id": str(full_document["id"]), # Конвертация `uuid.UUID` к `str`. "name": full_document["name"], "type": full_document["type"], } def send_message(message): """TODO""" ... async def listen_collection(collection: AsyncCollection, max_event_queue_size: int = 100000): db_name, collection_name = collection.database.name, collection.name events_queue = Queue(maxsize=int(max_event_queue_size)) EVENT_QUEUES_DICT[get_listener_label(db_name, collection_name)] = events_queue async def watch_collection_task_wrapper(queue: Queue, resume_token=None): while True: try: async for event, resume_token in watch_collection( WATCH_PIPELINE, params, app, resume_token=resume_token ): await queue.put(event) logger.info(f"Размер очереди '{db_name}.{collection_name}': {queue.qsize()}.") except Exception as ex: logger.error( f"При мониторинге {db_name}.{collection_name} возникла ошибка. " f"Мониторинг будет перезапущен{' с токеном возобновления'}.", exc_info=ex, ) logger.debug(f"Токен возобновления: {resume_token}.") continue await queue.put(None) async def event_loop(queue: Queue): while True: logger.info(f"Ожидаю событие для '{db_name}.{collection_name}'...") event = await queue.get() if event is None: break logger.debug(f"Получено событие:\n{event}") message = process_event(event) send_message(message) logger.info(f"Размер очереди '{db_name}.{collection_name}': {queue.qsize()}.") watch_task = create_task(watch_collection_task_wrapper(events_queue)) logger.debug("Запущен watch_task.") process_task = create_task(event_loop(events_queue)) logger.debug("Запущен process_task.") await gather(watch_task, process_task) logger.debug("Завершёны `watch_task` и `process_task`.") await events_queue.join() logger.debug(f"Очередь '{db_name}.{collection_name}' завершена.")
В нашем случае этого было недостаточно: некоторые сущности могут быть "не готовы" к отправке в момент генерации события в зависимости от состояния нашей базы данных. Добавим обработку исключения, в которой событие помещается обратно в конец очереди, и так до тех пор, пока оно не будет готово к отправке.
class Rule: ... isReady: bool # Условие готовности отправки объекта, для примера. async def listen_collection(collection: AsyncCollection, max_event_queue_size: int = 100000): ... async def event_loop(queue: Queue): while True: ... logger.debug(f"Получено событие:\n{event}") full_document = event["fullDocument"] is_ready = full_document.get("isReady", False) if not is_ready: await queue.put(event) logger.info(f"Объект {repr(full_document['id'])} не готов к отправке и помещён в конец очереди.") continue message = process_event(event) send_message(message) logger.info(f"Размер очереди '{db_name}.{collection_name}': {queue.qsize()}.") ...
Недостатком такого решения является то, что критическая неполадка на стороне источника может вызвать отказ, что обернётся потерей накопленной очереди. В таком случае можно предусмотреть сохранение идентификатора последнего отправленного документа (см. главу "Мне, пожалуйста, фильтрованного"), или реализовать что-то ещё более комплексное, но, дабы решить эти проблемы, поставим между нашими системами брокер сообщений.
Брокер сообщений
Брокер сообщений стоит между источником данных и их потребителем, выступая их посредником. Наше приложение в терминологии брокеров будет называться "продюсером". Не будем вдаваться в подробности, но посмотрим на основные проблемы, которые решает брокер в подобной интеграции:
Как реализовать транспорт между двумя системами? Брокер, являясь посредником, определяет методы записи и получения сообщений.
Что делать, если одна из систем ушла на техническое обслуживание? Вторая система может выполнять свои функции при доступности брокера.
Как обезопасить себя от потерь данных? Брокер решает задачи хранения очереди и повторных выдачей сообщений при исключениях.

В нашем случае используется RabbitMQ, он подходит для сравнительно небольших объёмов данных (биг-датой нас не назовёшь).
Обработка и отправка сообщения
Доработаем изначальное решение, добавив прослушку событий, их обработку с асинхронной FIFO-очередью и брокер сообщений. Изначальная выгрузка GET-запросом никуда не денется: она нам понадобится для инициализации справочников. Опционально, её можно доработать, прикрутив передачу через брокера, или реализовать подобную фичу вторым запросом, чтобы у потребителя остался выбор между двумя вариантами.
Реализуем функцию записи документа в RMQ:
RMQ_CREDS = { "username": "rmquser", "password": "rmquser", } RMQ_PROPS = { "host": "localhost", "port": "5672", "virtual_host": "/", "heartbeat": 36000, } RMQ_QUEUE_PARAMS = { "queue": "integration.events", "durable": True, } RMQ_PUBLISH_PARAMS = { "exchange": "INTEGRATION.EVENTS.EXCHANGE", "routing_key": "integration.events", } def send_message(message: dict, rmq_credentials: dict, rmq_properties: dict, rmq_queue_params: dict, rmq_publish_params: dict): """Преобразовать словарь сообщения в json-строку и отправить в RMQ.""" try: json_message = json.dumps(message, ensure_ascii=False) credentials = pika.PlainCredentials(**rmq_credentials) parameters = pika.ConnectionParameters( **rmq_properties, credentials=credentials, ) connection = pika.BlockingConnection(parameters) channel = connection.channel() channel.queue_declare(**rmq_queue_params) channel.basic_publish( **rmq_publish_params, body=json_message, ) connection.close() except (Exception, RuntimeError, TimeoutError) as ex: logger.error("Ошибка подключения к RabbitMQ.", exc_info=ex)
Итого
В данном решении мы постарались сфокусироваться на:
эффективности: функции фильтрации и триггерной выгрузки по событиям исключают лишнюю нагрузку на сеть;
высокодоступности: благодаря брокеру передача сообщений не останавливается при недоступности продюсера или потребителя, а обычные GET-запросы остаются запасным планом при неполадках на брокере;
надёжности: в случае, если брокер сам станет недоступным, у потребителя есть запасные варианты получения драгоценных метаданных.
С другой стороны, мы жертвуем:
простотой реализации: всё же, развёртывание и поддержка брокера — это дополнительные затраты;
простотой поддержки: реализация всех указанных фичей — это дополнительные работы и сложность при, казалось бы, простой выгрузке данных по требованию.
Надеюсь, данная статья станет полезной пищей для размышлений при разработке ваших продуктов.
Источники
1. https://www.mongodb.com/docs/manual/changeStreams/#open-a-change-stream.
2. https://www.mongodb.com/resources/products/capabilities/change-streams.