Aggregation Pipeline: dây chuyền xử lý dữ liệu và cái giá của từng trạm
Các bài trước xoay quanh một câu hỏi: làm sao để server tìm đúng document nhanh nhất. Bài 07 dựng index, bài 08 ghép index, bài 10 xem planner chọn plan. Nhưng màn hình dashboard của một hệ thống SaaS không hỏi "đơn nào?". Nó hỏi "doanh thu mỗi ngày của từng shop trong tháng 9 là bao nhiêu?". Muốn trả lời thì phải tìm, rồi gom nhóm, cộng, sắp xếp, đôi khi nối sang collection khác. find() không làm được những việc đó. Aggregation pipeline thì làm được.
Viết một pipeline chạy ra kết quả đúng thì không khó. Khó là biết nó tốn bao nhiêu. Cùng một câu hỏi, trong lab của bài này, có cách viết chạy 2 ms và có cách chạy 2,8 giây. Có pipeline đang chạy 2,5 giây bỗng chậm lên 10 giây chỉ vì nó cần 117 MB RAM, mà giới hạn là 100 MB. Bài này đi tìm lý do.
Bài này nằm ở đâu
Cần biết trước : bài 06 (query là một dây chuyền stage, cursor và batch)
bài 07–08 (IXSCAN, FETCH, compound index, covered query)
bài 10 (đọc explain, classic engine và SBE)
Giới thiệu : pipeline như dây chuyền sản xuất, stage streaming vs blocking,
tối ưu hoá pipeline mà server tự làm, index chỉ giúp đoạn đầu,
giới hạn 100 MB, allowDiskUse và spill, giá của $lookup / $unwind /
$facet, $merge để tính trước, đọc explain của aggregation
Dẫn tới : bài 13, Query Performance & PaginationMôi trường lab: MongoDB 8.3.11 chạy trong Docker (image
mongo:8, standalone), containermongo-lab-09giới hạn 2 CPU, 3 GB RAM, WiredTiger cache 1 GB, máy host Apple M4. Database riênglab09. Lúc đo, lab của các bài khác cũng đang chạy trên cùng máy host, nên thời gian có dao động. Mỗi phép đo được chạy lặp lại, bài ghi median của từng lượt và chạy ít nhất hai lượt. Hãy so sánh các con số với nhau, đừng coi chúng là latency production. Con số nào không đo thì được ghi là minh hoạ.
Như các bài trước, bài dùng ba nhãn. [tài liệu] là hành vi tài liệu MongoDB mô tả. [quan sát] là thứ đo được trong lab nhưng không phải cam kết của MongoDB. [hình dung] là mô hình đơn giản hoá để dễ nhớ.
Giải thích trong 30 giây
Aggregation pipeline là một danh sách các stage. Document đi vào stage đầu tiên, mỗi stage biến đổi dòng document rồi đẩy sang stage sau. $match lọc, $group gom nhóm và cộng dồn, $sort sắp xếp, $lookup nối sang collection khác, $merge ghi kết quả ra một collection.
Có hai điều quyết định chi phí. Thứ nhất, chỉ đoạn đầu pipeline dùng được index. Từ chỗ dòng document đã bị biến đổi trở đi, server làm việc trên dữ liệu "trần", không còn mục lục nào. Thứ hai, có những stage phải gom hết input rồi mới trả ra được kết quả đầu tiên ($group, $sort). Các stage này tốn RAM, và mỗi stage bị giới hạn khoảng 100 MB. Vượt mức đó thì hoặc lỗi, hoặc ghi tạm ra đĩa (spill) và chậm đi.
Hình dung trước: dây chuyền đóng gói trong một xưởng
Hãy tưởng tượng một xưởng xử lý đơn hàng. Hàng từ kho chạy trên băng chuyền qua từng trạm:
Kho hàng Trạm 1 Trạm 2 Trạm 3 Trạm 4
(collection) ─► kiểm tra ─────► dán nhãn ngày ────► gom theo shop ─► xếp thứ tự ─► xe tải
loại hàng lỗi (thêm thông tin) vào từng thùng các thùng
($match) ($set) ($group) ($sort)Vài điều rút ra từ cái xưởng này, và điều nào cũng đúng với MongoDB:
- Trạm kiểm tra nên đặt sát cửa kho. Loại hàng lỗi càng sớm thì các trạm sau càng ít việc. Hơn nữa, chỉ ở cửa kho mới có sổ vị trí hàng (index) để lấy đúng thứ cần mà không phải khuân hết ra.
- Có trạm làm từng món một, có trạm phải chờ. Trạm dán nhãn nhận một món, dán, đẩy đi ngay. Trạm gom thùng thì không thể gửi thùng "shop 42" đi khi chưa chắc món cuối cùng của shop 42 đã tới. Nó phải chờ hết hàng. Trong lúc chờ, các thùng chiếm chỗ trên sàn.
- Sàn của mỗi trạm có diện tích giới hạn. Hết chỗ thì hoặc dừng xưởng (lỗi), hoặc mang bớt thùng ra kho tạm ngoài sân rồi lát nữa mang vào (spill). Đi ra sân thì chậm.
- Trạm "chạy sang kho bên cạnh lấy thêm phụ kiện" ($lookup) tốn công theo từng món đi qua. Một nghìn món là một nghìn chuyến.
[hình dung] Đây chỉ là cách hình dung. MongoDB thực tế không có băng chuyền chạy song song giữa các trạm. Các stage được gọi nối nhau trong cùng một luồng xử lý, và với slot-based engine (SBE), nhiều stage còn được gộp lại thành một plan duy nhất. Nhưng chi phí của từng loại trạm thì đúng như trong xưởng.
Mental model: nửa đầu và nửa sau của pipeline
Khi bạn gửi aggregate, server không chạy từng stage theo đúng thứ tự bạn viết. Trước tiên nó tối ưu lại pipeline, rồi tách pipeline thành hai phần:
db.orders.aggregate([ $match, $group, $sort, $lookup, ... ])
│
▼
┌──────── optimizer ────────┐
│ sắp lại, gộp, đẩy xuống │
└───────────┬───────────────┘
▼
┌───────────────────────────────────────┐
│ NỬA ĐẦU: query layer │ ← planner, index, IXSCAN / COLLSCAN
│ $match (+ $sort, $limit, $project...) │ với SBE: có thể gồm cả $group, $lookup
│ explain: stages[0].$cursor │
└───────────────────┬───────────────────┘
│ dòng document
▼
┌───────────────────────────────────────┐
│ NỬA SAU: các stage còn lại │ ← không còn index của collection đầu vào
│ $set → $group → $sort → $lookup → ... │
│ explain: stages[1..n] │
└───────────────────────────────────────┘[tài liệu] Nửa đầu chính là một query bình thường, giống find() ở bài 06–10. Planner chọn index cho nó như bài 10 mô tả. Trong explain, nó hiện thành stage $cursor. Nếu toàn bộ pipeline chạy được bằng SBE, explain không còn mảng stages mà đưa mọi thứ vào queryPlanner (bài 10 đã gặp trường hợp này với $group).
Thứ hai cần phân biệt là streaming và blocking:
STREAMING (nhận 1, trả 1 ngay) BLOCKING (phải thấy hết input)
$match $set/$project $unwind $group $sort (không có index)
$limit $skip $lookup $bucket $bucketAuto $facet
↓ ↓
RAM gần như không đổi RAM tăng theo số nhóm / số document
trả document đầu tiên sớm document đầu tiên chỉ ra khi xong hết[tài liệu] Trang $group gọi thẳng $group là "blocking stage": pipeline phải chờ stage này nhận hết input. Trang giới hạn aggregation nói các stage không thể xuất document nào trước khi xử lý hết input là những stage phải giữ dữ liệu trong RAM, và có thể cần nhiều hơn 100 MB. $sort có index hỗ trợ thì không blocking, vì index đã có sẵn thứ tự.
Bảng giá từng stage
Trước khi đo, đây là bảng tóm tắt để tra. Phần sau sẽ kiểm chứng từng dòng.
| Stage | Loại | Dùng index của collection đầu vào? | Tốn gì |
|---|---|---|---|
$match | streaming | Có, nếu nằm ở đầu (sau tối ưu hoá) | đọc document / key |
$project, $set, $unset | streaming | Không | CPU tính expression |
$sort | blocking nếu không có index | Có, nếu không có $project, $unwind, $group đứng trước | RAM theo tổng kích thước input; với $limit thì chỉ n phần tử |
$limit, $skip | streaming | gộp vào stage trước | gần như không |
$group | blocking | Chỉ trường hợp đặc biệt ($first/$last, DISTINCT_SCAN) | RAM theo số nhóm |
$unwind | streaming | Không | nhân số document lên theo độ dài mảng |
$lookup | streaming | Dùng index của collection bị nối | một lần tra cứu cho mỗi document đi qua |
$facet | blocking | Không (bên trong sub-pipeline) | RAM, kết quả là một document ≤ 16 MiB |
$bucket, $bucketAuto | blocking | Không | như $group |
$merge, $out | cuối pipeline | Không | ghi dữ liệu, có index phải cập nhật |
Cột "dùng index" lấy từ trang Aggregation Pipeline Optimization [tài liệu]: $match dùng được index nếu là stage đầu tiên "sau các tối ưu hoá"; $sort dùng được index nếu không có $project, $unwind hay $group đứng trước; $group chỉ dùng index khi pipeline sort và group theo cùng field và chỉ dùng $first/$last; $lookup, $graphLookup, $unionWith dùng được index của collection mà chúng đọc.
Setup và dataset
MongoDB version : 8.3.11 (Docker image mongo:8, standalone)
Hardware : Apple M4 host; container 2 CPU, 3 GB RAM
Configuration : WiredTiger cache 1 GB (--wiredTigerCacheSizeGB 1)
Dataset : lab09.orders 1.000.000 document, 221,7 MB chưa nén
(avgObjSize 221 byte), 74,7 MB trên đĩa
lab09.customers 100.000 document
lab09.tenants 500 document
Indexes : ban đầu chỉ có _id_, các index khác tạo dần trong bàiDữ liệu giống bài 06, có thêm customerId để nối sang customers. Có 500 tenant, mỗi tenant 200 khách. Trạng thái: 70% completed, 10% pending, 10% cancelled, 10% refunded. createdAt rải đều trong 365 ngày trước 2026-10-01. Mỗi đơn có 1–3 món hàng. Số được sinh từ PRNG có seed cố định.
// trích gen.js: 100 lần insertMany, mỗi lần 10.000 đơn
docs.push({
tenantId: "t" + String(t).padStart(4, "0"), // t0000 … t0499
customerId: "c" + String(t * 200 + u).padStart(6, "0"), // nối sang customers.customerId
status: STATUS(rnd()),
createdAt: new Date(END - Math.floor(rnd() * 365 * DAY)),
total: items.reduce((s, it) => s + it.qty * it.price, 0), // VND
items // [{ sku, qty, price }]
});Đo thời gian bằng một hàm nhỏ trong mongosh. Nó chạy một lần làm nóng, rồi chạy n lần và lấy median của thời gian gọi aggregate(...).toArray(). Con số này gồm cả phần mongosh nhận kết quả, không chỉ thời gian server.
function bench(label, fn, n) {
fn(); // làm nóng
const t = [];
for (let i = 0; i < n; i++) { const a = performance.now(); fn(); t.push(performance.now() - a); }
// in median, min, max
}Ví dụ chính: doanh thu theo tenant theo ngày
Câu hỏi và pipeline
"Doanh thu từ đơn completed của mỗi tenant, mỗi ngày trong tháng 9/2026, theo giờ Việt Nam."
const P = [
{ $match: {
status: "completed",
createdAt: { $gte: ISODate("2026-08-31T17:00:00Z"), // 00:00 ngày 1/9 giờ VN
$lt: ISODate("2026-09-30T17:00:00Z") } } // 00:00 ngày 1/10 giờ VN
},
{ $group: {
_id: { tenantId: "$tenantId",
day: { $dateTrunc: { date: "$createdAt", unit: "day",
timezone: "Asia/Ho_Chi_Minh" } } },
revenue: { $sum: "$total" },
orders: { $sum: 1 } } },
{ $sort: { "_id.tenantId": 1, "_id.day": 1 } }
];
db.orders.aggregate(P)Pipeline trả về 14.650 dòng (500 tenant × 30 ngày, trừ những ngày tenant không có đơn nào). Hai dòng đầu:
{ _id: { tenantId: 't0000', day: ISODate('2026-08-31T17:00:00.000Z') }, revenue: 9280000, orders: 5 }
{ _id: { tenantId: 't0000', day: ISODate('2026-09-01T17:00:00.000Z') }, revenue: 9170000, orders: 3 }Một chi tiết đúng/sai trước khi nói hiệu năng: ngày nào là "ngày 1/9"? BSON Date luôn là UTC (bài 03). Nếu bạn cắt theo ngày UTC, đơn lúc 6 giờ sáng 1/9 ở Hà Nội (23:00 UTC ngày 31/8) sẽ bị tính vào ngày 31/8. [tài liệu] $dateTrunc nhận timezone, nên ở đây mỗi "ngày" bắt đầu lúc 17:00 UTC. Đó là lý do day hiển thị ...T17:00:00Z. Khoảng $match cũng phải dùng cùng ranh giới đó.
Trong SQL, đây là một câu quen thuộc:
SELECT tenant_id,
date_trunc('day', created_at AT TIME ZONE 'Asia/Ho_Chi_Minh') AS day,
sum(total) AS revenue, count(*) AS orders
FROM orders
WHERE status = 'completed' AND created_at >= '2026-08-31 17:00Z' AND created_at < '2026-09-30 17:00Z'
GROUP BY 1, 2
ORDER BY 1, 2;$match ứng với WHERE, $group ứng với GROUP BY cùng các hàm gộp, $sort ứng với ORDER BY. Khác biệt lớn là SQL là ngôn ngữ khai báo: bạn mô tả kết quả, planner tự quyết thứ tự. Pipeline thì có thứ tự: bạn viết các bước, optimizer chỉ được phép sắp lại trong một số trường hợp an toàn (phần sau).
Lần 1: không có index
winningPlan: ← nằm trong stages[0].$cursor
GROUP
COLLSCAN
cursor stats: { nReturned: 14650, totalKeysExamined: 0, totalDocsExamined: 1000000 }
group: { usedDisk: false, spills: 0, peakTrackedMemBytes: 1742993 }
later stage: { stage: '$sort', nReturned: 14650, usedDisk: false, peakTrackedMemBytes: 9991300 }Đọc từ dưới lên:
COLLSCAN,totalDocsExamined: 1.000.000. Không có index nào chostatus/createdAt, nên server đọc cả collection để lấy 57.464 đơn khớp (số này đọc được ở stagefiltertrong explain đầy đủ).GROUPnằm trongwinningPlan.$groupđược đẩy xuống query layer và chạy bằng SBE. [tài liệu] Từ 5.2,$groupchạy bằng SBE nếu nó là stage đầu tiên hoặc mọi stage đứng trước nó cũng chạy được bằng SBE.$sortvẫn là một stage riêng ở nửa sau, sắp xếp 14.650 dòng, dùng khoảng 10 MB.peakTrackedMemBytescủa$groupkhoảng 1,7 MB. Đây là trường mới từ 8.3 [tài liệu], cho biết bộ nhớ đỉnh mà stage dùng. RAM của$groupđi theo số nhóm (14.650), không theo số document (1 triệu).
Explain đầy đủ còn một chi tiết đáng xem. Stage scan của SBE ghi:
scanFieldNames: [ 'createdAt', 'tenantId', 'total', 'status' ]Đây là projection optimization [tài liệu]: pipeline tự xác định nó chỉ cần 4 field, nên chỉ lấy 4 field đó ra khỏi mỗi document, không mang items hay customerId đi tiếp. Vì vậy, thêm $project ở đầu pipeline để "bớt field" thường không giúp gì. Tài liệu khuyên đặt $project ở cuối, để định hình kết quả trả về.
Lần 2 và 3: thêm index, rồi làm index "phủ" luôn
Bài 08 đã nói về ESR. Ở đây cả status (equality) và createdAt (range) nằm trong $match:
db.orders.createIndex({ status: 1, createdAt: 1 }) // 1.152 ms để build
db.orders.createIndex({ status: 1, createdAt: 1, tenantId: 1, total: 1 }) // index phủindex { status, createdAt } index { status, createdAt, tenantId, total }
GROUP GROUP
FETCH PROJECTION_COVERED
IXSCAN status_1_createdAt_1 IXSCAN status_1_createdAt_1_tenantId_1_total_1
keys 57.464, docs 57.464 keys 57.464, docs 0Index thứ hai chứa đủ cả 4 field mà pipeline cần, nên server không phải mở document nào (totalDocsExamined: 0). Đây là covered query của bài 08, áp dụng cho aggregation.
Cuối cùng là câu hỏi mà màn hình dashboard của một tenant thực sự hỏi. Thêm tenantId: "t0042" vào $match và dùng index { tenantId: 1, status: 1, createdAt: 1 }:
GROUP
FETCH
IXSCAN tenantId_1_status_1_createdAt_1
keys 126, docs 126, nReturned 30Kết quả đo
| Cách chạy | Đọc gì | median lượt 1 | median lượt 2 |
|---|---|---|---|
| Không index (COLLSCAN) | 1.000.000 document | 422 ms | 530 ms |
Index {status, createdAt} | 57.464 key + 57.464 document | 281 ms | 267 ms |
Index phủ {status, createdAt, tenantId, total} | 57.464 key, 0 document | 126 ms | 165 ms |
Một tenant, index {tenantId, status, createdAt} | 126 key + 126 document | 3,3 ms | 1,7 ms |
Ba dòng đầu trả lời cùng một câu hỏi cho cả 500 tenant. Dòng cuối trả lời câu hỏi nhỏ hơn mà người dùng thường hỏi. Một lượt đo khác, ép plan bằng hint, cho kết quả cùng xu hướng: COLLSCAN 433 / 523 ms, FETCH 237 / 197 ms.
BEFORE AFTER (index) AFTER (index phủ)
COLLSCAN IXSCAN IXSCAN
↓ 1.000.000 document ↓ 57.464 key ↓ 57.464 key
filter FETCH (không mở document)
↓ 57.464 ↓ 57.464 document ↓
GROUP → 14.650 nhóm GROUP → 14.650 nhóm GROUP → 14.650 nhómHai điều đáng chú ý:
- Index giúp ít hơn bạn nghĩ khi câu hỏi phủ 5,7% collection. Từ COLLSCAN sang IXSCAN + FETCH chỉ nhanh khoảng 1,5–2 lần, vì 57.464 lần FETCH là 57.464 lần nhảy tới các document nằm rải rác. Bài 07 đã gặp hiện tượng này: index mạnh nhất khi lọc còn ít.
- Index phủ nhanh hơn nữa, nhưng có giá. Trong lab, nó nặng 24,2 MB, gấp đôi index
{status, createdAt}(11,9 MB).
Index phủ { status, createdAt, tenantId, total }
│
├── ✓ aggregation cả tháng không mở document nào
├── ✗ 24,2 MB RAM/đĩa nữa (collection chỉ có 74,7 MB trên đĩa)
├── ✗ mỗi insert, mỗi lần đổi status hay total phải sửa thêm một index
└── ⚠ đổi pipeline (cần thêm field) là hết "phủ"Còn một lựa chọn rẻ hơn cả hai: đừng tính lại mỗi lần. Phần $merge ở cuối bài sẽ đo nó.
Server tự sắp xếp lại pipeline
30 giây
Trước khi chạy, optimizer viết lại pipeline theo một số luật cố định [tài liệu]. Mục tiêu luôn là: lọc càng sớm càng tốt, và đưa được $match/$sort lên đầu để dùng index. Tài liệu nói rõ các tối ưu hoá này "có thể thay đổi giữa các phiên bản", và cách xem pipeline sau khi tối ưu là dùng explain.
Những luật chính theo trang Aggregation Pipeline Optimization:
| Luật | Trước | Sau |
|---|---|---|
$match sau $project/$set/$addFields | $set → $match | phần filter không phụ thuộc field tính ra được tách lên trước |
$match sau $sort | $sort → $match | $match → $sort |
$skip sau $project | $project → $skip | $skip → $project |
$sort + $limit | hai stage | một $sort chỉ giữ top-n |
$sort + $skip + $limit | ba stage | $sort giữ top-(skip+limit), rồi $skip |
$limit + $limit, $skip + $skip | hai stage | một stage (min / tổng) |
$match + $match | hai stage | một $match với $and |
$lookup + $unwind (+ $match) | hai, ba stage | $lookup tự unwind (và lọc) bên trong |
Thí nghiệm A: viết ngược thứ tự
Giả sử ai đó viết pipeline lấy 10 đơn pending mới nhất của tenant t0042, kèm ngày giờ VN, theo đúng thứ tự họ nghĩ ra:
const A = [
{ $set: { day: { $dateTrunc: { date: "$createdAt", unit: "day", timezone: "Asia/Ho_Chi_Minh" } } } },
{ $sort: { createdAt: -1 } },
{ $match: { tenantId: "t0042", status: "pending" } },
{ $limit: 10 }
];Đọc nguyên văn thì đây là thảm hoạ: tính day cho 1 triệu đơn, sắp xếp cả 1 triệu đơn, rồi mới lọc. Explain cho thấy server làm khác:
stages[0].$cursor parsedQuery: { $and: [ { status: 'pending' }, { tenantId: 't0042' } ] }
FETCH
IXSCAN tenantId_1_status_1_createdAt_1 keys 203, docs 203
stages[1] $set nReturned 203
stages[2] $sort { createdAt: -1 }, limit: 10 nReturned 10, peakTrackedMemBytes 4686$matchđược kéo lên trước cả$sortlẫn$set, và vì đã nằm ở đầu nên dùng được index.$limitđược gộp vào$sort(limit: 10). Stage sort chỉ giữ 10 phần tử tốt nhất, chỉ tốn 4,7 KB.
Nhưng có một việc optimizer không làm: nó không đưa $set xuống sau $sort. Vì $set vẫn nằm trước $sort, $sort không dùng được index (luật "không có $project đứng trước"). Server phải lấy đủ 203 đơn pending của tenant rồi sắp trong RAM. Viết lại theo đúng thứ tự:
const A2 = [
{ $match: { tenantId: "t0042", status: "pending" } },
{ $sort: { createdAt: -1 } },
{ $limit: 10 },
{ $set: { day: { $dateTrunc: { date: "$createdAt", unit: "day", timezone: "Asia/Ho_Chi_Minh" } } } }
];stages[0].$cursor
LIMIT 10
FETCH
IXSCAN tenantId_1_status_1_createdAt_1 keys 10, docs 10 ← không còn SORT
stages[1] $set nReturned 10Index {tenantId, status, createdAt} đã có thứ tự createdAt bên trong mỗi cặp (tenant, status). Server đi ngược index 10 bước là xong, $set chỉ tính cho 10 document.
| Pipeline | keys / docs | median lượt 1 | median lượt 2 |
|---|---|---|---|
| A (viết ngược) | 203 / 203 + sort trong RAM | 1,5 ms | 1,1 ms |
| A2 (đúng thứ tự) | 10 / 10 | 0,7 ms | 0,6 ms |
Chênh 1 ms nghe không đáng kể. Nhưng 203 là số đơn pending của một tenant nhỏ. Với một tenant có 200.000 đơn pending, A phải đọc và sắp 200.000 document, còn A2 vẫn đọc 10. Optimizer cứu bạn được phần lớn, không phải toàn bộ.
Thí nghiệm B: $match trên field vừa tính ra
[ { $set: { day: { $dateTrunc: { ... } } } },
{ $match: { day: ISODate("2026-09-14T17:00:00Z"), tenantId: "t0042" } } ]stages[0].$cursor parsedQuery: { tenantId: 't0042' } keys 2022, docs 2022
stages[1] $set nReturned 2022
stages[2] $match { day: ... } nReturned 7[tài liệu] $match bị tách đôi. Phần tenantId không phụ thuộc $set nên được đưa lên đầu và dùng index. Phần day phải chờ $set tính xong. Kết quả: server đọc 2.022 đơn của tenant để giữ lại 7. Nếu bạn viết điều kiện thẳng trên createdAt ($gte/$lt theo ranh giới ngày VN), cả điều kiện sẽ nằm trong index.
Phiên bản: trang tối ưu hoá hiện tại (9.0) có thêm một luật mới: đưa cả phép tính lẫn
$matchlên trước những stage đắt như$lookup. Luật đó ghi "new in 9.0", không áp dụng cho 8.3.11 trong lab này.
Thí nghiệm C và D: lọc sau $group
Trong SQL có HAVING. Trong pipeline, đó là $match đặt sau $group. Có hai trường hợp rất khác nhau.
// C: lọc theo khoá nhóm
[ { $group: { _id: { tenantId: "$tenantId", status: "$status" }, revenue: { $sum: "$total" } } },
{ $match: { "_id.tenantId": "t0042" } } ]
// D: lọc theo giá trị đã cộng dồn
[ { $group: { _id: "$tenantId", revenue: { $sum: "$total" } } },
{ $match: { revenue: { $gt: 9e9 } } },
{ $sort: { revenue: -1 } } ]C: GROUP D: stages[0] GROUP
FETCH COLLSCAN docs 1.000.000
IXSCAN tenantId_1_status_1_createdAt_1 stages[1] $match revenue nReturned 2
keys 2022, docs 2022, nReturned 4 stages[2] $sort- D là điều hiển nhiên: không ai biết
revenuecủa một tenant trước khi cộng xong, nên$matchnày buộc phải chờ$groupduyệt hết 1 triệu đơn. - C thú vị hơn. [quan sát] Trên 8.3.11, server nhận ra
_id.tenantIdchính làtenantIdcủa input, nên đưa điều kiện lên trước$groupvà dùng index: chỉ 2.022 key. Trang tối ưu hoá mà tôi đọc không liệt kê luật này. Vì vậy đừng dựa vào nó: hãy tự đặt$match { tenantId }ở đầu. Viết vậy thì rõ ràng hơn, và không phụ thuộc phiên bản.
Index chỉ giúp được đoạn đầu dây chuyền
Gom các thí nghiệm lại thì được một bức tranh:
┌────────── vùng index dùng được ──────────┐
pipeline: $match → $sort → $limit → $set → $group → $sort → $lookup
IXSCAN theo gộp ─────── từ đây: dữ liệu "trần" ───────
index vào sort không index, không thứ tự sẵn
↑
dùng index của collection KHÁCBa hệ quả thực tế:
$sortsau$groupluôn là sort trong RAM. Kết quả của$grouplà dữ liệu mới, không có index nào. [tài liệu]$groupcũng không đảm bảo thứ tự output. Muốn có thứ tự thì phải có$sortsau nó. Tin tốt: số nhóm thường nhỏ hơn rất nhiều so với số document.- Mọi thứ bạn lọc được bằng field gốc nên nằm trong
$matchđầu tiên, kể cả khi bạn sẽ lọc lại sau đó. $lookuplà ngoại lệ duy nhất "ở giữa dây chuyền" dùng được index, nhưng là index của collection bị nối, không phải collection đầu vào.
Giới hạn 100 MB và chuyện tràn ra đĩa
30 giây
[tài liệu] Mỗi stage cần giữ dữ liệu trong RAM bị giới hạn 100 MB. Từ MongoDB 6.0, tham số server allowDiskUseByDefault quyết định chuyện gì xảy ra khi vượt mức. Nếu nó là true, stage ghi file tạm ra đĩa. Nếu false, stage báo lỗi. Bạn có thể đổi cho từng lệnh bằng option allowDiskUse. Các stage có thể spill gồm $group, $sort (khi không có index hỗ trợ), $bucket, $bucketAuto, $setWindowFields, $sortByCount. Riêng $facet thì không spill được, vượt 100 MB là lỗi.
[quan sát] Trong lab, allowDiskUseByDefault là true. Hai tham số nội bộ internalQueryMaxBlockingSortMemoryUsageBytes và internalDocumentSourceGroupMaxMemoryBytes đều là 104857600 (100 MiB).
Quay lại cái xưởng: spill là khi sàn của trạm hết chỗ, công nhân mang một phần thùng ra kho tạm ngoài sân rồi lát nữa mang vào ghép lại. Xưởng không dừng, nhưng mỗi chuyến ra sân là thêm thời gian.
Thí nghiệm 1: sort cả collection, và sort top-k
// SORTALL: sắp 1 triệu đơn theo total, bỏ 999.990 cái đầu → phải giữ cả document
[ { $sort: { total: -1, createdAt: 1 } }, { $skip: 999990 } ]
// TOPK: 10 đơn lớn nhất
[ { $sort: { total: -1, createdAt: 1 } }, { $limit: 10 } ]Chặn spill để xem giới hạn:
db.orders.aggregate(SORTALL, { allowDiskUse: false })QueryExceededMemoryLimitNoDiskUseAllowed: ... Sort exceeded memory limit of 104857600 bytes,
but did not opt in to external sorting.Với mặc định (cho phép spill), explain("executionStats"):
SORTALL TOPK
SKIP SORT limit=10
SORT { total: -1, createdAt: 1 } COLLSCAN
COLLSCAN
usedDisk: true usedDisk: false
spills: 4 spills: 0
spilledRecords: 1.000.000 peakTrackedMemBytes: 3.720
spilledBytes: 240.685.951
spilledDataStorageSize: 62.769.246
peakTrackedMemBytes: 103.808.946Đọc các trường spill [tài liệu, các trường có từ 8.1–8.3]:
usedDisk: true,spills: 4: stage ghi ra đĩa 4 lần, mỗi lần bộ nhớ chạm khoảng 100 MB (peakTrackedMemBytes103,8 MB).spilledByteslà ước lượng số byte ghi ra trước khi nén: 240,7 MB, tức gần như toàn bộ 1 triệu document.spilledDataStorageSizelà dung lượng đĩa thật sự dùng: 62,8 MB. Dữ liệu spill được nén khoảng 3,8 lần.- TOPK đọc đúng 1 triệu document như SORTALL, nhưng chỉ giữ 3.720 byte. [tài liệu] Khi
$limitđược gộp vào$sort, sort chỉ cần giữ n phần tử tốt nhất trong lúc chạy.
| Pipeline | median lượt 1 | median lượt 2 |
|---|---|---|
| SORTALL, giới hạn 100 MB (spill) | 2.881 ms | 2.747 ms |
| SORTALL, giới hạn sort nâng lên 1 GB (không spill) | 1.936 ms | 2.343 ms |
| TOPK (sort + limit 10) | 288 ms | 465 ms |
Để thấy giá của spill, tôi đổi tham số nội bộ internalQueryMaxBlockingSortMemoryUsageBytes cho cùng pipeline. Ở mức 400 MB, explain cho thấy sort không spill và dùng 325,7 MB. Spill làm SORTALL chậm hơn khoảng 15–50% trong lab này. Nhưng chênh lệch lớn nhất nằm ở chỗ khác: TOPK nhanh hơn SORTALL 6–10 lần, chỉ nhờ một $limit.
Một sự cố thật trong lab. Trong một lượt đo với giới hạn sort nâng lên 1 GB, container 3 GB bị kernel kill vì hết bộ nhớ (
OOMKilled: true), mất kết nối giữa chừng. Tôi không xác định được chính xác stage nào đẩy nó qua ngưỡng. Nhưng bài học đủ rõ: giới hạn 100 MB tồn tại để một câu query không giết cả server. Tham sốinternal*là công cụ cho thí nghiệm, đừng chỉnh trên production.
Thí nghiệm 2: $group với gần 1 triệu nhóm
RAM của $group đi theo số nhóm. Gom theo khách × ngày trên cả năm thì gần như mỗi đơn là một nhóm:
const GRP = [
{ $group: { _id: { c: "$customerId",
d: { $dateTrunc: { date: "$createdAt", unit: "day", timezone: "Asia/Ho_Chi_Minh" } } },
revenue: { $sum: "$total" }, n: { $sum: 1 } } },
{ $count: "groups" } // → 986.470 nhóm
];allowDiskUse: false → Exceeded memory limit for $group, but didn't allow external spilling;
pass allowDiskUse:true to opt in
mặc định:
GROUP ← $count, được SBE dịch thành một GROUP nữa
GROUP
COLLSCAN
group: { usedDisk: true, spills: 2, spilledRecords: 989.056,
spilledBytes: 49.452.800, spilledDataStorageSize: 22.372.352,
peakTrackedMemBytes: 104.857.802 }Nhóm này chạy bằng SBE. [quan sát] Ngưỡng spill của $group trong SBE là một tham số nội bộ khác, internalQuerySlotBasedExecutionHashAggApproxMemoryUseInBytesBeforeSpill, mặc định cũng 100 MiB. Nâng riêng nó lên 500 MB thì pipeline không spill, và explain cho biết nó thực sự cần 117,3 MB, chỉ hơn giới hạn 17%.
| GRP (986.470 nhóm) | median lượt 1 | median lượt 2 |
|---|---|---|
| giới hạn 100 MB → spill 2 lần | 10.501 ms | 10.248 ms |
| giới hạn 500 MB → không spill | 2.708 ms | 2.507 ms |
Vượt giới hạn 17% làm pipeline chậm khoảng 4 lần. [quan sát] Lượng dữ liệu ghi ra không lớn (22 MB trên đĩa), nên phần chậm khó có thể chỉ do ghi đĩa. [hình dung] Cách giải thích hợp lý nhất (tôi chưa kiểm chứng trong mã nguồn): khi spill, các nhóm bị chia thành nhiều phần, cuối cùng server phải đọc lại các phần đó và gộp những nhóm trùng khoá. Với gần 1 triệu nhóm, phần gộp lại này rất đắt.
Đây là kiểu sự cố khó chịu nhất trên production. Pipeline chạy ổn nhiều tháng, rồi dữ liệu tăng thêm một chút là thời gian nhảy vọt, dù không ai sửa code.
Khi pipeline chạm giới hạn: sửa thế nào
Theo thứ tự nên thử:
- Lọc sớm hơn.
$matchtheo thời gian, theo tenant. Ít document đi vào$groupthì ít nhóm hơn. - Giảm số nhóm hoặc kích thước mỗi nhóm. Có thật sự cần khách × ngày trên cả năm? Tránh
$push: "$$ROOT"hay$addToSettrên mảng lớn trong$group, vì mỗi nhóm sẽ phình theo số document. - Thêm
$limitngay sau$sortkhi chỉ cần top-n. - Có index cho
$sortở đầu pipeline để sort không còn blocking. - Tính trước bằng
$merge(phần cuối bài). - Chỉ sau đó mới chấp nhận spill như một trạng thái bình thường, và theo dõi nó: slow query log và profiler có ghi
usedDisk.
allowDiskUse: true không phải cách sửa. Nó chỉ là van an toàn để query không lỗi.
$lookup: mỗi document là một chuyến sang kho bên cạnh
Ôn lại và mở rộng bài 04
Bài 04 đã đo $lookup cho một đơn: không có index trên foreignField thì server quét gần 500.000 dòng hàng, có index thì chỉ 7 key. Ở đây ta xem ba chuyện bài 04 chưa đo: chi phí nhân lên theo số document, ba "chiến lược" join mà server chọn, và vị trí đặt $lookup trong pipeline.
Thí nghiệm: doanh thu tháng 9 của tenant t0042 theo phân khúc khách
const L = [
{ $match: { tenantId: "t0042", status: "completed", createdAt: SEPT } }, // 126 đơn
{ $lookup: { from: "customers", localField: "customerId",
foreignField: "customerId", as: "customer" } },
{ $unwind: "$customer" },
{ $group: { _id: "$customer.segment", revenue: { $sum: "$total" }, orders: { $sum: 1 } } }
];Explain (verbosity queryPlanner) cho thấy luật gộp $lookup + $unwind:
{ '$lookup': { from: 'customers', as: 'customer', localField: 'customerId',
foreignField: 'customerId',
unwinding: { preserveNullAndEmptyArrays: false } } }[tài liệu] $unwind đã được nhập vào $lookup (trường unwinding), để không phải dựng mảng customer rồi tách ra.
customers.customerId | $lookup stage trong explain | median lượt 1 | median lượt 2 |
|---|---|---|---|
| không có index | totalDocsExamined: 12.700.000, collectionScans: 127 | 2.137 ms | 2.832 ms |
có index customerId_1 | totalDocsExamined: 127, totalKeysExamined: 127, indexesUsed: ['customerId_1'] | 2,2 ms | 2,0 ms |
Không có index: mỗi đơn trong 126 đơn kéo theo một lần quét cả 100.000 khách. 127 lần quét × 100.000 = 12,7 triệu document để trả lời một câu hỏi về 126 đơn. Thêm index thì nhanh hơn khoảng 1.000 lần.
BEFORE (không index) AFTER (index customerId_1)
126 đơn 126 đơn
↓ mỗi đơn: ↓ mỗi đơn:
quét 100.000 customers seek 1 key → 1 document
↓ ↓
12,7 triệu document 127 document[quan sát] Pipeline này chạy $lookup bằng classic engine (nó hiện như một stage riêng $lookup ở nửa sau). Tài liệu không nói rõ lý do. Khi bỏ $unwind ra (lấy segment bằng $first: "$customer.segment" trong $group), explain hiện EQ_LOOKUP customerId_1 strategy=IndexedLoopJoin của SBE như bài 04. Vậy nhiều khả năng chính việc gộp $unwind vào $lookup đã đưa stage này về classic engine.
Ba chiến lược join
Khi $lookup chạy bằng SBE (stage EQ_LOOKUP), explain ghi thêm strategy. Bài 04 đã gặp hai chiến lược. Lab này gặp chiến lược thứ ba khi nối 126 đơn sang tenants qua một field không có index (tenantKey):
EQ_LOOKUP strategy=HashJoin from=lab09.tenants
FETCH
IXSCAN tenantId_1_status_1_createdAt_1
totalDocsExamined: 626 ← 126 đơn + 500 tenant, đọc tenants đúng MỘT lần
hash_lookup: { usedDisk: false, peakTrackedMemBytes: 58386 }| Strategy | Khi nào (quan sát) | Chi phí |
|---|---|---|
IndexedLoopJoin | có index trên foreignField | mỗi document: một lần seek index |
HashJoin | không có index, collection bị nối nhỏ | đọc collection bị nối một lần, dựng bảng băm trong RAM |
NestedLoopJoin | không có index, collection bị nối lớn | mỗi document: quét cả collection bị nối |
[quan sát] Trang tài liệu $lookup không mô tả các strategy này. Trong lab, các tham số nội bộ internalQueryCollectionMaxNoOfDocumentsToChooseHashJoin (10.000 document) và internalQueryCollectionMaxDataSizeBytesToChooseHashJoin (100 MB) cho thấy giới hạn mà server dùng khi cân nhắc HashJoin. customers có 100.000 document, vượt ngưỡng này, nên rơi vào nested loop. Đây là chi tiết cài đặt, có thể thay đổi giữa các phiên bản. Điều tài liệu cam kết chỉ là: $lookup so khớp bằng nhau chạy tốt hơn khi có index trên foreignField, và "nhiều khả năng có hiệu năng kém" nếu không có.
Đặt $lookup ở đâu: trước hay sau $group?
Báo cáo "top 20 tenant theo doanh thu tháng 9, kèm tên và gói dịch vụ". Tên nằm ở tenants. Có hai cách viết:
// BEFORE: nối từng đơn, rồi mới gom
[ M, { $lookup: { from: "tenants", localField: "tenantId", foreignField: "_id", as: "t" } },
{ $unwind: "$t" },
{ $group: { _id: { tenantId: "$tenantId", name: "$t.name", plan: "$t.plan" }, revenue: { $sum: "$total" } } },
{ $sort: { revenue: -1 } }, { $limit: 20 } ]
// AFTER: gom trước, cắt top 20, rồi mới nối
[ M, { $group: { _id: "$tenantId", revenue: { $sum: "$total" } } },
{ $sort: { revenue: -1 } }, { $limit: 20 },
{ $lookup: { from: "tenants", localField: "_id", foreignField: "_id", as: "t" } },
{ $unwind: "$t" }, { $project: { revenue: 1, name: "$t.name", plan: "$t.plan" } } ]Cả hai đều join qua _id (luôn có index), nên mỗi lần tra đều rẻ. Khác nhau là số lần tra:
| Pipeline | số lần $lookup tra | median lượt 1 | median lượt 2 |
|---|---|---|---|
| BEFORE | 57.464 | 494 ms | 428 ms |
AFTER, $lookup sau $sort+$limit | 20 | 22,8 ms | 21,7 ms |
AFTER, $lookup trước $sort+$limit | 500 | 28,5 ms | 21,9 ms |
Nhanh hơn khoảng 20 lần, và không cần index mới. Có thêm một lợi ích: ở bản AFTER, $group nằm ngay sau $match nên được đẩy xuống SBE và dùng index phủ (PROJECTION_COVERED, 0 document). Ở bản BEFORE, $lookup + $unwind chen vào giữa nên $group chạy ở nửa sau.
Nguyên tắc: $lookup tốn công theo số document đi qua nó, nên hãy đặt nó ở chỗ dòng document hẹp nhất. Phần lớn trường hợp là sau $match, sau $group, sau $limit.
$unwind: tách mảng thành nhiều document
$unwind biến một đơn có 3 món thành 3 document, mỗi document mang một món. Nó streaming và rẻ cho từng document. Nhưng nó nhân số document cho mọi stage phía sau.
Thí nghiệm nhỏ: tổng số lượng hàng bán được của từng tenant trong tháng 9, hai cách viết cho cùng kết quả (đã kiểm tra trùng khớp):
// U1: tách mảng rồi cộng
[ M, { $unwind: "$items" }, { $group: { _id: "$tenantId", qty: { $sum: "$items.qty" } } } ]
// U2: cộng ngay trên mảng bằng expression
[ M, { $group: { _id: "$tenantId", qty: { $sum: { $sum: "$items.qty" } } } } ]U1: $cursor 57.464 đơn → $unwind 114.833 document → $group 500 nhóm
U2: $cursor 57.464 đơn → $group 500 nhóm| median lượt 1 | median lượt 2 | |
|---|---|---|
U1 $unwind + $group | 135 ms | 119 ms |
| U2 expression trên mảng | 122 ms | 122 ms |
[quan sát] Với trung bình 2 món mỗi đơn, chênh lệch nhỏ đến mức lẫn trong dao động đo. Tôi ghi lại kết quả này vì nó đi ngược lời khuyên "đừng bao giờ $unwind" hay gặp. $unwind không đắt tự thân. Nó đắt khi mảng dài (một document 1.000 phần tử thành 1.000 document), khi theo sau là $lookup hay $sort trên dòng đã bị nhân lên, hoặc khi bạn $unwind rồi $group theo _id chỉ để "ráp lại" document. Trong trường hợp cuối, các array expression ($sum, $filter, $reduce trên mảng) làm được cùng việc mà không tách document.
Còn khi bạn thật sự cần đơn vị là từng món hàng (doanh thu theo SKU), thì $unwind là cách đúng.
$facet và $bucket: nhiều báo cáo từ một lần đọc
30 giây
Một trang dashboard thường cần nhiều con số trên cùng một tập dữ liệu: số đơn theo trạng thái, doanh thu theo ngày, phân bố giá trị đơn, top khách. $facet chạy nhiều sub-pipeline trên cùng một input và trả về một document chứa tất cả kết quả. $bucket gom giá trị vào các khoảng (histogram), giống $group nhưng theo khoảng.
const subs = {
byStatus: [ { $group: { _id: "$status", n: { $sum: 1 } } } ],
byDay: [ { $match: { status: "completed" } },
{ $group: { _id: { $dateTrunc: { date: "$createdAt", unit: "day", timezone: "Asia/Ho_Chi_Minh" } },
revenue: { $sum: "$total" } } }, { $sort: { _id: 1 } } ],
totals: [ { $bucket: { groupBy: "$total",
boundaries: [0, 500000, 1000000, 2000000, 5000000, 1e12],
default: "other", output: { n: { $sum: 1 } } } } ],
topCustomers: [ { $match: { status: "completed" } },
{ $group: { _id: "$customerId", spent: { $sum: "$total" } } },
{ $sort: { spent: -1 } }, { $limit: 5 } ]
};
const F1 = [ { $match: { tenantId: "t0042", createdAt: SEPT } }, { $facet: subs } ];Kết quả là một document 1.537 byte, ví dụ:
byStatus: [{"_id":"refunded","n":14},{"_id":"completed","n":126},{"_id":"cancelled","n":17},{"_id":"pending","n":17}]
totals: [{"_id":0,"n":8},{"_id":500000,"n":12},{"_id":1000000,"n":29},{"_id":2000000,"n":62},{"_id":5000000,"n":63}]
byDay: 30 dòng$facet đứng ở đâu quyết định tất cả
[tài liệu] Nếu $facet là stage đầu tiên, nó luôn COLLSCAN. Các sub-pipeline không dùng được index của collection. Chỉ những stage đứng trước $facet ($match, $sort) mới dùng được index. Để thấy rõ, ta so sánh ba cách:
- F1:
$matchcó index, rồi$facet. - F2: bốn aggregation riêng, mỗi cái có
$matchriêng. - F3:
$facetđứng đầu,$matchnằm trong từng sub-pipeline.
F1 $cursor: FETCH ← IXSCAN tenantId_1_status_1_createdAt_1 keys 180, docs 174
$facet nReturned 1
F3 $cursor: COLLSCAN docs 1.000.000, nReturned 1.000.000
$facet nReturned 1| median lượt 1 | median lượt 2 | |
|---|---|---|
F1 $match → $facet | 1,3 ms | 1,6 ms |
| F2 bốn aggregation riêng | 3,2 ms | 3,8 ms |
F3 $facet đứng đầu | 1.225 ms | 916 ms |
F3 đưa toàn bộ 1 triệu đơn vào cả bốn sub-pipeline, chỉ để mỗi cái lọc lại còn 174. Chậm hơn F1 khoảng 600–900 lần. F1 nhanh hơn F2 vì đọc 174 đơn một lần thay vì bốn lần, và chỉ tốn một round-trip.
Các giới hạn cần nhớ [tài liệu]:
$facet
│
├── ✓ một lần đọc, một round-trip, nhiều báo cáo
├── ✗ sub-pipeline không dùng được index → lọc chung phải nằm TRƯỚC $facet
├── ✗ output là MỘT document ≤ 16 MiB → đừng trả danh sách dài trong facet
├── ✗ mỗi kết quả trung gian ≤ 100 MB, và $facet KHÔNG spill (allowDiskUse không giúp)
└── ⚠ không dùng được $out, $merge, $facet lồng nhau, $search, $vectorSearch... bên trongKhi các báo cáo cần những tập dữ liệu khác nhau (khoảng thời gian khác, tenant khác), các aggregation riêng thường tốt hơn, vì mỗi cái có $match và index riêng.
$merge và $out: tính một lần, đọc nhiều lần
30 giây
Câu hỏi "doanh thu từng ngày của tenant t0042 trong tháng 9" có một điểm đặc biệt: doanh thu của ngày 14/9 không thay đổi sau khi ngày 14/9 kết thúc (bỏ qua chuyện hoàn tiền về sau). Vậy tại sao mỗi lần người dùng mở dashboard, ta lại cộng lại từ đầu?
$merge ghi kết quả của pipeline vào một collection. Nó phải là stage cuối cùng. Đây chính là materialized view theo kiểu bạn tự quản lý.
Dựng bảng daily_revenue
db.daily_revenue.createIndex({ tenantId: 1, day: 1 }, { unique: true });
const ROLL = (from, to) => [
{ $match: { status: "completed", createdAt: { $gte: from, $lt: to } } },
{ $group: { _id: { tenantId: "$tenantId",
day: { $dateTrunc: { date: "$createdAt", unit: "day", timezone: "Asia/Ho_Chi_Minh" } } },
revenue: { $sum: "$total" }, orders: { $sum: 1 } } },
{ $project: { _id: 0, tenantId: "$_id.tenantId", day: "$_id.day",
revenue: 1, orders: 1, refreshedAt: "$$NOW" } },
{ $merge: { into: "daily_revenue", on: ["tenantId", "day"],
whenMatched: "replace", whenNotMatched: "insert" } }
];[tài liệu] $merge cần một unique index đúng trên các field on (ở đây tenantId, day). Với collection đã tồn tại, index đó phải có sẵn từ trước. whenMatched: "replace" thay dòng cũ bằng dòng mới tính. whenNotMatched: "insert" thêm dòng chưa có. Nếu lệnh lỗi giữa chừng, những gì đã ghi không được rollback.
| Thao tác | Đọc gì | median lượt 1 | median lượt 2 |
|---|---|---|---|
| Dựng lại cả năm (178.513 dòng, 18 MB) | 700.262 đơn completed | 4.252 ms | 5.219 ms |
| Cập nhật một ngày vừa qua (494 dòng) | 1.983 key, 0 document (index phủ) | 23,1 ms | 29,4 ms |
Dashboard một tenant, đọc tháng 9 từ daily_revenue | 30 key, 30 document | 0,8 ms | 0,8 ms |
So sánh: cùng câu hỏi tính thẳng từ orders | 126 key, 126 document | 3,3 ms | 1,7 ms |
Với một tenant, chênh lệch nhỏ vì index trên orders đã tốt. Lợi ích lớn nằm ở chỗ khác:
- Câu hỏi càng rộng thì rollup càng lợi. Một năm của tenant t0042 là 1.416 đơn completed nếu tính thẳng, nhưng chỉ 359 dòng trong rollup. Báo cáo toàn hệ thống cả năm là 700.262 đơn, so với 178.513 dòng.
- Chi phí được dời sang chỗ bạn kiểm soát được. Job cập nhật một ngày tốn 23–29 ms, chạy định kỳ, thay vì mỗi lượt mở dashboard lại tốn vài trăm ms trên
orders.
Rollup bằng $merge
│
├── ✓ đọc dashboard: vài chục key thay vì hàng nghìn / hàng trăm nghìn document
├── ✓ cập nhật tăng dần: chỉ tính lại khoảng thời gian vừa đổi
├── ✗ dữ liệu trễ đến lần chạy job kế tiếp
├── ✗ thêm một collection, một index, một job phải giám sát
└── ⚠ đơn bị sửa muộn (refund sau 10 ngày) → phải tính lại cả ngày cũ, không chỉ "hôm qua"Điểm cuối cùng là chỗ hay sai nhất. Nếu một đơn ngày 14/9 chuyển sang refunded vào ngày 24/9, job "chỉ tính hôm qua" sẽ không bao giờ sửa dòng 14/9. Có hai cách thường gặp: tính lại một cửa sổ trượt (ví dụ 30 ngày gần nhất) mỗi lần chạy, hoặc theo dõi thay đổi bằng change streams (bài 29) để biết ngày nào cần tính lại.
$out khác $merge thế nào? [tài liệu] $out thay toàn bộ collection đích bằng kết quả mới, và không ghi được vào sharded collection. $merge hoà kết quả vào collection có sẵn (insert, replace, merge, giữ nguyên, hoặc báo lỗi) và ghi được vào sharded collection. Dùng $out khi mỗi lần bạn muốn dựng lại toàn bộ. Dùng $merge cho cập nhật tăng dần.
So với PostgreSQL
Hai hệ thống giải cùng một bài toán và gặp cùng những giới hạn vật lý. Khác nhau nằm ở chỗ ai quyết định thứ tự và ai chọn cách join.
| Chủ đề | MongoDB aggregation | PostgreSQL |
|---|---|---|
| Thứ tự xử lý | Pipeline có thứ tự. Optimizer chỉ sắp lại theo một số luật an toàn | SQL khai báo. Planner tự chọn thứ tự dựa trên chi phí ước lượng |
WHERE / HAVING | $match trước / sau $group | hai mệnh đề riêng |
| Join | $lookup: IndexedLoopJoin, HashJoin (collection nhỏ), NestedLoopJoin | nested loop, hash join, merge join, chọn theo thống kê và chi phí |
| Bộ nhớ cho sort/hash | 100 MB mỗi stage, spill khi allowDiskUse | work_mem cho mỗi node sort/hash (mặc định 4 MB), vượt thì ghi file tạm ("external merge" trong EXPLAIN ANALYZE) |
| Materialized view | tự làm bằng $merge + job | CREATE MATERIALIZED VIEW + REFRESH (refresh tính lại toàn bộ) |
Bài học chung: trong PostgreSQL, viết JOIN trước hay sau GROUP BY trong một subquery thường không đổi plan nhiều, vì planner tự sắp xếp lại. Trong MongoDB, như thí nghiệm $lookup ở trên, thứ tự bạn viết là một phần của plan. Bạn chịu trách nhiệm đặt stage đắt ở chỗ dòng dữ liệu hẹp nhất. Đổi lại, bạn biết chính xác pipeline sẽ làm gì, và explain phản ánh đúng thứ bạn viết. Với MongoDB, mô hình document (bài 04) thường giúp tránh join ngay từ đầu. Khi vẫn phải join, $lookup có index chạy tốt cho những tập nhỏ đã được lọc.
Những lỗi thường gặp
$matchkhông đứng đầu, hoặc lọc trên field tính ra. Optimizer kéo được phần lớn điều kiện lên, nhưng điều kiện trên field do$settạo ra phải chờ. Lọc trên field gốc (createdAttheo ranh giới ngày) để điều kiện nằm trong index.$set/$projectđặt trước$sort+$limit.$sortmất khả năng dùng index. Hãy định hình kết quả ở cuối pipeline.- Thêm
$projectở đầu "để bớt dữ liệu". Server đã tự chỉ lấy những field cần dùng. - Cắt ngày theo UTC trong khi người dùng nghĩ theo giờ Việt Nam. Hãy truyền
timezonecho$dateTrunc/$dateToStringvà dùng cùng ranh giới trong$match. $lookupkhông có index trênforeignFieldsang một collection lớn: mỗi document đi vào là một lần quét cả collection (12,7 triệu document cho 126 đơn trong lab).$lookupđặt trước$groupkhi chỉ cần thông tin cho kết quả cuối: 57.464 lần tra thay vì 20.$facetđứng đầu pipeline: COLLSCAN cả collection vào mọi sub-pipeline.- Sort lớn không có
$limit: phải giữ toàn bộ input trong RAM hoặc spill. Có$limitthì chỉ giữ top-n. - Coi
allowDiskUse: truelà cách sửa: pipeline vượt giới hạn 17% đã chậm 4 lần trong lab. Spill là van an toàn, không phải giải pháp. - Tin rằng
$grouptrả kết quả có thứ tự. Nó không đảm bảo thứ tự. Cần thứ tự thì thêm$sort. - Tính lại báo cáo không đổi ở mỗi request. Dữ liệu của ngày đã qua có thể tính trước bằng
$merge, miễn là bạn xử lý được dữ liệu sửa muộn.
Tóm tắt
- Pipeline là một dây chuyền: nửa đầu là một query bình thường (
$cursor, dùng index), nửa sau làm việc trên dòng document không còn index của collection đầu vào. Với SBE,$groupvà$lookupở đầu có thể được gộp vào nửa đầu. - Optimizer tự kéo
$matchlên, tách$matchtheo field, gộp$sort+$limit, gộp$lookup+$unwind. Nó không đưa$setra sau$sort. Trên 8.3.11 còn quan sát được việc đẩy$matchtheo khoá nhóm lên trước$group, nhưng tài liệu không liệt kê luật này. - Ví dụ doanh thu tenant × ngày: COLLSCAN 422–530 ms, IXSCAN + FETCH 267–281 ms, index phủ 126–165 ms, một tenant 1,7–3,3 ms.
$grouptốn RAM theo số nhóm,$sortkhông index tốn RAM theo tổng kích thước input,$sort+$limitchỉ giữ top-n (3.720 byte so với 100+ MB).- Giới hạn 100 MB mỗi stage. Vượt mức thì lỗi hoặc spill (
usedDisk,spills,spilledBytes,spilledDataStorageSize). Trong lab,$groupcần 117 MB đã chậm từ 2,5–2,7 s lên 10,2–10,5 s vì spill. $lookuptốn công theo số document đi qua: không index thì 12,7 triệu document cho 126 đơn, có index thì 127. HashJoin xuất hiện khi collection bị nối nhỏ. Đặt$lookupsau$group/$limitlàm báo cáo nhanh hơn khoảng 20 lần.$facetphải đứng sau một$matchcó index. Đứng đầu thì COLLSCAN, chậm hơn 600–900 lần trong lab.$mergebiến báo cáo đắt thành một collection nhỏ, cập nhật tăng dần: 23–29 ms mỗi ngày, đọc dashboard 0,8 ms.
Tự kiểm tra
- Pipeline
[ { $group: { _id: "$tenantId", n: { $sum: 1 } } }, { $sort: { n: -1 } }, { $limit: 5 } ]có dùng được index{ tenantId: 1 }cho$sortkhông? Stage sort cần bao nhiêu bộ nhớ? (Không.$sortđứng sau$groupnên sắp trên dữ liệu mới. Nhưng$limitđược gộp vào nên chỉ giữ top 5, rất ít bộ nhớ.) - Vì sao
[ $set: {day: ...}, $match: { day: X, tenantId: "t1" } ]vẫn dùng được index{ tenantId: 1 }? (Optimizer tách$match. PhầntenantIdkhông phụ thuộc$setnên được đưa lên đầu. Phầndayvẫn phải chờ.) - Một job nightly chạy 3 giây suốt nửa năm, rồi một tuần nhảy lên 12 giây dù dữ liệu chỉ tăng 10%. Bạn sẽ kiểm tra trường nào trong explain? (
usedDisk,spills,peakTrackedMemBytescủa$group/$sort. Rất có thể nó vừa vượt 100 MB và bắt đầu spill.) - Báo cáo cần tên khách hàng cho top 10 khách theo doanh thu năm. Đặt
$lookupsangcustomersở đâu? (Sau$grouptheocustomerId,$sortvà$limit: 10. Khi đó chỉ có 10 lần tra.)
Nếu phải giải thích bài này mà không dùng thuật ngữ MongoDB nào: dữ liệu chạy qua một dây chuyền trạm. Chỉ trạm ở cửa kho mới có sổ vị trí hàng. Trạm gom thùng phải chờ hết hàng và cần chỗ trên sàn, hết chỗ thì phải chạy ra kho tạm ngoài sân. Mỗi lần chạy sang kho bên cạnh tốn một chuyến cho từng món, nên hãy chạy sau khi đã gom và cắt bớt hàng. Còn báo cáo không đổi thì tính một lần rồi dán lên bảng.
Bài tiếp theo
Bài này và các bài 07–10 cho bạn bộ công cụ cho một query: index, planner, explain, pipeline. Bài 13, Query Performance & Pagination, đưa chúng vào bối cảnh hệ thống. Vì sao skip lớn làm trang 500 của danh sách đơn chậm dần, và phân trang theo (tenantId, createdAt, _id) giải quyết ra sao. Thiết kế index cho cả hệ thống multi-tenant thay vì cho từng câu query. Và khi production chậm thì bắt đầu từ đâu: slow query log, profiler, những trường như usedDisk vừa gặp ở đây. Đó cũng là chỗ để phân biệt query nhanh với hệ thống khoẻ.