Настройка синхронизации изменений данных MongoDB Atlas в режиме реального времени с Microsoft Fabric

Microsoft Fabric

В этой статье описывается, как интегрировать MongoDB Atlas с Microsoft Fabric с помощью открытой функции зеркального отображения в Fabric. Мы рекомендуем эту архитектуру для приема операционных данных с высокой точностью, низкой задержкой в OneLake.

Architecture

Ускоритель зеркального отображения MongoDB Atlas to Fabric реализует открытую функцию зеркального отображения. Открытая зеркальная база данных в Fabric раскрывает зону посадки в OneLake. Приложения записывают файлы Parquet, содержащие данные об изменении MongoDB, в эту зону посадки, следуя спецификации открытого зеркального отображения.

На следующей схеме показано, как приложение для зеркального отображения данных, развернутое в службе приложений Azure, передает события изменений MongoDB Atlas в Fabric.

Схема архитектуры, демонстрирующая интеграцию Fabric с MongoDB Atlas с помощью открытого зеркального отображения.

Скачайте файл Visio для этой архитектуры.

Поток данных

Следующий поток данных соответствует предыдущей схеме:

  1. Создайте открытую зеркальную базу данных в Fabric с помощью REST API или портала Fabric.

  2. Получите URL-адрес целевой зоны, связанный с зеркальной базой данных.

  3. Разверните акселератор зеркального отображения с помощью Terraform или Azure Resource Manager.

  4. Приложение выполняет несколько ключевых функций:

    • Выполняет начальную загрузку исторических данных из MongoDB
    • Подписывается на потоки изменений MongoDB для записи текущих операций вставки, обновления и удаления
    • Записывает захваченные данные об изменениях в файлы Parquet в зону посадки Fabric
  5. Структура автоматически выполняет следующие действия:

    1. Обнаруживает новые файлы Parquet в зоне загрузки
    2. Преобразует новые файлы Parquet в таблицы Delta, поддерживающие изменения схемы
    3. Сохраняет зеркальные таблицы, синхронизированные с исходными коллекциями MongoDB
    4. Создает семантику по умолчанию для Power BI
  6. Ресурсы Power BI, Lakehouse в Fabric и Хранилища данных Fabric могут использовать синхронизированные данные для аналитики и отчетности.

Components

  • Открытое зеркальное отображение в Fabric — это возможность управляемой репликации данных, которая синхронизирует внешние источники данных в Fabric с помощью открытых форматов таблиц. В этой архитектуре он постоянно загружает данные об изменениях MongoDB, преобразуя файлы Parquet в таблицы Delta и оставляя их синхронизированными с изменениями MongoDB Atlas.

  • OneLake — это унифицированное открытое озеро данных для Fabric, которое предоставляет централизованное хранилище для всех рабочих нагрузок Fabric. В этой архитектуре она выступает в качестве стартовой зоны для данных изменений и общего уровня хранения, который используется нижестоящими службами Fabric.

  • Lakehouses в Fabric — это унифицированные платформы данных, которые объединяют озеро данных с возможностями аналитики и выполнения SQL-запросов по таблицам Delta. В этой архитектуре озеро-хранилище предоставляет зеркальные таблицы Delta и встроенный доступ к T‑SQL для аналитики и выполнения запросов.

  • Семантические модели в Fabric — это наборы данных Power BI, определяющие удобные для бизнеса метаданные и связи для поддержки аналитических запросов и отчетов. В этой архитектуре lakehouse автоматически генерирует их для ускорения аналитики и отчетности в Power BI.

  • Служба приложений — это полностью управляемая платформа для размещения веб-приложений и фоновых служб. В этой архитектуре размещается приложение репликации на основе Python, которое организует обработку изменений в Fabric.

  • Потоки изменений MongoDB предоставляют механизм для отслеживания изменений данных в режиме реального времени из коллекций MongoDB. В этой архитектуре они записывают операции вставки, обновления и удаления из MongoDB Atlas для непрерывной синхронизации данных.

  • Terraform — это инструмент "инфраструктура как код" (IaC), используемый для декларативной подготовки облачных ресурсов. В этой архитектуре шаблоны автоматизируют развертывание необходимых ресурсов Azure и Fabric.

  • Power BI — это платформа бизнес-аналитики для создания интерактивных панелей мониторинга и отчетов. В этой архитектуре она визуализирует зеркальные таблицы Delta с помощью Direct Lake для высокопроизводительной аналитики в режиме реального времени.

На следующей схеме показана архитектура интеграции зеркального отображения.

Схема, показывющая архитектуру интеграции зеркального отображения.

Скачайте файл PowerPoint данной архитектуры.

Открытое зеркальное отображение позволяет вашему приложению записывать изменения данных MongoDB непосредственно в Fabric, преобразовывать их в формат Delta и сразу же предоставлять доступ к озеру данных, хранилищу данных, аналитике в режиме реального времени и Power BI.

Альтернативы

Fabric поддерживает другие шаблоны интеграции с MongoDB Atlas. Эти альтернативные варианты могут соответствовать рабочей нагрузке в зависимости от пороговых значений задержки, операционных ограничений или существующей инфраструктуры.

В следующих разделах описываются альтернативные подходы к интеграции данных для разнообразных нужд оперативной и аналитической обработки.

Интеллект в реальном времени с потоками событий и хранилищами событий

Система Fabric Real-Time Intelligence предоставляет встроенный, без необходимости программирования путь для приема данных с использованием коннектора захвата изменений данных (CDC) MongoDB для потока событий. Коннектор передаёт события изменений из MongoDB Atlas непосредственно в Fabric.

Системы Eventstreams направляют события CDC в следующие пункты назначения:

  • Хранилище событий (база данных KQL) для аналитики в режиме реального времени, обнаружения аномалий и наблюдаемости
  • OneLake для обработки нижнего озера или хранилища данных
  • Папки Lakehouse для Apache Spark, аналитики SQL или рабочих нагрузок машинного обучения

Используйте этот подход, если вам потребуются следующие возможности:

  • Операционные панели мониторинга в режиме реального времени
  • Журнал высокой пропускной способности и обработка событий
  • Мониторинг и обнаружение на основе KQL
  • Сценарии с низкой задержкой без пользовательского кода

Триггеры Atlas, функции и OneLake (модель push)

Этот подход использует триггеры MongoDB Atlas для вызова функций пользовательских данных Fabric или Функций Azure. Функция записывает обновленный документ в OneLake с помощью API, совместимого с Azure Data Lake Storage.

Схема архитектуры с триггерами и функциями MongoDB Atlas, которые помещают данные в OneLake.

Поток данных
  1. Триггер Atlas обнаруживает операцию вставки, обновления или удаления.

  2. Триггер вызывает функцию Atlas.

  3. Функция Atlas отправляет нагрузку в пользовательскую функцию Fabric или функцию Функции Azure.

  4. Функция пользовательских данных Fabric или Функции Azure записывает документ JSON в OneLake.

  5. При необходимости конвейер Fabric преобразует и загружает документ в хранилище данных типа "озеро" или в стандартное хранилище данных.

Используйте модель push-отправки, если требуется прием практически в режиме реального времени, но не удается развернуть акселератор зеркального отображения или если вы предпочитаете бессерверную интеграцию на основе событий.

Каналы Fabric и коннекторы MongoDB (модель pull)

Конвейеры передачи данных Fabric включают соединитель MongoDB, который поддерживает установленные локально MongoDB и MongoDB Atlas.

Обычно соединители используются для следующих задач:

  • Исторические нагрузки данных
  • Ежедневные или почасовые запланированные синхронизации
  • Инкрементное прием данных с помощью запросов MongoDB
  • Многооблачная или гибридная интеграция

Конвейеры могут загружать документы в озеро, использующее Delta, Parquet, Avro, JSON или CSV. Они также могут загружать документы в хранилище данных.

Мы рекомендуем эту модель для пакетных рабочих нагрузок или сценариев, которые не требуют приема в режиме реального времени. Например, можно использовать конвейер Fabric для выполнения ночной копии из MongoDB Atlas в хранилище данных для создания исполнительной отчетности. Данные должны отражать конечное состояние в конце дня, поэтому запланированная с помощью соединителя MongoDB система достаточна и позволяет избежать сложности непрерывного ввода данных.

Apache® и Apache Spark™ являются зарегистрированными товарными знаками или товарными знаками Apache Software Foundation в США и/или других странах. Использование этих меток не подразумевает подтверждения от Apache Software Foundation.

Архитектура соединителя для интеграции MongoDB с конвейерами Fabric.

В левом углу расположена прямоугольная область с надписью "Потребители", включающая внутренние приложения, клиентские службы и API для использования вне экосистемы Майкрософт в любом канале. Двунаправленная стрелка соединяет потребителей с полем, помеченным уровнем операционных данных. Это поле содержит MongoDB Atlas в Azure, модель данных документов, архитектуру распределенных систем и облачные или локальные компоненты. Линии подключают уровень операционных данных и реляционную систему управления базами данных (RDBMS) в табличном формате в верхней части и журналы в неструктурированном формате внизу. Значок конвейера, представляющий соединитель источника MongoDB, подключает рабочий уровень данных к другому поле в правом хранилище корпоративных данных. В этом поле строки подключают Spark, Fabric и SQL к машинному обучению, аналитике больших данных и панели мониторинга бизнес-аналитики. Линия соединяет блок корпоративного хранилища данных, значок соединителя приемника MongoDB и хранилище аналитики MongoDB.

Вы также можете использовать потоки данных в Fabric для интеграции с MongoDB Atlas. Модель запросов подходит для аналитико‑орентированных и самостоятельных сценариев бизнес‑аналитики.

Потоки данных предоставляют возможность без кода, ориентированную на Power BI, для загрузки данных из MongoDB Atlas в OneLake. Если вы используете Power BI в качестве основного средства аналитики, можно использовать соединитель Dataflow 2-го поколения MongoDB для следующих задач:

  • Запланированное прием
  • Извлечение отфильтрованных исторических данных
  • Формирование документов и преобразование без кода
  • Результаты, записанные непосредственно в OneLake с помощью стандартного вывода Dataflows Gen2

Пакетная интеграция

Вы можете использовать пакетную или микро-пакетную интеграцию для перемещения исторических или фильтруемых данных из MongoDB Atlas в OneLake. Организации могут использовать конвейеры Fabric и соединитель Spark MongoDB версии 10.x, который поддерживает шаблоны приема пакетов и потоковой передачи.

  • Используйте соединитель Spark для пакетной интеграции и потоковой передачи. Коннектор MongoDB Spark позволяет гибко загружать данные на основе DataFrame в хранилище данных Fabric. Он поддерживает следующие возможности:

    • Полная историческая загрузка пакета
    • Фильтрованные или инкрементальные нагрузки при помощи запросов MongoDB
    • Микро пакетная или непрерывная потоковая передача с помощью структурированной потоковой передачи Spark
    • Запись в таблицы Delta или Parquet в OneLake

    Этот подход является оптимальным для команд, которые в основном работают с Apache Spark, или для рабочих нагрузок, требующих преобразований в рамках загрузки данных.

  • При необходимости используйте конвейеры Fabric для пакетной загрузки данных. Конвейеры Fabric могут оркестрировать пакетную приемку MongoDB через коннектор MongoDB, но коннектор Spark обеспечивает большую гибкость для сценариев пакетной и потоковой передачи.

Подробности сценария

MongoDB Atlas — это общее рабочее хранилище для внутренних приложений, клиентских служб и интеграции, отличных от Майкрософт. Организации могут использовать Fabric для объединения данных Atlas с реляционными, потоковыми и неструктурированными источниками для анализа данных, бизнес-аналитики и машинного обучения в большом масштабе.

Потенциальные варианты использования

Розничная торговля:

  • Оптимизация комплектования продуктов и продвижения.
  • Клиент 360 и гипер-персонализация
  • Прогнозирование спроса и запаса
  • Интеллектуальный поиск и рекомендации

Банковские услуги и финансы:

  • Обнаружение и предотвращение мошенничества
  • Персонализированные финансовые продукты и предложения

Телекоммуникации:

  • Аналитика качества сети и оптимизация
  • Агрегирование телеметрии периферийных вычислений

Автомобильной:

  • Аналитика подключенных транспортных средств
  • Обнаружение аномалий в коммуникации Интернета вещей (IoT)

Производство:

  • Прогнозное обслуживание
  • Оптимизация инвентаризации и склада

Пример: пакет продуктов для розничной торговли

Используйте данные шаблона продаж для разработки пакетов, которые увеличивают размер корзины и маржу.

Источники данных:

  • Каталог продуктов в MongoDB Atlas
  • Данные о продажах в SQL Azure

Поток:

  1. Используйте конвейер Fabric для загрузки данных о продуктах и продажах в хранилище данных или в лейкхаус.

  2. Примените обновления CDC или на основе событий для синхронизации почти в режиме реального времени в дополнение к начальной загрузке.

  3. Сопоставление моделей и шаблонов совместных покупок, таких как анализ корзины рынка, и отображение метрик через Power BI.

Снимок экрана этапов конвейера и диаграмм для группирования продуктов, включая продажи по продуктам, годам, регионам и сходствам.

Рекомендации по анализу:

  • Комплект товаров, такие как ручка и запасной стержень для чернил.
  • Продвижение комплекта в регионах с высоким уровнем соответствия.

Пример: Продвижение продуктов для розничной торговли

Для рекомендации дополнительных продуктов используйте данные о поведении клиентов, рентабельности и сопоставлении продуктов.

Подход:

  • Тренируйте модели машинного обучения в записных книжках Fabric Spark или интегрируйте с Машинное обучение Azure.
  • Используйте OneLake в качестве хранилища функций и обслуживайте прогнозы для приложений Power BI или нижестоящих приложений.

Снимок экрана потока данных и рабочего процесса машинного обучения для продвижения продукта на основе поведения клиентов и характеристик продукта.

Если модель достигает высокой точности, она создает приоритетный набор альтернативных рекомендаций по продукту для каждого клиента или сегмента.

Рекомендации

Эти рекомендации реализуют основные принципы платформы Azure Well-Architected Framework, которая представляет собой набор руководящих принципов, которые можно использовать для улучшения качества рабочей нагрузки. Для получения дополнительной информации см. Well-Architected Framework.

Безопасность

Безопасность обеспечивает гарантии от преднамеренного нападения и неправильного использования ценных данных и систем. Для получения дополнительной информации см. контрольный список по проверке проектирования на безопасность.

  • Используйте ПРОТОКОЛ HTTPS и последние версии TLS для функций пользовательских данных Fabric и конечных точек Функций Azure.
  • Проверьте входящие полезные данные в виде событий изменения MongoDB.
  • Настройте проверки подлинности Microsoft Entra ID и управления доступом на основе принципа наименьших привилегий (RBAC) для Fabric.
  • Настройте OneLake для наследования моделей безопасности и проверки подлинности Data Lake Storage.
  • Используйте MongoDB Atlas для встроенных элементов управления для доступа, сетевой изоляции, шифрования и аудита.

Оптимизация затрат

Оптимизация затрат фокусируется на способах сокращения ненужных расходов и повышения эффективности работы. Дополнительные сведения см. в контрольном списке проектной экспертизы для оптимизации затрат.

  • Оптимизация емкости Rightsize Fabric и консолидация рабочих нагрузок, где это целесообразно.
  • Пакетное изменение документов для уменьшения количества вызовов функций и объёма небольших файлов.
  • Сжимайте и оптимизируйте таблицы lakehouse и запланируйте нагрузочные преобразования вне пиковых периодов.
  • Используйте сжатие Parquet или Delta, например Snappy, для уменьшения использования хранилища и повышения производительности сканирования.
  • Оптимизация размера кластеров Atlas и оценка шардирования и уровней хранилища.

Эффективность работы

Эффективность производительности — это способность рабочей нагрузки эффективно масштабироваться в соответствии с требованиями пользователей. Для получения дополнительной информации см. контрольный список проверки проектного решения на эффективность производительности .

  • Объединение нескольких событий изменения для уменьшения накладных расходов на небольшие файлы.
  • Используйте Delta для оптимизированных запросов Spark, SQL и BI.
  • Выберите стратегии секционирования и распределения, которые подходят для больших складов.
  • Настройте параллелизм конвейера и примените фильтры pushdown к соединителям MongoDB.
  • Отслеживайте задержку приема данных и реализуйте повторные попытки и идемпотентные объединения вставок и обновлений.
  • Запланируйте OPTIMIZE и VACUUM операции по обслуживанию lakehouse.
  • Используйте модель push-отправки для приема на основе событий практически в режиме реального времени.
  • Используйте пуллинговую модель для запланированных, пакетных или микропакетных нагрузок.
  • Используйте хранилище данных для управляемых реляционных моделей и корпоративной бизнес-аналитики.
  • Используйте конечные точки SQL Lakehouse для упрощенного использования SQL над Delta Lake без развертывания хранилища.

Соавторы

Корпорация Майкрософт поддерживает эту статью. Следующие авторы написали эту статью.

Основные авторы:

Другие участники:

  • Сунил Сабат | Главный руководитель программы — команда фабрики данных Azure

Чтобы видеть непубличные профили на LinkedIn, войдите в аккаунт LinkedIn.

Дальнейшие действия