Aggregation Pipeline: dây chuyền xử lý dữ liệu và cái giá của từng trạm

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

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 & Pagination

Môi trường lab: MongoDB 8.3.11 chạy trong Docker (image mongo:8, standalone), container mongo-lab-09 giới hạn 2 CPU, 3 GB RAM, WiredTiger cache 1 GB, máy host Apple M4. Database riêng lab09. 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.

StageLoạiDùng index của collection đầu vào?Tốn gì
$matchstreamingCó, nếu nằm ở đầu (sau tối ưu hoá)đọc document / key
$project, $set, $unsetstreamingKhôngCPU tính expression
$sortblocking nếu không có indexCó, nếu không có $project, $unwind, $group đứng trướcRAM theo tổng kích thước input; với $limit thì chỉ n phần tử
$limit, $skipstreaminggộp vào stage trướcgần như không
$groupblockingChỉ trường hợp đặc biệt ($first/$last, DISTINCT_SCAN)RAM theo số nhóm
$unwindstreamingKhôngnhân số document lên theo độ dài mảng
$lookupstreamingDùng index của collection bị nốimột lần tra cứu cho mỗi document đi qua
$facetblockingKhông (bên trong sub-pipeline)RAM, kết quả là một document ≤ 16 MiB
$bucket, $bucketAutoblockingKhôngnhư $group
$merge, $outcuối pipelineKhôngghi 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ài

Dữ 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 cho status/createdAt, nên server đọc cả collection để lấy 57.464 đơn khớp (số này đọc được ở stage filter trong explain đầy đủ).
  • GROUP nằm trong winningPlan. $group được đẩy xuống query layer và chạy bằng SBE. [tài liệu] Từ 5.2, $group chạ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.
  • $sort vẫn là một stage riêng ở nửa sau, sắp xếp 14.650 dòng, dùng khoảng 10 MB.
  • peakTrackedMemBytes của $group khoả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 0

Index 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 30

Kết quả đo

Cách chạyĐọc gìmedian lượt 1median lượt 2
Không index (COLLSCAN)1.000.000 document422 ms530 ms
Index {status, createdAt}57.464 key + 57.464 document281 ms267 ms
Index phủ {status, createdAt, tenantId, total}57.464 key, 0 document126 ms165 ms
Một tenant, index {tenantId, status, createdAt}126 key + 126 document3,3 ms1,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óm

Hai điều đáng chú ý:

  1. 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.
  2. 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ậtTrướcSau
$match sau $project/$set/$addFields$set → $matchphầ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 + $limithai stagemột $sort chỉ giữ top-n
$sort + $skip + $limitba stage$sort giữ top-(skip+limit), rồi $skip
$limit + $limit, $skip + $skiphai stagemột stage (min / tổng)
$match + $matchhai stagemộ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ả $sort lẫ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 10

Index {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.

Pipelinekeys / docsmedian lượt 1median lượt 2
A (viết ngược)203 / 203 + sort trong RAM1,5 ms1,1 ms
A2 (đúng thứ tự)10 / 100,7 ms0,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 $match lê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 revenue của một tenant trước khi cộng xong, nên $match này buộc phải chờ $group duyệ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.tenantId chính là tenantId của input, nên đưa điều kiện lên trước $group và 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ÁC

Ba hệ quả thực tế:

  1. $sort sau $group luôn là sort trong RAM. Kết quả của $group là dữ liệu mới, không có index nào. [tài liệu] $group cũng không đảm bảo thứ tự output. Muốn có thứ tự thì phải có $sort sau nó. Tin tốt: số nhóm thường nhỏ hơn rất nhiều so với số document.
  2. 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 đó.
  3. $lookup là 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 (peakTrackedMemBytes 103,8 MB).
  • spilledBytes là ướ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. spilledDataStorageSize là 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.
Pipelinemedian lượt 1median lượt 2
SORTALL, giới hạn 100 MB (spill)2.881 ms2.747 ms
SORTALL, giới hạn sort nâng lên 1 GB (không spill)1.936 ms2.343 ms
TOPK (sort + limit 10)288 ms465 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 1median lượt 2
giới hạn 100 MB → spill 2 lần10.501 ms10.248 ms
giới hạn 500 MB → không spill2.708 ms2.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ử:

  1. Lọc sớm hơn. $match theo thời gian, theo tenant. Ít document đi vào $group thì ít nhóm hơn.
  2. 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 $addToSet trên mảng lớn trong $group, vì mỗi nhóm sẽ phình theo số document.
  3. Thêm $limit ngay sau $sort khi chỉ cần top-n.
  4. Có index cho $sort ở đầu pipeline để sort không còn blocking.
  5. Tính trước bằng $merge (phần cuối bài).
  6. 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 explainmedian lượt 1median lượt 2
không có indextotalDocsExamined: 12.700.000, collectionScans: 1272.137 ms2.832 ms
có index customerId_1totalDocsExamined: 127, totalKeysExamined: 127, indexesUsed: ['customerId_1']2,2 ms2,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 }
StrategyKhi nào (quan sát)Chi phí
IndexedLoopJoincó index trên foreignFieldmỗi document: một lần seek index
HashJoinkhô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
NestedLoopJoinkhông có index, collection bị nối lớnmỗ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:

Pipelinesố lần $lookup tramedian lượt 1median lượt 2
BEFORE57.464494 ms428 ms
AFTER, $lookup sau $sort+$limit2022,8 ms21,7 ms
AFTER, $lookup trước $sort+$limit50028,5 ms21,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 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 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: $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 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 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

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 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 ở 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 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
Join$lookup: IndexedLoopJoin, HashJoin (collection nhỏ), NestedLoopJoinnested loop, hash join, merge join, chọn theo thống kê và chi phí
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, 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

  • $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.
  • 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 timezone cho $dateTrunc/$dateToString và dùng cùng ranh giới trong $match.
  • $lookup không có index trên foreignField sang 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 $group khi 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ó $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.
  • Tin rằng $group trả 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, $group và $lookup ở đầu có thể được gộp vào nửa đầu.
  • Optimizer tự kéo $match lên, tách $match theo field, gộp $sort + $limit, gộp $lookup + $unwind. Nó không đưa $set ra sau $sort. Trên 8.3.11 còn quan sát được việc đẩy $match theo 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.
  • $group tốn RAM theo số nhóm, $sort không index tốn RAM theo tổng kích thước input, $sort + $limit chỉ 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, $group cần 117 MB đã chậm từ 2,5–2,7 s lên 10,2–10,5 s vì spill.
  • $lookup tố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 $lookup sau $group/$limit làm báo cáo nhanh hơn khoảng 20 lần.
  • $facet phải đứng sau một $match có index. Đứng đầu thì COLLSCAN, chậm hơn 600–900 lần trong lab.
  • $merge biế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

  1. Pipeline [ { $group: { _id: "$tenantId", n: { $sum: 1 } } }, { $sort: { n: -1 } }, { $limit: 5 } ] có dùng được index { tenantId: 1 } cho $sort không? Stage sort cần bao nhiêu bộ nhớ? (Không. $sort đứng sau $group nê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ớ.)
  2. Vì sao [ $set: {day: ...}, $match: { day: X, tenantId: "t1" } ] vẫn dùng được index { tenantId: 1 }? (Optimizer tách $match. Phần tenantId không phụ thuộc $set nên được đưa lên đầu. Phần day vẫn phải chờ.)
  3. 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, peakTrackedMemBytes của $group/$sort. Rất có thể nó vừa vượt 100 MB và bắt đầu spill.)
  4. Báo cáo cần tên khách hàng cho top 10 khách theo doanh thu năm. Đặt $lookup sang customers ở đâu? (Sau $group theo customerId, $sort và $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ẻ.

Tài liệu tham khảo