Aggregation (P3/3): `$unwind`, `$facet`, `$merge`
Ở phần trước: server tự sắp xếp lại pipeline theo các luật cố định, và index chỉ giúp đoạn đầu. Mỗi stage bị giới hạn 100 MB trong RAM, và việc vượt mức phụ thuộc cấu hình allowDiskUse.
- Cần đọc trước: Server sắp xếp lại, index ở đầu, giới hạn 100 MB
- Dẫn tới: $lookup & Joins, bài này khép lại ở đây.
$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 lẫn trong dao động đo, trái với lời khuyên "đừng bao giờ $unwind" hay gặp. $unwind đắt khi mảng dài (một document 1.000 phần tử thành 1.000 document), khi theo sau là stage đắt 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; trường hợp cuối, 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
Ý chính
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, trong 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
Ý chính
Doanh thu của ngày 14/9 không thay đổi sau khi ngày 14/9 kết thúc (tạm bỏ qua hoàn tiền về sau). Vậy tại sao mỗi lần 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 ở câu hỏi rộng: một năm của tenant t0042 là 1.416 đơn completed nhưng chỉ 359 dòng rollup; cả hệ thống là 700.262 đơn so với 178.513 dòng. Chi phí dời sang một job 23–29 ms chạy định kỳ.
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 để 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
| 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 |
| 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, planner tự sắp lại thứ tự xử lý. Trong MongoDB, như Thí nghiệm A ở phần trước (Server sắp xếp lại, index ở đầu, giới hạn 100 MB) và $facet ở 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, thứ tự xử lý dễ đoán, và explain cho thấy chính xác pipeline sau khi sắp lại. Join và các chiến lược join của hai bên được so sánh ở bài $lookup & Joins.
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.- 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. $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. - 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.
Cột mốc: Bạn đã biết $unwind nhân số document cho các stage sau, $facet chạy nhiều báo cáo từ một lần đọc, và $merge giúp tính một lần rồi đọc nhiều lần. Bài Aggregation Pipeline khép lại ở đây.
Hỏi & đáp
Dashboard cần bốn con số cho tháng 9 của tenant t0042. Cách viết nào nhanh nhất trong lab?
Bảng daily_revenue được job $merge cập nhật mỗi đêm, chỉ tính lại "hôm qua". Rủi ro chính của thiết kế này là gì?
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. 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 bỏ qua một trạm: trạm chạy sang kho bên cạnh. Bài $lookup & Joins dành trọn cho nó. Vì sao thiếu index ở collection bị nối biến 126 đơn thành 12,7 triệu document phải đọc, hai cú pháp $lookup chạy trên hai engine khác nhau ra sao, server gộp $lookup + $unwind + $match thế nào và khi nào nó không gộp được, vì sao đặt $lookup sau $group và $limit nhanh hơn khoảng 20 lần. Và câu hỏi lớn hơn mọi mẹo tối ưu: khi nào nên embed thay vì join.