Выполнение конвейеров агрегации между источниками
Вы напишете конвейеры агрегации, объединяющие данные коллекций Atlas с файлами JSON или Parquet, хранящимися в S3, в одном запросе.
«Выполнение конвейеров агрегации между источниками» — бесплатный урок 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 — локальная установка не требуется.
Все уроки этого курса
- Что такое Atlas Data Federation
- Сопоставление источников S3 и Atlas с виртуальным пространством имён
- Выполнение конвейеров агрегации между источниками
- Секционирование данных S3 для производительности запросов