MongoDB: 集計パイプライン入門:aggregateメソッドと基本ステージ
最終更新:2026-08-26
集計パイプラインはMongoDB最強のデータ分析ツールです。マスターすればSQL分析シナリオの90%をMongoDBで代替できます。
本コースでは、集計パイプラインの基礎、ステージの概念、実行順序を解説します。
1. このレッスンで学ぶこと
aggregate()メソッドの基本構文- パイプラインステージの概念と実行順序
- 基本ステージ:$match、$group、$project、$sort、$limit、$skip、$count
- 実践的なシンプル集計
2. 集計パイプラインとは
概念説明: 集計パイプラインはMongoDB最強のデータ分析フレームワークです。データ処理を複数のステージに分解し、順次実行します。各ステージは前段の出力を入力として受け取り、フィルタリング・変換・グループ化・統計計算を行い、最終結果を生成します。この「パイプライン」モデルはUNIXパイプの概念に基づき、複雑なデータ分析を組み合わせ可能でデバッグしやすくします。
動作原理: 集計パイプラインは配列の順序でステージごとに実行されます。最初のステージがコレクションの全ドキュメントを取得し、各後続ステージがドキュメントストリームを変換(フィルタリング・射影・グループ化等)して次のステージに渡します。MongoDBクエリオプティマイザは$matchと$sortをパイプライン先頭に移動してインデックスを活用しようとします。単一ステージのメモリ制限は100MBで、超過時はallowDiskUse: trueが必要です。
graph LR
A[コレクション<br/>1000ドキュメント] --> B[$match<br/>フィルタ]
B --> C[$group<br/>グループ集計]
C --> D[$sort<br/>ソート]
D --> E[$limit<br/>制限]
E --> F[結果]
style B fill:#fff3cd
style C fill:#d4edda
sequenceDiagram
participant DB as MongoDB
participant S1 as $match
participant S2 as $group
participant S3 as $sort
participant S4 as $limit
DB->>S1: 1000ドキュメント
Note over S1: フィルタ: category=Electronics
S1->>S2: 250ドキュメント
Note over S2: ブランドでグループ化<br/>$sum, $avg
S2->>S3: 15グループ
Note over S3: 合計の降順でソート
S3->>S4: 15グループ
Note over S4: 上位10件取得
S4-->>DB: 10グループ結果
UNIXパイプとの類似:
# 評価4以上の製品を検索し、価格降順でソート、上位10件
cat products.json | jq 'select(.rating > 4)' | jq 'sort_by(-.price)' | head -10
# MongoDB集計
db.products.aggregate([
{ $match: { rating: { $gt: 4 } } },
{ $sort: { price: -1 } },
{ $limit: 10 }
]);
3. aggregate()の基本構文
概念説明: db.collection.aggregate(pipeline, options)が集計パイプラインを実行するエントリメソッドです。pipelineはステージオブジェクトの配列、optionsは実行動作を制御(メモリ制限、タイムアウト、インデックスヒント)。集計関数はカーソルを返し、toArray()で一括取得またはforEach()で順次反復できます。
| パラメータ | 型 | 説明 |
|---|---|---|
pipeline |
Array | ステージオブジェクトの配列。順次実行 |
options.allowDiskUse |
Boolean | ディスクへの一時書き込みを許可(デフォルト: false) |
options.maxTimeMS |
Number | タイムアウト(ミリ秒) |
options.batchSize |
Number | 1バッチあたりの返却ドキュメント数 |
options.hint |
String/Object | 指定インデックスの強制使用 |
// === 基本的な使用方法 ===
db.products.aggregate([
{ $match: { category: 'Electronics' } },
{ $group: { _id: '$brand', total: { $sum: 1 } } }
]);
// === カーソルを返す ===
const cursor = db.products.aggregate([...]);
const results = await cursor.toArray();
// === 集計オプション ===
db.products.aggregate(
[{ $match: {} }],
{
allowDiskUse: true, // ディスク使用許可(大規模データセット処理)
maxTimeMS: 30000, // 30秒タイムアウト
batchSize: 100, // バッチサイズ
hint: 'category_1' // インデックス強制使用
}
);
4. 基本ステージ詳細
概念概要: 集計パイプラインは30以上のステージ演算子を提供します。本節では最も基本的な7つに焦点を当てます:$match(フィルタ)、$group(グループ集計)、$project(射影・計算)、$sort(ソート)、$limit(制限)、$skip(スキップ)、$count(カウント)。これらを組み合わせれば、日常的な分析シナリオの80%をカバーします。
graph TB
A[基本ステージ] --> B[$match<br/>ドキュメントフィルタ<br/>findのフィルタと同様]
A --> C[$group<br/>グループ集計<br/>$sum/$avg/$min/$max]
A --> D[$project<br/>フィールド射影<br/>選択/計算フィールド]
A --> E[$sort<br/>ソート<br/>1 昇順 -1 降順]
A --> F[$limit/$skip<br/>ページネーション<br/>制限/スキップ]
A --> G[$count<br/>カウント<br/>出力ドキュメント数]
B --> H["✅ 先頭に配置<br/>インデックス使用でデータ量削減"]
style B fill:#fff3cd
style C fill:#d4edda
style H fill:#d4edda
(1) $match:ドキュメントフィルタ
$matchの最適化原則: $matchは集計パイプラインで唯一インデックスを利用できるステージのため、パイプラインの先頭に配置することが最も重要な最適化原則です。理由:1. $matchが$groupより先に実行され、インデックスを直接使用してフィルタリングすることで下流の処理データ量を削減;2. MongoDBクエリオプティマイザは$matchを$lookupより前に押し下げようとしますが、ユーザー定義ステージの順序を変更しません;3. $groupの後に$matchを実行すると、全データを集計してからフィルタリングすることになり、計算リソースを無駄に消費します。
// === $match は find のフィルタと同様 ===
db.products.aggregate([
{ $match: { category: 'Electronics', price: { $gte: 100 } } }
]);
| $matchの利点 | 説明 |
|---|---|
| パイプライン前段でのフィルタリング | 後続ステージの処理データ量を削減 |
| インデックスを使用可能 | 高パフォーマンス(findと同様) |
(2) $group:グループ集計
$groupのグループ化原則: $groupは集計パイプラインで最も重要なステージです。ドキュメントストリームを_idフィールドでグループ化し、各グループ独立にアキュムレータ計算を行い、単一の集計結果を出力します。_idの値がグループ粒度を決定します:フィールド参照('_id: $category')は単一フィールドでグループ化、オブジェクト式('_id: {year, month}')はフィールド組み合わせでグループ化、nullはグループ化なし(コレクション全体の統計)。
// === フィールドでグループ化 ===
db.products.aggregate([
{ $group: {
_id: '$category', // グループ化フィールド
count: { $sum: 1 }, // カウント
avgPrice: { $avg: '$price' },
maxPrice: { $max: '$price' },
minPrice: { $min: '$price' }
}}
]);
出力:
TEXT 📖 参照専用[ { _id: 'Electronics', count: 250, avgPrice: 599, maxPrice: 1999, minPrice: 99 }, { _id: 'Books', count: 200, avgPrice: 29, maxPrice: 79, minPrice: 9 }, ... ]
// === 複数フィールドでのグループ化 ===
db.orders.aggregate([
{ $group: {
_id: { year: { $year: '$createdAt' }, month: { $month: '$createdAt' } },
total: { $sum: '$total' },
count: { $sum: 1 }
}}
]);
// === コレクション全体のグループ化 ===
db.products.aggregate([
{ $group: {
_id: null, // グループ化なし、コレクション全体の統計
totalProducts: { $sum: 1 },
avgPrice: { $avg: '$price' }
}}
]);
(3) $project:フィールド射影
$projectの2つの用途: $projectには2つの目的があります—1. フィールド選択(SQLのSELECTに相当)、1または0でフィールド出力を制御;2. フィールド計算(SQLのASに相当)、式を使って新しいフィールドを作成。注意:$projectの1/0ルールでは、_idはデフォルトで出力されます(非表示にするには明示的に0に設定)。他のフィールドは、いずれか1つでも1に設定すれば他はデフォルトで0(ホワイトリストモード)、いずれか1つでも0に設定すれば他はデフォルトで1(ブラックリストモード)。
// === 出力フィールド選択 ===
db.products.aggregate([
{ $project: {
sku: 1,
title: 1,
price: 1,
discountedPrice: { $multiply: ['$price', 0.9] } // 新フィールド計算
}}
]);
// === フィールド名変更 ===
db.products.aggregate([
{ $project: {
productName: '$title', // title → productName に改名
price: 1,
category: 1
}}
]);
(4) $sort:ソート
$sortのメモリ制限と最適化: $sortはメモリ内でソートを実行し、デフォルトのメモリ制限は100MBです—超過するとエラーになります(allowDiskUse: trueを設定しない限り)。最適化戦略:1. $sortの前に$matchを使用してデータ量を削減(100,000件のソートより1,000件のソートの方が遥かに高速);2. ソートフィールドにインデックスを作成(インデックスは既にソート済みのため、MongoDBはインデックス順で直接結果を返せます);3. $sortと$limitを組み合わせると、MongoDBはTop Nヒープのみ維持(全データセットをソートせず)、メモリ使用量をO(N)からO(limit)に削減;4. 複合インデックス{category: 1, price: -1}は$matchと$sortの両方を満たし、「カバークエリ + ソート」を実現。
// === 1 昇順、-1 降順 ===
db.products.aggregate([
{ $sort: { price: -1 } }
]);
// === 複数フィールドでのソート ===
db.products.aggregate([
{ $sort: { category: 1, price: -1 } }
]);
(5) $limit / $skip:ページネーション
$skipのパフォーマンス問題: $skipの値が大きいほどパフォーマンスが低下します—$skip(10000)はMongoDBに最初の10,000ドキュメントをスキャンして破棄させるため、クライアントには返却されなくてもスキャンとソートのオーバーヘッドが存在します。高度ページネーションの最適化:1. カーソルベースページネーション(skipの代わりに_id: {$gt: lastId}を使用、パフォーマンスは一定);2. 最大ページ番号を制限(例:50ページまでのみナビゲーション許可、超過時は検索範囲の絞り込みを促す);3. $sort + $limitによるTop Nモード(skipなし、最初のN件のみ取得)。オフセットページネーションはデータ量が少ない場合(< 1,000ページ)のみ許容されます。
// === $limit 結果数制限 ===
db.products.aggregate([
{ $sort: { price: -1 } },
{ $limit: 10 }
]);
// === $skip スキップ ===
db.products.aggregate([
{ $sort: { price: -1 } },
{ $skip: 20 },
{ $limit: 10 } // 21-30件目
]);
(6) $count:カウント
$countの使用例: $countは最もシンプルな集計ステージです—Nドキュメントを入力として受け取り、count値を含む単一ドキュメントを出力します。$group({_id: null, total: {$sum: 1}})と同等ですが、より簡潔な構文です。一般的な用途:1. フィルタ済みドキュメント数のカウント($match + $count);2. $facet内のサブパイプラインとして総数を取得(ページネーション用);3. $matchと組み合わせて条件付きカウント(例:「Electronics製品は何件か?」)。
// === シンプルなカウント ===
db.products.aggregate([
{ $match: { category: 'Electronics' } },
{ $count: 'totalElectronics' }
]);
// [ { totalElectronics: 250 } ]
// === countDocuments と同等 ===
db.products.countDocuments({ category: 'Electronics' });
5. アキュムレータ
概念説明: アキュムレータは$groupステージで使用される集計関数で、各ドキュメントグループに対して計算を行い、単一の結果を出力します。$sum合計、$avg平均、$min/$max極値、$first/$last最初と最後の値、$push/$addToSet配列蓄積。アキュムレータは統計分析の核となるツールです。
graph LR
A[グループ化フィールド _id] --> B[グループ1]
A --> C[グループ2]
A --> D[グループ3]
B --> E["$sum: {$sum:1}<br/>$avg: {$avg:'$price'}<br/>$push: {$push:'$name'}"]
C --> F["$sum: {$sum:1}<br/>$avg: {$avg:'$price'}<br/>$push: {$push:'$name'}"]
D --> G["$sum: {$sum:1}<br/>$avg: {$avg:'$price'}<br/>$push: {$push:'$name'}"]
E --> H[グループごとに1行出力<br/>集計結果]
F --> H
G --> H
style H fill:#d4edda
(1) アキュムレータ一覧
| アキュムレータ | 意味 | 例 |
|---|---|---|
$sum |
合計 | { $sum: '$price' } |
$avg |
平均 | { $avg: '$price' } |
$min |
最小値 | { $min: '$price' } |
$max |
最大値 | { $max: '$price' } |
$first |
最初の値 | { $first: '$name' } |
$last |
最後の値 | { $last: '$name' } |
$push |
配列蓄積 | { $push: '$name' } |
$addToSet |
配列に重複排除で蓄積 | { $addToSet: '$name' } |
$count |
カウント(トップレベル用) | { $count: 'total' } |
▶ サンプル 1:アキュムレータの実践(難易度 ⭐)
// === カテゴリごとの商品数と価格集計 ===
db.products.aggregate([
{
$group: {
_id: '$category',
count: { $sum: 1 },
totalStock: { $sum: '$stock' },
avgPrice: { $avg: '$price' },
maxPrice: { $max: '$price' },
minPrice: { $min: '$price' },
topProducts: { $push: '$title' } // 全商品名
}
}
]);
// === $addToSet 重複排除蓄積 ===
db.users.aggregate([
{
$group: {
_id: '$city',
uniqueRoles: { $addToSet: '$role' } // 各都市の役割セット
}
}
]);
6. 実行順序の最適化
概念説明: 集計パイプラインのパフォーマンスはステージ順序に大きく依存します。核心的な最適化原則:$matchはできるだけ早い段階に配置、$projectは$group後にフィールド数を削減、$sortと$matchは隣接させてインデックスを活用。MongoDBクエリオプティマイザは一部の最適化を自動実行しますが、ユーザー定義のフェーズ順序は変更しません。
graph LR
subgraph "✅ 最適な順序"
A1[$match<br/>先にフィルタ] --> B1[$sort<br/>ソート] --> C1[$limit<br/>制限]
end
subgraph "⚠️ アンチパターン"
A2[$sort<br/>全件ソート] --> B2[$limit<br/>制限] --> C2[$match<br/>後でフィルタ]
end
style A1 fill:#d4edda
style C2 fill:#f8d7da
| 最適化戦略 | 効果 | 適用条件 |
|---|---|---|
$matchを先頭に配置 |
後続データ量を削減 | フィルタフィールドにインデックスあり |
$match + $sortを隣接 |
インデックススキャンに統合 | ソートフィールドに複合インデックスあり |
$projectでフィールド削減 |
メモリ使用量削減 | $group後に使用 |
hint()でインデックス強制 |
全件スキャン回避 | オプティマイザが誤ったインデックスを選択時 |
// ✅ 最適化:$matchを先頭に配置
db.products.aggregate([
{ $match: { category: 'Electronics' } }, // 先にフィルタ(インデックス使用)
{ $sort: { price: -1 } },
{ $limit: 10 }
]);
// ⚠️ 反例:$matchを後ろに配置
db.products.aggregate([
{ $sort: { price: -1 } },
{ $limit: 10 },
{ $match: { category: 'Electronics' } } // 処理後にフィルタ、リソース無駄
]);
7. 総合実践
(1) Eコマース売上統計
// === 月次売上統計 ===
db.orders.aggregate([
{ $match: { status: 'paid', createdAt: { $gte: new Date('2026-01-01') } } },
{
$group: {
_id: {
year: { $year: '$createdAt' },
month: { $month: '$createdAt' }
},
totalSales: { $sum: '$total' },
orderCount: { $sum: 1 },
avgOrderValue: { $avg: '$total' }
}
},
{ $sort: { '_id.year': 1, '_id.month': 1 } }
]);
(2) 商品カテゴリ分析
// === 商品カテゴリ + 価格帯分析 ===
db.products.aggregate([
{
$bucket: {
groupBy: '$price',
boundaries: [0, 100, 500, 1000, 5000, 10000],
default: 'Other',
output: {
count: { $sum: 1 },
products: { $push: '$title' }
}
}
}
]);
▶ サンプル 2:ShopHub Eコマースカテゴリ分析パイプライン(難易度 ⭐)
// シナリオ:ShopHub運営チームが商品データを多角的に分析する必要がある
// データ準備
db.products.insertMany([
{ sku: 'PHONE-001', title: 'Smartphone X', category: 'Electronics', brand: 'TechCorp', price: 599, stock: 120, rating: 4.5 },
{ sku: 'PHONE-002', title: 'Smartphone Y', category: 'Electronics', brand: 'DataFlow', price: 399, stock: 80, rating: 4.0 },
{ sku: 'LAPTOP-001', title: 'Laptop Pro', category: 'Electronics', brand: 'TechCorp', price: 1299, stock: 50, rating: 4.8 },
{ sku: 'BOOK-001', title: 'MongoDB Guide', category: 'Books', brand: 'AppVenture', price: 29, stock: 500, rating: 4.2 },
{ sku: 'BOOK-002', title: 'Node.js Mastery', category: 'Books', brand: 'AppVenture', price: 39, stock: 300, rating: 4.6 }
]);
// 1. カテゴリ別統計:商品数、平均価格、最高価格、ブランドリスト
db.products.aggregate([
{
$group: {
_id: '$category',
count: { $sum: 1 },
avgPrice: { $avg: '$price' },
maxPrice: { $max: '$price' },
minPrice: { $min: '$price' },
totalStock: { $sum: '$stock' },
brands: { $addToSet: '$brand' }
}
},
{ $sort: { count: -1 } }
]);
// 2. 価格帯分布($bucket)
db.products.aggregate([
{
$bucket: {
groupBy: '$price',
boundaries: [0, 100, 500, 1000, 5000],
default: '5000+',
output: { count: { $sum: 1 }, titles: { $push: '$title' } }
}
}
]);
// 3. 高評価商品(rating >= 4.5)
db.products.aggregate([
{ $match: { rating: { $gte: 4.5 } } },
{ $project: { sku: 1, title: 1, price: 1, rating: 1, _id: 0 } },
{ $sort: { rating: -1, price: 1 } }
]);
出力: 1) カテゴリ別統計(Electronics: 3商品、平均価格765.67;Books: 2商品、平均価格34);2) 価格バケット(0-100: 2;500-1000: 1;1000-5000: 1);3) 高評価商品リスト。
▶ サンプル 3:Eコマース売上分析パイプライン(難易度 ⭐⭐)
// シナリオ:ShopHub月次売上ダッシュボード - 総売上、注文数、上位商品
db.orders.insertMany([
{ orderNumber: 'ORD-001', userId: 'user_001', items: [{ sku: 'PHONE-001', qty: 2, price: 599 }], total: 1198, status: 'paid', createdAt: new Date('2026-07-01') },
{ orderNumber: 'ORD-002', userId: 'user_002', items: [{ sku: 'LAPTOP-001', qty: 1, price: 1299 }], total: 1299, status: 'paid', createdAt: new Date('2026-07-05') },
{ orderNumber: 'ORD-003', userId: 'user_001', items: [{ sku: 'PHONE-001', qty: 1, price: 599 }, { sku: 'BOOK-001', qty: 2, price: 29 }], total: 657, status: 'paid', createdAt: new Date('2026-07-10') }
]);
// 完全パイプライン:フィルタ → グループ → ソート → 制限
const salesReport = db.orders.aggregate([
// ステージ1:2026年7月の支払済み注文のみフィルタ
{ $match: {
status: 'paid',
createdAt: {
$gte: new Date('2026-07-01'),
$lt: new Date('2026-08-01')
}
}},
// ステージ2:ユーザーごとにグループ化して合計を計算
{ $group: {
_id: '$userId',
totalSpent: { $sum: '$total' },
orderCount: { $sum: 1 },
avgOrderValue: { $avg: '$total' },
orders: { $push: '$orderNumber' }
}},
// ステージ3:支出に基づいて顧客ティアを計算
{ $addFields: {
tier: {
$switch: {
branches: [
{ case: { $gte: ['$totalSpent', 2000] }, then: 'Platinum' },
{ case: { $gte: ['$totalSpent', 1000] }, then: 'Gold' },
{ case: { $gte: ['$totalSpent', 500] }, then: 'Silver' }
],
default: 'Bronze'
}
}
}},
// ステージ4:総支出の降順でソート
{ $sort: { totalSpent: -1 } },
// ステージ5:上位10顧客に制限
{ $limit: 10 }
]);
// 結果出力
salesReport.forEach(r => {
console.log(`ユーザー: ${r._id}, 支出: $${r.totalSpent}, 注文数: ${r.orderCount}, ティア: ${r.tier}`);
});
出力:
TEXT 📖 参照専用ユーザー: user_001, 支出: $1855, 注文数: 2, ティア: Gold ユーザー: user_002, 支出: $1299, 注文数: 1, ティア: Silver
❓ よくある質問
findよりどれくらいパフォーマンスが悪い?findを、複雑な分析にはaggregateを使用してください。hint()オプションで特定インデックスの使用を強制できます。allowDiskUse: trueで一時的にディスクに書き込む必要があります。aggregateが返すカーソルを反復するには?cursor.hasNext()とcursor.next()を使用するか、mongoshで直接反復してください。📖 まとめ
- 集計パイプライン = 複数ステージを直列に通じてデータ処理
- 基本ステージ:$match、$group、$project、$sort、$limit、$skip、$count
- アキュムレータ:$sum、$avg、$min、$max、$first、$last、$push、$addToSet
$matchを先頭に配置して処理負荷を軽減- メモリ制限:100MB、
allowDiskUseで超過可能 - 集計パイプラインはインデックスを使用可能
📝 練習問題
- 基本問題(⭐):$groupを使用して各カテゴリの商品数をカウントし、平均価格を計算せよ。
- 基本問題(⭐):$match + $sort + $limitを使用してElectronics商品の価格上位10件をクエリせよ。
- 応用問題(⭐⭐):月ごとの注文総数と総売上を計算せよ($year/$month + $group)。
- 応用問題(⭐⭐):$projectを使用してdiscountedPrice(price × 0.9)を計算せよ。
- チャレンジ問題(⭐⭐⭐):完全な売上ダッシュボード(月、カテゴリ、顧客別の売上分析)を作成せよ。