Query Execution Engine (P1/3): Cây stage, `work()` và SORT

13 phút đọcSeries: MongoDB: từ gốc đến internals1 lượt xem

Bài này trả lời ba câu hỏi. Một: một query dừng sớm được khi nào? Hai: khi nào thì phải đọc hết? Ba: vì sao kết quả có thể lệch khi dữ liệu đổi giữa chừng?

Bài có 3 Part:

  • P1 (Part này): Cây stage, work() và SORT.
  • P2: PROJECTION, yielding và SBE.
  • P3: EXPRESS và bảng tra explain.

Suốt Phần Index & Query, ta đã hỏi chọn cách chạy nào: index nào, compound ra sao, planner chạy thử hay CBR ước lượng. Bài này hỏi câu còn lại: chọn xong rồi thì plan chạy như thế nào?

Nghe như chuyện nội bộ, nhưng nó giải thích nhiều hiện tượng thực tế. Vì sao limit(5) có lúc chỉ đọc 6 document, có lúc phải đọc 1 triệu? Vì sao một câu sort() không index chạy được với 10 kết quả nhưng lỗi với 700.000 kết quả? Vì sao một cursor đang đọc lại trả về cùng một document hai lần? Và vì sao cùng một pipeline $group, explain lúc thì có works, lúc thì không?

Bài này nằm ở đâu

  • Cần biết trước: CRUD & Query Model (cursor, batch, COLLSCAN 1 triệu document), Index Fundamentals (IXSCAN + FETCH), Compound Indexes & ESR (sort bằng index, covered query), Query Planner & Plan Cache (cây plan, works trong trial), Aggregation Pipeline (stage streaming và blocking, spill)
  • Giới thiệu: cây stage kiểu pull, work() và bốn trạng thái ADVANCED / NEED_TIME / NEED_YIELD / IS_EOF, đọc executionStats theo từng stage, SORT (top-k, giới hạn 100 MB, spill, SORT_MERGE, SORT_KEY_GENERATOR), các biến thể PROJECTION, yielding và chuyện đọc không phải snapshot, classic engine vs slot-based engine (SBE), EXPRESS_IXSCAN, bảng tra explain
  • Dẫn tới: Transactions & Atomicity (Phần Transactions)

Môi trường lab: MongoDB 8.3.11 trong Docker (image mongo:8, standalone), container riêng mongo-lab-15 giới hạn 2 CPU, 3 GB RAM, WiredTiger cache 1 GB, máy host Apple M4, database lab15. Lab của các bài khác chạy song song trên cùng máy host nên thời gian dao động; bài ghi median của nhiều lần chạy. Hãy so các con số với nhau, đừng coi chúng là latency production. Số nào không đo thì ghi là minh hoạ.

Nhãn: [tài liệu] theo tài liệu chính thức; [quan sát] đo trong lab này (8.3.11); [chi tiết cài đặt] cách server đang làm, không phải cam kết API; [hình dung] mô hình để dễ nhớ.

Ý chính

Plan là một cây stage. Thực thi plan nghĩa là stage gốc liên tục hỏi stage con: "cho tôi kết quả tiếp theo". Stage con hỏi tiếp xuống dưới, tới tận stage lá đọc index hay collection. Mỗi lần hỏi là một work, và (theo mô hình của classic engine) câu trả lời là một trong bốn: "đây, một kết quả", "chưa có, hỏi lại đi", "đợi, tôi cần nhường", "hết rồi".

Vì dữ liệu được kéo từng cái, một query dừng được ngay khi đủ kết quả. Trừ khi giữa đường có một stage blocking như SORT, vốn phải nuốt hết input rồi mới nhả ra cái đầu tiên. Trong lúc chạy, query còn nhường (yield) định kỳ. Mỗi lần nhường, nó buông snapshot dữ liệu, nên một query dài không nhìn dữ liệu ở một thời điểm duy nhất. Cuối cùng, cùng một cây plan có thể được chạy bởi hai engine: classic (cây stage gọi nhau như trên) và SBE (dịch plan thành chương trình làm việc trên các "slot").

Hình dung trước: quầy bánh mì làm theo đơn

Một quầy bánh mì có bốn người đứng thành hàng. Người cuối là người đóng hộp. Khách cần 5 ổ.

Người đóng hộp quay sang người kẹp nhân: "cho anh một ổ". Người kẹp nhân quay sang người chọn bánh: "cho một cái bánh". Người chọn bánh lấy phiếu kế tiếp, ra kệ lấy bánh, thấy cháy, vứt đi, nói "chưa có, hỏi lại". Lần sau thì có bánh đẹp, chuyền lên. Đủ 5 ổ, người đóng hộp nói "xong", và không ai làm thêm ổ thứ sáu.

Giờ đổi đơn: "5 ổ to nhất trong cả mẻ". Người đóng hộp không thể lấy 5 ổ đầu tiên được nữa. Anh phải chờ cả mẻ đi qua, so kích thước, rồi mới có ổ đầu tiên để giao. Mẻ càng lớn thì bàn để bánh chờ so càng phải rộng.

Thêm một chi tiết. Thỉnh thoảng người chọn bánh rời chỗ để nhường kệ cho người khác xếp lại, rồi quay về đi tiếp từ "phiếu kế tiếp". Nếu lúc anh vắng, ai đó dời một ổ đã giao về cuối kệ, anh sẽ giao nó lần nữa.

Quầy bánh mì                                MongoDB
──────────────────────────────────────      ──────────────────────────────────────
mỗi người một việc, đứng thành hàng         cây stage: LIMIT ← PROJECTION ← FETCH ← IXSCAN
"cho anh một ổ"                             work(): stage cha gọi stage con
"đây" / "chưa có, hỏi lại" / "hết bánh"     ADVANCED / NEED_TIME / IS_EOF
đủ 5 ổ thì cả quầy dừng                     pull-based: LIMIT dừng được sớm
"5 ổ to nhất": chờ cả mẻ                    SORT là stage blocking
bàn chờ so kích thước                       bộ nhớ của SORT, giới hạn 100 MB
rời chỗ một chút rồi quay lại               yield: saveState / restoreState
ổ bị dời về cuối kệ, giao hai lần           cursor không phải snapshot

[hình dung] Không có bốn người chạy song song: mọi stage được gọi nối nhau trên cùng một luồng, và "nhường" chỉ kéo dài trong chốc lát.

Mental model: cây stage và nhịp kéo

                    client: find / getMore
                             │  "cho tôi một batch"
                             ▼
                      ┌─────────────┐
                      │   LIMIT     │  work() ──► ADVANCED  (1 kết quả đi lên)
                      └──────┬──────┘         ├─► NEED_TIME (đã làm, chưa có gì)
                             │ work()         ├─► NEED_YIELD(phải nhường storage)
                      ┌──────▼──────┐         └─► IS_EOF    (hết)
                      │ PROJECTION  │
                      └──────┬──────┘
                             │ work()
                      ┌──────▼──────┐
                      │   FETCH     │  đọc document theo RecordId, áp filter
                      └──────┬──────┘
                             │ work()
                      ┌──────▼──────┐
                      │   IXSCAN    │  đọc một index key
                      └─────────────┘

    giữa các work, định kỳ:  saveState ──► (yield: buông snapshot, lock) ──► restoreState

Ba ý giữ suốt bài:

  1. Dữ liệu được kéo từ trên xuống, kết quả đi từ dưới lên, mỗi lần một cái. Với classic engine, số works ở mỗi stage xấp xỉ số lần nó bị gọi.
  2. Stage streaming nhả kết quả ngay khi có. Stage blocking (SORT không có index, GROUP) phải nuốt hết input trước. Một stage blocking ở giữa cây làm mất lợi thế "dừng sớm" của mọi thứ bên dưới nó.
  3. Một query dài không chạy trên một snapshot duy nhất. Nó nhường định kỳ, và mỗi lần nhường, thế giới bên dưới có thể đã đổi.

Setup và dataset

Dataset giống bài CRUD & Query Model: đơn hàng của một SaaS nhiều tenant.

MongoDB version : 8.3.11 (mongo:8, standalone)
Hardware        : Apple M4 host; container 2 CPU, 3 GB RAM
Configuration   : --wiredTigerCacheSizeGB 1
Dataset         : lab15.orders, 1.000.000 document, avgObjSize 215 byte
                  215,7 MB chưa nén, 71,6 MB trên đĩa
Indexes         : _id_ lúc đầu; sau đó thêm { tenantId: 1, createdAt: -1 }
                  và { tenantId: 1, status: 1, total: 1 }
// trích gen.js: PRNG mulberry32 seed 15, 100 lần insertMany x 10.000 document
docs.push({
  tenantId: "t" + String(Math.floor(rnd() * 500)).padStart(4, "0"),   // 500 tenant
  userId:   "u" + String(Math.floor(rnd() * 200)).padStart(4, "0"),
  status:   pickStatus(rnd()),       // 70% completed, 10% pending, 10% cancelled, 10% refunded
  createdAt: new Date(END - Math.floor(rnd() * 365 * DAY)),
  total:    items.reduce((s, it) => s + it.qty * it.price, 0),          // VND
  items                                                                 // 1–3 món { sku, qty, price }
});
inserted 1000000 in 28721 ms

Seed khác bài CRUD nên số lệch một chút: query "t0042, pending, từ 01/09" ở đây trả 17 đơn (bài CRUD: 21).

Thêm một collection nhỏ cho phần EXPRESS: lab15.customers, 100.000 khách, email có unique index, phone có index thường. Cả hai field đều không trùng giá trị.

work(): bốn câu trả lời của một stage

Classic engine chạy plan bằng cách gọi work() lặp lại

[tài liệu] Tài liệu explain mô tả works là số "đơn vị công việc" mà stage đã làm: "query execution chia việc thành những đơn vị nhỏ". Một đơn vị có thể là đọc một index key, fetch một document, áp projection lên một document, hoặc làm một chút việc nội bộ. Trường này chỉ có khi query chạy bằng classic engine.

[chi tiết cài đặt, mô hình đơn giản hoá] Mỗi lần được gọi, stage trả về một trong bốn trạng thái:

Trạng tháiNghĩaĐếm vào trường
ADVANCEDCó một kết quả đưa lên stage chaadvanced
NEED_TIMEĐã làm việc nhưng chưa có kết quả (document không khớp filter, hoặc stage đang gom input)needTime
NEED_YIELDStorage engine cần query nhường (ví dụ xung đột ghi, dữ liệu chưa sẵn)needYield
IS_EOFHết dữ liệuisEOF: 1

Tên bốn trạng thái lấy từ mã nguồn classic engine (PlanStage::StageState), không phải cam kết API. Explain chỉ cho thấy các bộ đếm tương ứng; phép cộng ra works ở dưới là suy từ số đo, không phải công thức được tài liệu hoá.

Đọc lại COLLSCAN 1 triệu document theo nhịp works

Query quen thuộc của bài CRUD, chưa có index nào ngoài _id:

db.orders.find({
  tenantId: "t0042",
  status: "pending",
  createdAt: { $gte: ISODate("2026-09-01T00:00:00Z") }
}).sort({ createdAt: -1 }).explain("executionStats")
SORT       works 1000019  advanced 17  needTime 1000001  needYield 0  saveState 9  isEOF 1
           memLimit 104857600  type 'simple'  totalDataSizeSorted 4136  usedDisk false
  COLLSCAN works 1000001  advanced 17  needTime 999983   needYield 0  saveState 9  isEOF 1
           docsExamined 1000000

Median 9 lần explain: 180 ms ở lượt 1, 144 ms ở lượt 2 (lượt 1 nhiễu hơn vì lab khác đang nạp dữ liệu).

Đọc COLLSCAN trước, vì nó là lá:

  • works 1.000.001 = 1.000.000 lần đọc document + 1 lần cuối trả IS_EOF.
  • advanced 17: 17 document khớp filter, được đưa lên SORT.
  • needTime 999.983: những document đọc lên rồi vứt. Kiểm tra: 17 + 999.983 + 1 = 1.000.001. Với một stage lá, works ≈ advanced + needTime + 1 (cộng thêm needYield nếu có).

Rồi tới SORT:

  • needTime 1.000.001: mỗi lượt gọi con, SORT chưa có gì để trả lên vì phải gom hết.
  • advanced 17, rồi 1 lần IS_EOF: 1.000.001 + 17 + 1 = 1.000.019 works.

Công nằm ở lá: một triệu lượt để 17 kết quả đi lên, tỉ lệ docsExamined / nReturned ≈ 58.800 : 1, giờ thấy được tới từng stage.

Pull-based nghĩa là dừng được sớm

Cùng COLLSCAN (ép bằng hint({ $natural: 1 })), ba câu limit(5):

find({ status: "completed" }).limit(5)
  LIMIT    works 7        advanced 5   isEOF 1
  COLLSCAN works 6        advanced 5   needTime 1   isEOF 0   docsExamined 6

find({ tenantId: "t0042" }).limit(5)
  LIMIT    works 2256     advanced 5   isEOF 1
  COLLSCAN works 2255     advanced 5   needTime 2250  isEOF 0  docsExamined 2255

find({ tenantId: "t0042" }).sort({ total: -1 }).limit(5)
  SORT     works 1000007  advanced 5   needTime 1000001  isEOF 1
  COLLSCAN works 1000001  advanced 1986                  isEOF 1  docsExamined 1000000

[quan sát] Câu đầu: 70% đơn là completed, 6 document là đủ 5 kết quả, COLLSCAN dừng ở isEOF: 0. Câu thứ hai: t0042 chỉ chiếm 0,2% dữ liệu, nên phải đọc 2.255 document mới gặp đủ 5 đơn. Câu thứ ba thêm sort() không có index: COLLSCAN phải đi tới cuối (1.000.000 document, 139 ms), vì đơn có total lớn nhất có thể nằm ở bất cứ đâu.

Câu thứ ba không có stage LIMIT: limit đã được gộp vào SORT (limitAmount: 5).

Bài học: limit() chỉ rẻ khi không có stage blocking bên dưới và kết quả khớp đủ dày. Vì vậy bài Pagination đòi index cung cấp luôn thứ tự sort.

Streaming và blocking: SORT

Top-k: limit gộp vào sort

[tài liệu] Khi sort() đi với limit() mà index không cho được thứ tự, MongoDB dùng thuật toán top-k: chỉ giữ k kết quả tốt nhất đã thấy tới lúc đó. Nếu k kết quả này vượt 100 MB thì mới phải ghi file tạm.

Thử trên 700.217 đơn completed, sort theo total giảm dần, không có index hỗ trợ:

QueryPlanDữ liệu SORT giữSpillMedian (ms)
.limit(10)SORT(limit 10) ← COLLSCAN2.940 byteKhông194
Không limitSORT ← COLLSCAN173,7 MB (peak 103,8 MB)Có, 2 lần, 154,8 MB ra đĩa541
Không limit, nâng giới hạn sort lên 300 MB (chỉ trong lab)SORT ← COLLSCAN173,7 MBKhông410
Không limit, .allowDiskUse(false)Lỗi sau khi đọc 603.776 document

[quan sát] Top-k đọc đúng 1 triệu document như bản không limit, nhưng chỉ giữ 2,9 KB. Gần hết 194 ms của nó là thời gian COLLSCAN. Bản không limit phải giữ 173,7 MB, vượt 100 MB, nên ghi ra đĩa hai lần và chậm hơn bản không spill khoảng 30% (541 so với 410 ms).

Giới hạn 100 MB và allowDiskUse với find

[tài liệu] Nếu plan có stage SORT, MongoDB sort trong bộ nhớ và bị giới hạn 100 MB. Từ 6.0, tham số server allowDiskUseByDefault mặc định là true: vượt giới hạn thì tự ghi file tạm. cursor.allowDiskUse(false) cấm việc đó cho một query. allowDiskUse() không có tác dụng với sort được index trả lời hoặc sort dưới 100 MB.

[quan sát] Lab đọc được allowDiskUseByDefault: true và tham số nội bộ internalQueryMaxBlockingSortMemoryUsageBytes: 104857600. Khi cấm spill:

MongoServerError[QueryExceededMemoryLimitNoDiskUseAllowed]: Executor error during find command
  :: caused by :: Sort exceeded memory limit of 104857600 bytes, but did not opt in to external sorting.

Ở stage SORT, đọc memLimit, limitAmount (top-k), totalDataSizeSorted, usedDisk; [tài liệu] từ 8.2 có thêm spills, spilledBytes, spilledRecords, spilledDataStorageSize, từ 8.3 có peakTrackedMemBytes. Chi phí spill của $sort và $group đã được đo ở bài Aggregation Pipeline.

Projection đặt dưới SORT làm sort nhẹ đi

Thêm projection chỉ lấy total vào câu sort không limit:

find({ status: "completed" }, { total: 1, _id: 0 }).sort({ total: -1 })

SORT ← PROJECTION_SIMPLE ← COLLSCAN
totalDataSizeSorted 33.610.416   usedDisk false   spills 0     median 521 ms

[quan sát] Projection được đặt bên dưới SORT. SORT chỉ còn phải giữ các document đã cắt gọn: 33,6 MB thay vì 173,7 MB, nằm dưới giới hạn nên không spill. Thời gian thì chỉ giảm nhẹ trong lab này (521 so với 541 ms của bản spill), vì PROJECTION vẫn phải chạy trên 700.217 document: lợi ích chính là bộ nhớ và không ghi đĩa. Sort lớn đang spill thì kiểm tra xem có đang kéo cả document qua SORT chỉ để dùng vài field không.

SORT_MERGE: trộn nhiều luồng đã có thứ tự

Với index { tenantId: 1, createdAt: -1 } và query nhiều tenant:

find({ tenantId: { $in: ["t0042", "t0043", "t0044"] } }).sort({ createdAt: -1 }).limit(20)
LIMIT ← FETCH ← SORT_MERGE ← [ IXSCAN t0042 | IXSCAN t0043 | IXSCAN t0044 ]
keys 22   docs 20   nReturned 20

find({ tenantId: { $in: ["t0042", "t0043", "t0044"] } }).sort({ total: -1 }).limit(20)
FETCH ← SORT(limit 20) ← IXSCAN tenantId_1_status_1_total_1
keys 5930   docs 20   nReturned 20

[quan sát] Câu đầu: mỗi tenant là một đoạn index đã xếp sẵn theo createdAt. Planner tách $in thành ba IXSCAN, và SORT_MERGE trộn ba luồng như trộn ba chồng phiếu đã xếp: mỗi lần nhìn đầu ba chồng, lấy cái mới nhất. Đây là stage streaming: đủ 20 thì dừng, chỉ đọc 22 key. Câu thứ hai sort theo total, mà không index nào cho sẵn thứ tự total xuyên qua ba tenant (trong { tenantId, status, total }, status đứng giữa), nên quay về SORT blocking (top-k 20), đọc 5.930 key.

Bài Compound Indexes & ESR có ngưỡng của chiêu này: $in quá dài thì planner không tách nữa. Bài Pagination gặp SORT_MERGE ở một $or keyset.

SORT_KEY_GENERATOR: khi phải mang sort key lên trên

Trên standalone, lab chỉ thấy stage này khi query yêu cầu chính sort key làm metadata:

find({ tenantId: "t0042" }, { createdAt: 1, k: { $meta: "sortKey" } }).sort({ createdAt: -1 }).limit(3)
LIMIT ← PROJECTION_DEFAULT ← SORT_KEY_GENERATOR ← FETCH ← IXSCAN

[quan sát + chi tiết cài đặt] Thứ tự đã có từ index, nên không cần SORT. Nhưng index chỉ cho thứ tự, không cho giá trị sort key gắn kèm từng document, nên SORT_KEY_GENERATOR tính lại key đó. Stage streaming, rẻ. Nó cũng có thể xuất hiện trên shard của cluster, nơi mongos cần sort key để trộn kết quả (lab standalone không kiểm tra được).

Cột mốc: Bạn đã đọc được works, advanced và needTime trong explain, biết vì sao SORT chặn stage phía trên, và biết giới hạn 100 MB của sort. Tiếp theo: PROJECTION, yielding và SBE.

Hỏi & đáp

Explain một COLLSCAN ghi works: 50001, advanced: 0, needTime: 50000, isEOF: 1. Collection có bao nhiêu document và query trả bao nhiêu kết quả?

  1. 50.001 document, 1 kết quả

    Work cuối cùng không đọc document nào: nó trả IS_EOF. Và advanced: 0 nghĩa là không document nào được đưa lên stage cha. Xem mục "Đọc lại COLLSCAN 1 triệu document theo nhịp works".

  2. 50.000 document, 0 kết quả

    Đúng. Với stage lá, works ≈ advanced + needTime + 1: 0 + 50.000 + 1 = 50.001. Mỗi needTime là một document đọc lên rồi vứt, work cuối là IS_EOF. Giống COLLSCAN của lab: 17 + 999.983 + 1 = 1.000.001. Xem mục "Đọc lại COLLSCAN 1 triệu document theo nhịp works".

  3. 50.000 document, 50.000 kết quả

    needTime đếm những lần stage làm việc mà chưa có kết quả, như document không khớp filter. Kết quả đi lên được đếm ở advanced, ở đây là 0. Xem mục "work(): bốn câu trả lời của một stage".

  4. Không suy ra được: works không liên quan tới số document

    Với COLLSCAN, mỗi work là một lần đọc document, cộng một lần cuối trả IS_EOF; lab đối chiếu được works 1.000.001 với docsExamined 1.000.000. Xem mục "Đọc lại COLLSCAN 1 triệu document theo nhịp works".

Không có index nào ngoài _id. Query find({ tenantId: "t0042" }).sort({ total: -1 }).limit(5) đọc bao nhiêu document?

  1. Khoảng 6, vì limit(5) cho phép dừng sớm như find({ status: "completed" }).limit(5)

    Câu status: "completed" dừng ở 6 document vì không có stage blocking và 70% đơn khớp. Thêm sort() không index thì không dừng sớm được nữa. Xem mục "Pull-based nghĩa là dừng được sớm".

  2. Khoảng 2.255, như find({ tenantId: "t0042" }).limit(5) không có sort

    2.255 là khi không có sort: đọc tới khi gặp đủ 5 đơn của t0042. Với sort theo total, đơn lớn nhất có thể nằm ở bất cứ đâu, nên phải đọc hết. Xem mục "Pull-based nghĩa là dừng được sớm".

  3. 1.000.000: SORT là stage blocking, limit chỉ được gộp vào SORT (limitAmount: 5)

    Đúng. COLLSCAN phải đi tới cuối (1.000.000 document, 139 ms) vì SORT phải thấy hết input trước khi nhả kết quả đầu tiên. limit() chỉ rẻ khi không có stage blocking bên dưới. Xem mục "Pull-based nghĩa là dừng được sớm".

  4. 5, vì limit được đẩy xuống COLLSCAN

    Explain cho thấy ngược lại: không có stage LIMIT, limit nằm trong SORT, còn COLLSCAN ghi docsExamined 1000000 và isEOF 1. Xem mục "Pull-based nghĩa là dừng được sớm".

Một find sort 700.217 đơn theo total, không limit, không index hỗ trợ: SORT giữ 173,7 MB, spill 2 lần, 541 ms. Cách sửa nào đúng hướng nhất?

  1. Gọi .allowDiskUse(false) để buộc query chạy trong bộ nhớ

    Cấm spill không làm SORT nhỏ đi: lab báo lỗi QueryExceededMemoryLimitNoDiskUseAllowed sau khi đọc 603.776 document. Xem mục "Giới hạn 100 MB và allowDiskUse với find".

  2. Nâng internalQueryMaxBlockingSortMemoryUsageBytes lên 300 MB để hết spill

    Lab thấy nhanh hơn (410 ms) nhưng đó là tham số nội bộ, và nâng giới hạn sort là đổi lỗi lấy nguy cơ OOM. Xem mục "Top-k: limit gộp vào sort".

  3. Không cần sửa gì: allowDiskUse mặc định đã bật nên query vẫn chạy

    allowDiskUse chỉ đổi lỗi thành chậm: 541 ms so với 410 ms khi không spill. Query chạy được không có nghĩa là đã tốt. Xem mục "Top-k: limit gộp vào sort".

  4. Giảm dữ liệu đi qua SORT: projection chỉ lấy field cần, thêm limit, hoặc index cho thứ tự

    Đúng. Projection đặt dưới SORT đưa dữ liệu sort xuống 33,6 MB và hết spill; thêm limit(10) biến nó thành top-k chỉ giữ 2.940 byte; index cho thứ tự thì không còn SORT. Xem mục "Projection đặt dưới SORT làm sort nhẹ đi" và "Top-k: limit gộp vào sort".