MongoDB Academy · Урок

Выполнение конвейеров агрегации между источниками

Вы напишете конвейеры агрегации, объединяющие данные коллекций Atlas с файлами JSON или Parquet, хранящимися в S3, в одном запросе.

Урок 3 из 413 шагов

«Выполнение конвейеров агрегации между источниками» — бесплатный урок MongoDB Academy на CoddyKit. Это урок 3 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения MongoDB Academy, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс MongoDB Academy содержит 4 уроков всего.

Что особенного в конвейерах из разных источников

Конвейер агрегации из разных источников в Atlas Data Federation выполняет те же этапы агрегации, которые Вы знаете по MongoDB, но данные для каждого этапа могут поступать из разных физических систем — сегмента S3, работающего кластера Atlas или обоих источников одновременно. Объединённый механизм запросов прозрачно обрабатывает маршрутизацию, распределение запросов и объединение результатов. С точки зрения приложения это выглядит как запрос к одной коллекции MongoDB.

Простой поиск в разных источниках

Самый простой запрос к разным источникам — это вызов find() для виртуальной коллекции, основанной на файлах S3. Механизм запросов читает файлы, разбирает их и применяет фильтр. Поля фильтра, совпадающие с атрибутами разделов в пути, автоматически приводят к исключению ненужных файлов. Поля, которые не совпадают с атрибутами разделов, применяются как фильтр после чтения данных.

// Virtual collection 'events' backed by S3 JSON files
// Path: /data/events/{year int}/{month int}/*.json

// This query prunes to /data/events/2025/1/ only
const jan2025 = await db.collection('events').find({
  year: 2025,
  month: 1,
  eventType: 'purchase'   // post-read filter (not a partition attr)
}).toArray()

Агрегация данных S3 как данных работающей коллекции

К виртуальным коллекциям на основе S3 можно применять любые этапы агрегации: $match, $group, $project, $sort, $limit. По возможности механизм запросов переносит выполнение этапов ближе к источнику, особенно $match для отсечения разделов и столбцов в Parquet, а остальные этапы выполняет в собственном вычислительном слое после чтения данных.

// Group S3-archived events by region and count
const summary = await db.collection('events_2024').aggregate([
  { $match: { year: 2024, month: { $in: [10, 11, 12] } } },  // pruning
  { $group: { _id: '$region', total: { $sum: 1 }, revenue: { $sum: '$amount' } } },
  { $sort: { revenue: -1 } },
  { $limit: 10 }
]).toArray()

Объединение Atlas и S3 с помощью $lookup

Самый мощный шаблон работы с разными источниками — использование $lookup для объединения работающей коллекции Atlas с коллекцией, архивированной в S3. Начните конвейер с работающей коллекции (она будет ведущей) и выполните поиск в виртуальной коллекции на основе S3. Всегда размещайте $match в начале конвейера, чтобы свести к минимуму количество выполняемых поисков.

// Join live customers (Atlas) with archived orders (S3)
const result = await db.collection('customers').aggregate([
  { $match: { tier: 'gold', region: 'EU' } },   // filter live data first
  { $lookup: {
    from: 'orders_archive',   // virtual S3-backed collection
    let: { custId: '$_id' },
    pipeline: [
      { $match: { $expr: { $eq: ['$customerId', '$$custId'] } } },
      { $project: { orderId: 1, amount: 1, date: 1 } }
    ],
    as: 'orderHistory'
  }},
  { $addFields: { totalSpend: { $sum: '$orderHistory.amount' } } },
  { $sort: { totalSpend: -1 } },
  { $limit: 50 }
]).toArray()

Агрегация данных из нескольких кластеров Atlas

Если в Вашем объединённом экземпляре есть несколько хранилищ кластеров Atlas, Вы можете объединять коллекции из разных кластеров Atlas в одном конвейере. Это полезно для многопользовательских или мультирегиональных развертываний, где данные распределены по отдельным кластерам и требуются отчёты между кластерами без их объединения или создания отдельной базы данных для отчётности.

// Virtual collections pointing to different Atlas clusters
// 'orders_us' -> Atlas cluster in US
// 'orders_eu' -> Atlas cluster in EU

// Union results from two clusters
db.orders_us.aggregate([
  { $match: { date: { $gte: ISODate('2025-01-01') } } },
  { $unionWith: {
    coll: 'orders_eu',
    pipeline: [{ $match: { date: { $gte: ISODate('2025-01-01') } } }]
  }},
  { $group: { _id: '$status', count: { $sum: 1 } } }
])

Запись результатов в Atlas с помощью $out и $merge

После выполнения агрегации из разных источников Вы можете записать результаты обратно в работающую коллекцию Atlas с помощью $out или $merge. Это шаблон ETL: прочитать исторические данные из S3, объединить их с актуальными данными, вычислить агрегаты и записать результаты в материализованную коллекцию-представление в Atlas, из которой приложения смогут быстро и недорого читать данные.

// ETL: aggregate S3 archive + Atlas, write result to Atlas
db.events_2024.aggregate([
  { $match: { year: 2024 } },
  { $group: {
    _id: { region: '$region', month: '$month' },
    sessions: { $sum: 1 },
    revenue: { $sum: '$amount' }
  }},
  { $merge: {
    into: { db: 'reporting', coll: 'monthly_summary' },
    whenMatched: 'replace',
    whenNotMatched: 'insert'
  }}
])

Отсечение столбцов в Parquet: чтение только нужных данных

При запросе файлов Parquet Data Federation применяет отсечение столбцов: если на этапе $project указаны только определённые поля, механизм запросов читает из файла Parquet только соответствующие столбцы (данные в нём хранятся по столбцам). Это может сократить объём прочитанных данных более чем на 90 % по сравнению с чтением всех столбцов. Размещайте $project как можно раньше в конвейере, чтобы максимально эффективно отсекать столбцы.

// Column pruning: only reads 'region', 'amount', 'date' columns from Parquet
db.events_2024.aggregate([
  { $project: { region: 1, amount: 1, date: 1, _id: 0 } },  // early project
  { $match: { region: 'EU' } },
  { $group: { _id: '$region', totalRevenue: { $sum: '$amount' } } }
])
// Other columns (userId, sessionId, metadata, etc.) are never read from disk

Мониторинг производительности объединённых запросов

Atlas Data Federation записывает подробности выполнения запросов в интерфейсе Atlas на вкладке История запросов. Для каждого запроса отображаются: объём обработанных данных, время выполнения и количество просканированных файлов и разделов. Большой объём обработанных данных обычно означает, что атрибуты разделов отсутствуют или запрос не соответствует ни одному ключу раздела. Используйте этот журнал, чтобы настроить конфигурацию хранилища и шаблоны запросов.

// Get query stats via the admin DB on the federated instance
db.adminCommand({ currentOp: 1 })
// Shows active federated queries with bytes read, duration

// In Atlas UI: Data Federation > Query History
// Shows past queries, duration, data processed, and cost estimate

Обработка различий схем в разных источниках

Файлы S3 из разных периодов или систем могут иметь разные схемы: имена полей, типы и структуру. Data Federation корректно обрабатывает такие различия: отсутствующие поля возвращают значение null, а дополнительные поля включаются в результат. В конвейере можно использовать $ifNull, $cond и $convert, чтобы привести различающиеся схемы к единому виду перед группировкой или объединением.

// Normalise schema variations across old and new S3 file formats
db.events.aggregate([
  { $addFields: {
    // Old format: 'user_id', New format: 'userId'
    userId: { $ifNull: ['$userId', '$user_id'] },
    // Old format: string amount, New format: number
    amount: { $convert: { input: '$amount', to: 'double', onError: 0 } }
  }},
  { $group: { _id: '$userId', total: { $sum: '$amount' } } }
])

Кэширование результатов объединённых запросов

Data Federation не кэширует результаты между запросами — каждый запрос заново читает исходные данные. Для панелей мониторинга, которые регулярно выполняют один и тот же отчёт, используйте шаблон ETL: запланируйте агрегацию, записывающую результаты в коллекцию Atlas с помощью $merge, а затем настройте панель на запросы к быстрой коллекции Atlas. Atlas Triggers могут планировать такое обновление с любым интервалом cron.

Ограничения конвейеров из разных источников

Учитывайте текущие ограничения: 1) транзакции не поддерживаются в объединённых экземплярах; 2) использование индексов применяется только к коллекциям на основе Atlas, но не к файлам S3; 3) очень большие наборы результатов могут превысить время ожидания — используйте $out/$merge, чтобы записать результаты, вместо передачи их потоком обратно; 4) задержка выше, чем у запроса к работающему Atlas, из-за операций ввода-вывода S3, поэтому такой подход не подходит для запросов в реальном времени, обращённых непосредственно к пользователям.

Быстрая проверка

Проверьте, насколько хорошо Вы усвоили концепции MongoDB и баз данных NoSQL из этого урока.

Итоги урока

В этом уроке Вы узнали, что конвейеры агрегации из разных источников используют те же этапы MongoDB для виртуальных коллекций на основе S3 или кластеров Atlas, $lookup позволяет объединять актуальные данные Atlas с архивами S3 в одном конвейере, а ранний этап $project позволяет отсекать столбцы в файлах Parquet и значительно уменьшать объём просматриваемых данных. Далее мы рассмотрим разбиение данных S3 для повышения производительности запросов.

Можно начать бесплатно

Изучай JavaScript с ИИ-репетитором — бесплатно

Пиши и запускай код прямо в браузере, получай мгновенную помощь от ИИ-репетитора 24/7 и продолжи учиться на сайте или в приложении.

Курсы
30
Уроки
120

Часто задаваемые вопросы

Урок «Выполнение конвейеров агрегации между источниками» бесплатный?

Да — полный текст урока «Выполнение конвейеров агрегации между источниками» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс MongoDB Academy, подпишись на CoddyKit PRO. Курс MongoDB Academy содержит 4 уроков всего.

Чему я научусь в уроке «Выполнение конвейеров агрегации между источниками»?

Вы напишете конвейеры агрегации, объединяющие данные коллекций Atlas с файлами JSON или Parquet, хранящимися в S3, в одном запросе. Ты практикуешь MongoDB Academy с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.

Нужен ли мне опыт, чтобы начать MongoDB Academy?

Предыдущий опыт не требуется. MongoDB Academy на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 3 из 4.

Сколько времени занимает урок «Выполнение конвейеров агрегации между источниками»?

Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.

Можно ли писать и запускать код в этом уроке MongoDB Academy?

Да. Каждый урок MongoDB Academy включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.

Все уроки этого курса

  1. Что такое Atlas Data Federation
  2. Сопоставление источников S3 и Atlas с виртуальным пространством имён
  3. Выполнение конвейеров агрегации между источниками
  4. Секционирование данных S3 для производительности запросов
← Назад к MongoDB Academy