Aggregation (P3/3): `$unwind`, `$facet`, `$merge`

9 phút đọcSeries: MongoDB: từ gốc đến internals

Ở 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.

$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 1median lượt 2
U1 $unwind + $group135 ms119 ms
U2 expression trên mảng122 ms122 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: $match có index, rồi $facet.
  • F2: bốn aggregation riêng, mỗi cái có $match riêng.
  • F3: $facet đứng đầu, $match nằ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 1median lượt 2
F1 $match → $facet1,3 ms1,6 ms
F2 bốn aggregation riêng3,2 ms3,8 ms
F3 $facet đứng đầu1.225 ms916 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 trong

Khi 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 1median lượt 2
Dựng lại cả năm (178.513 dòng, 18 MB)700.262 đơn completed4.252 ms5.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 ms29,4 ms
Dashboard một tenant, đọc tháng 9 từ daily_revenue30 key, 30 document0,8 ms0,8 ms
So sánh: cùng câu hỏi tính thẳng từ orders126 key, 126 document3,3 ms1,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 aggregationPostgreSQL
Thứ tự xử lýPipeline có thứ tự. Optimizer chỉ sắp lại theo một số luật an toànSQL khai báo. Planner tự chọn thứ tự dựa trên chi phí ước lượng
WHERE / HAVING$match trước / sau $grouphai mệnh đề riêng
Bộ nhớ cho sort/hash100 MB mỗi stage, spill khi allowDiskUsework_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 viewtự làm bằng $merge + jobCREATE 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

  • $match khô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 $set tạo ra phải chờ. Lọc trên field gốc (createdAt theo ranh giới ngày) để điều kiện nằm trong index.
  • $set/$project đặt trước $sort + $limit. $sort mấ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 timezone cho $dateTrunc/$dateToString và 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ó $limit thì chỉ giữ top-n.
  • Coi allowDiskUse: true là 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?

  1. $facet đứng đầu, mỗi sub-pipeline có $match tenant + tháng 9 riêng

    $facet đứng đầu luôn COLLSCAN: sub-pipeline không dùng được index. F3 đưa cả 1.000.000 đơn vào bốn nhánh, mất 916–1.225 ms. Xem mục "$facet đứng ở đâu quyết định tất cả".

  2. Một $match có index trên tenant + tháng, rồi $facet với bốn nhánh

    Chỉ stage đứng trước $facet dùng được index. F1 đọc 174 đơn một lần, trong một round-trip: 1,3–1,6 ms, so với 3,2–3,8 ms của bốn aggregation riêng và 916–1.225 ms của $facet đứng đầu. Xem mục "$facet đứng ở đâu quyết định tất cả".

  3. Bốn aggregation riêng, vì $facet luôn chậm hơn do không spill được

    Bốn aggregation riêng (F2) mất 3,2–3,8 ms, chậm hơn F1 vì đọc 174 đơn bốn lần. Việc $facet không spill chỉ là giới hạn khi kết quả trung gian vượt 100 MB. Xem mục "$facet đứng ở đâu quyết định tất cả".

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ì?

  1. $merge ghi đè toàn bộ collection mỗi lần, mất dữ liệu các ngày cũ

    Đó là $out. $merge hoà kết quả vào collection có sẵn: whenMatched: "replace" chỉ thay dòng trùng tenantId + day. Xem mục "$merge và $out: tính một lần, đọc nhiều lần".

  2. Đơn refund muộn sau 10 ngày không bao giờ được tính lại cho ngày cũ

    Job "chỉ tính hôm qua" không sửa dòng 14/9 khi một đơn ngày 14/9 thành refunded vào 24/9. Cách thường gặp: tính lại một cửa sổ trượt (vd. 30 ngày) hoặc theo dõi thay đổi bằng change streams. Xem mục "Dựng bảng daily_revenue".

  3. Đọc dashboard từ daily_revenue chậm hơn tính thẳng từ orders

    Ngược lại: đọc tháng 9 của một tenant từ rollup mất 0,8 ms (30 key) so với 1,7–3,3 ms tính thẳng; với câu hỏi rộng, chênh lệch còn lớn hơn nhiều. Xem mục "Dựng bảng daily_revenue".

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.

Tài liệu tham khảo