BigQuery 14 — Các pha & thông số của stage
Vì sao bài này là lõi của cả series
Ở BigQuery — Query plan & Execution Details bạn đã biết một query được cắt thành stage, và ở BigQuery — Đọc execution graph bạn đã biết các stage nối nhau thành một cây phụ thuộc (DAG) qua shuffle. Nhưng khi mở một stage ra, bạn gặp một rừng con số: waitMsAvg, computeMsMax, slotMs, shuffleOutputBytesSpilled, parallelInputs… Nếu không đọc được rừng con số này, bạn chỉ biết query chậm mà không biết chậm ở đâu và vì sao.
Bài này là bài lõi quan trọng nhất của series: nó giải nghĩa từng thông số của một stage, và quan trọng hơn — dạy bạn cách đọc chúng thành chẩn đoán. Đọc xong bài này, mở một stage bất kỳ bạn sẽ trả lời được ba câu: (1) stage này tiêu thời gian vào pha nào? (2) tải có bị lệch giữa các worker không? (3) nó có đang thiếu bộ nhớ phải tràn ra đĩa không? Ba câu đó là toàn bộ nền tảng để chẩn đoán và tinh chỉnh ở BigQuery — Chẩn đoán & tinh chỉnh.
Mô hình tinh thần: một stage giống một ca làm việc của một đội công nhân song song (worker). Cả đội cùng làm một loại việc, mỗi người xử một phần dữ liệu. Bốn pha là bốn động tác trong ca; avg là "người trung bình" của đội, max là "người làm lâu nhất". Đọc stage tức là đọc ca làm việc đó tắc ở động tác nào và có ai bị dồn việc quá tải không.
Bốn pha thời gian của một stage
Mỗi stage, với mỗi worker, BigQuery đo thời gian chia làm bốn pha tuần tự. Đây là chìa khóa đầu tiên để biết nút thắt:
- WAIT — worker chờ slot rảnh để bắt đầu làm việc. Đây không phải thời gian tính toán; nó là thời gian xếp hàng. WAIT cao nghĩa là stage bị nghẽn tài nguyên, chưa kịp vào việc.
- READ — worker đọc dữ liệu đầu vào: hoặc đọc từ storage (stage lá, đọc bảng trên Colossus) hoặc đọc shuffle từ stage trước (stage ở giữa/trên cây). READ cao nghĩa là đang phải kéo nhiều dữ liệu vào.
- COMPUTE — worker xử lý dữ liệu: filter, join, aggregate, sort, hàm phân tích… Đây là pha "làm việc thật". COMPUTE cao là stage nặng về tính toán.
- WRITE — worker ghi kết quả ra: ghi shuffle cho stage sau, hoặc ghi kết quả cuối. WRITE cao thường đi kèm lượng shuffle output lớn.
Cách đọc tỷ lệ bốn pha để tìm nút thắt — nhìn pha nào chiếm phần lớn thời gian của stage:
| Pha chiếm ưu thế | Diễn giải | Hướng xử lý (chi tiết ở bq-15) |
|---|---|---|
| WAIT cao | Thiếu slot / bị query khác tranh tài nguyên / stage phụ thuộc chưa xong | Xem reservation/slot, lịch chạy, giảm concurrency |
| READ cao | Đọc quá nhiều dữ liệu vào (quét rộng, shuffle lớn) | Cắt cột, partition pruning, giảm dữ liệu vào stage |
| COMPUTE cao | Phép tính nặng (join lớn, sort, regex, UDF) | Đổi kiểu join, thêm cluster, đơn giản hóa biểu thức |
| WRITE cao | Sinh ra shuffle output khổng lồ đẩy sang stage sau | Giảm cardinality trước khi shuffle, lọc sớm |
Bảng lớn — giải nghĩa TỪNG thông số của một stage
Đây là phần tra cứu cốt lõi. Tên field lấy đúng theo cấu trúc query plan của BigQuery (trong job_stages của INFORMATION_SCHEMA.JOBS hoặc REST API Job.statistics.query.queryPlan):
| Field thật | Ý nghĩa | Đọc thế nào | Dấu hiệu bất thường / ngưỡng |
|---|---|---|---|
id | Số thứ tự stage | Định danh để nối phụ thuộc | — |
name | Tên stage (vd S02: Join+) | Gợi ý loại việc chính của stage | — |
status | Trạng thái (COMPLETE, RUNNING…) | Stage đã xong chưa | Kẹt RUNNING lâu = nghi treo |
inputStages | Danh sách id stage nạp shuffle vào stage này | Dựng cây phụ thuộc (xem bq-13) | — |
waitMsAvg / waitMsMax | Thời gian chờ slot: trung bình / lớn nhất theo worker | Pha xếp hàng, không phải tính toán | waitMsAvg chiếm tỷ lệ lớn = nghẽn slot |
readMsAvg / readMsMax | Thời gian đọc input/shuffle: avg / max | Pha kéo dữ liệu vào | Cao = đọc quá nhiều; max>>avg = skew đọc |
computeMsAvg / computeMsMax | Thời gian xử lý: avg / max | Pha tính toán thật | Cao = việc nặng; max>>avg = skew tính |
writeMsAvg / writeMsMax | Thời gian ghi output/shuffle: avg / max | Pha ghi ra | Cao = shuffle output lớn; max>>avg = skew ghi |
slotMs | Tổng slot-milliseconds stage tiêu thụ | Tổng công tính toán = (số slot × thời gian) | Stage slotMs lớn nhất = "kẻ ngốn slot" cần tối ưu trước |
recordsRead | Số bản ghi đọc vào stage | Khối lượng dòng đầu vào | So với stage trước để thấy dữ liệu phình/co |
recordsWritten | Số bản ghi ghi ra | Khối lượng dòng đầu ra | recordsWritten >> recordsRead = join làm nổ dòng |
shuffleOutputBytes | Bytes ghi vào shuffle để đẩy sang stage sau | Lượng dữ liệu trung chuyển qua mạng | Rất lớn = shuffle nặng, tốn network + bộ nhớ |
shuffleOutputBytesSpilled | Bytes shuffle phải tràn ra đĩa vì hết bộ nhớ | Chỉ báo áp lực bộ nhớ | > 0 là cảnh báo: thiếu bộ nhớ / shuffle quá lớn |
parallelInputs | Số đơn vị đầu vào có thể xử song song (≈ số worker tối đa) | Mức song song lý thuyết của stage | Quá nhỏ = không tận dụng được song song |
completedParallelInputs | Số đơn vị đầu vào đã xử lý xong | So với parallelInputs để biết tiến độ | Chênh lớn khi đang chạy = còn nhiều việc dở |
steps[].kind | Loại thao tác trong stage: READ, WRITE, COMPUTE, FILTER, SORT, AGGREGATE, JOIN, ANALYTIC_FUNCTION, LIMIT, CROSS_JOIN… | Cho biết stage làm gì | CROSS_JOIN bất ngờ = join thiếu điều kiện |
Đơn vị: các field
...Ms...tính bằng milliseconds;slotMslà slot-milliseconds (khác hẳn — nó nhân thêm số slot).shuffleOutputBytes*tính bằng bytes.records*là số dòng.
avg vs max — hai con số phát hiện data skew
Điểm tinh tế nhất khi đọc stage: mỗi pha có cả avg lẫn max theo worker, và khoảng cách giữa chúng nói lên phân bố tải.
computeMsAvg= worker trung bình xử lý trong bao lâu.computeMsMax= worker lâu nhất xử lý trong bao lâu.
Nếu cả đội chia việc đều, avg và max gần bằng nhau. Nhưng nếu một worker bị dồn phần dữ liệu lớn bất thường (ví dụ một customer_id chiếm 30% giao dịch, hoặc một khóa join lệch), worker đó làm mãi không xong trong khi cả đội đã rảnh — max vọt lên rất cao so với avg. Đó chính là data skew (dữ liệu lệch). Cả stage phải chờ worker chậm nhất, nên max mới là thứ quyết định thời gian thực của stage, không phải avg.
Quy tắc ngón tay cái: nếu ...Max của một pha lớn hơn ...Avg của pha đó vài lần trở lên (ví dụ ≥ 3–4×), nghi ngay data skew ở pha đó. Skew ở computeMs thường do khóa join/group lệch; skew ở readMs do một số worker phải đọc phần dữ liệu lớn hơn hẳn.
Panel thông số một stage — mô phỏng có chú thích
Dưới đây là cách một stage hiện ra khi bạn mở Execution Details (console) hoặc đọc từ job_stages. Con số là mẫu minh hoạ để tập đọc:
┌────────────────────────────────────────────────────────────────────────────┐
│ STAGE S03: Join+ status: COMPLETE │
├────────────────────────────────────────────────────────────────────────────┤
│ inputStages: [S01, S02] ← nhận shuffle từ 2 stage con (xem bq-13) │
│ │
│ PHA avg max ← đọc trục thời gian tương đối bên dưới │
│ ───────────────────────────────── │
│ WAIT 8 ms 15 ms ← chờ slot: NHỎ, không nghẽn tài nguyên │
│ READ 40 ms 55 ms ← đọc shuffle vào: bình thường │
│ COMPUTE 120 ms 460 ms 🔥 ← max ≈ 3.8x avg → DATA SKEW ở compute │
│ WRITE 25 ms 30 ms ← ghi shuffle ra: đều │
│ │
│ slotMs (tổng công): 2,480,000 slot-ms ← stage ngốn slot nhất │
│ recordsRead: 84,000,000 │
│ recordsWritten: 152,000,000 ← nổ dòng: join làm phình│
│ shuffleOutputBytes: 9.4 GB │
│ shuffleOutputBytesSpilled: 1.8 GB 🔥 ← >0: TRÀN ĐĨA, thiếu RAM │
│ parallelInputs: 2,000 │
│ completedParallelInputs: 2,000 ← đã xong 100% │
│ │
│ steps: READ → JOIN(INNER) → AGGREGATE → WRITE │
├────────────────────────────────────────────────────────────────────────────┤
│ CHẨN ĐOÁN NHANH: │
│ • COMPUTE max>>avg → skew ở khóa join/group (1 khóa nóng dồn 1 worker) │
│ • spilled 1.8GB → shuffle vượt bộ nhớ, tràn đĩa → chậm thêm │
│ • recordsWritten > recordsRead → join làm nổ dòng, cần soi điều kiện join │
└────────────────────────────────────────────────────────────────────────────┘
Ba luật đọc quan trọng nhất
Ba tình huống này chiếm phần lớn các ca chẩn đoán thực tế. Học thuộc phản xạ đọc chúng:
Nếu WAIT cao thì…
Stage tiêu nhiều thời gian ở pha chờ slot, tức nó không thiếu tính toán mà thiếu tài nguyên. Nguyên nhân thường gặp: project đang chạy on-demand chạm trần slot, hoặc reservation quá nhỏ so với tải, hoặc nhiều query nặng chạy đồng thời tranh slot. Cách xử lý không nằm ở viết lại SQL mà ở tài nguyên & lịch: tăng/điều chỉnh reservation, giãn lịch các job nặng, giảm concurrency giờ cao điểm. Lưu ý: một chút WAIT ở đầu stage là bình thường (worker mới được cấp slot); chỉ lo khi WAIT chiếm tỷ trọng lớn trong tổng thời gian stage.
Nếu COMPUTE max >> avg thì…
Đây là chữ ký kinh điển của data skew. Worker trung bình xong nhanh nhưng một (vài) worker ôm phần dữ liệu lệch nên kéo dài cả stage. Nguồn skew phổ biến trong dữ liệu ngân hàng: một merchant_id/customer_id "khủng" chiếm phần lớn giao dịch; khóa join có nhiều giá trị NULL bị gom về một worker; GROUP BY trên cột phân bố cực lệch. Hướng gỡ (chi tiết ở bq-15): tách/lọc riêng khóa nóng, thêm bước tiền tổng hợp, xử lý NULL trước khi join, hoặc thay đổi khóa phân phối. Đừng cố tối ưu avg — thời gian thực của stage bị max quyết định.
Nếu spilled > 0 thì…
shuffleOutputBytesSpilled > 0 nghĩa là dữ liệu shuffle không đủ chỗ trong bộ nhớ nên phải tràn ra đĩa — đọc/ghi đĩa chậm hơn RAM nhiều lần, làm stage chậm hẳn. Hai nguyên nhân gốc thường đi cùng nhau: (1) shuffle quá lớn (join/sort/group trên khối dữ liệu khổng lồ chưa được lọc), và/hoặc (2) skew khiến một worker phải giữ phần dữ liệu vượt bộ nhớ của nó. Hướng gỡ: lọc và cắt cột sớm hơn để giảm dữ liệu vào bước shuffle, giảm cardinality trước khi join/sort, xử lý skew như trên; nếu dùng reservation, nhiều slot hơn cũng đồng nghĩa nhiều bộ nhớ shuffle hơn. Nguyên tắc: spilled > 0 là cờ đỏ luôn đáng điều tra, kể cả khi query vẫn ra kết quả.
Timeline / thời gian tương đối của stage
Ngoài con số từng pha, console còn vẽ timeline tương đối: mỗi stage là một thanh trên trục thời gian của cả job, cho thấy stage nào chạy khi nào, cái nào chồng lấn, cái nào phải chờ cái trước. Đọc timeline giúp phân biệt hai kiểu chậm: chậm vì một stage nặng (một thanh dài) hay chậm vì chuỗi phụ thuộc dài (nhiều thanh nối đuôi, mỗi cái chờ cái trước xong). Kết hợp timeline (khi nào) với bốn pha (tắc ở đâu trong mỗi stage) cho bức tranh đầy đủ. Cách dựng cây phụ thuộc từ inputStages xem lại bq-13.
Đọc thông số stage trực tiếp bằng SQL
Bạn không cần console — có thể moi job_stages từ INFORMATION_SCHEMA.JOBS và tự tính các chỉ số chẩn đoán. Câu dưới tìm những stage nghi skew (compute max lớn hơn avg nhiều lần) và có spill trong các job gần đây:
-- Soi stage nghi data skew hoặc tràn đĩa trong 1 ngày qua
SELECT
job_id,
stage.id AS stage_id,
stage.name AS stage_name,
stage.slot_ms AS slot_ms,
stage.compute_ms_avg AS compute_avg_ms,
stage.compute_ms_max AS compute_max_ms,
SAFE_DIVIDE(stage.compute_ms_max,
stage.compute_ms_avg) AS compute_skew_ratio, -- >~4 = nghi skew
stage.wait_ms_avg AS wait_avg_ms,
stage.records_read AS records_read,
stage.records_written AS records_written,
stage.shuffle_output_bytes AS shuffle_bytes,
stage.shuffle_output_bytes_spilled AS spilled_bytes -- >0 = tràn đĩa
FROM `region-asia-southeast1`.INFORMATION_SCHEMA.JOBS_BY_PROJECT AS job,
UNNEST(job.job_stages) AS stage
WHERE job.creation_time > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 1 DAY)
AND job.job_type = 'QUERY'
AND job.state = 'DONE'
AND (
SAFE_DIVIDE(stage.compute_ms_max, stage.compute_ms_avg) > 4 -- skew
OR stage.shuffle_output_bytes_spilled > 0 -- spill
)
ORDER BY stage.slot_ms DESC
LIMIT 50;
Trong job_stages, tên field ở dạng snake_case (compute_ms_avg, shuffle_output_bytes_spilled…) tương ứng camelCase ở REST API (computeMsAvg, shuffleOutputBytesSpilled…) — cùng một thông số, chỉ khác quy ước đặt tên theo giao diện. Sắp theo slot_ms giảm dần để ưu tiên stage ngốn công nhất trước.
Use case thực tế
Đội Data NCB có một job hằng đêm gộp transactions với customers để tính hạn mức rủi ro theo khách. Job chạy ~14 phút, chậm bất thường. Mở Execution Details, hầu hết stage lành, nhưng stage join S03 nổi bật:
computeMsAvg = 120 msnhưngcomputeMsMax = 460 ms→ max ≈ 3.8× avg, chữ ký data skew ở compute.shuffleOutputBytesSpilled = 1.8 GB(> 0) → shuffle vượt bộ nhớ, tràn đĩa.recordsWritten (152M) > recordsRead (84M)→ join làm nổ dòng.slotMscủaS03chiếm ~70% tổng job → đây là "kẻ ngốn slot".
Điều tra sâu: một nhóm nhỏ khách hàng doanh nghiệp (~0.1%) có customer_id chiếm gần 40% giao dịch, dồn vào vài worker → vừa gây skew compute vừa làm shuffle của các worker đó phồng vượt RAM → spill. Ba can thiệp: (1) tiền tổng hợp giao dịch theo customer_id ở một stage trước rồi mới join, giảm mạnh dòng vào bước join; (2) xử lý riêng nhóm khách "khủng" bằng một nhánh query tách; (3) lọc cột và ngày sớm hơn để thu nhỏ shuffle. Kết quả: computeMsMax về sát avg, spilled về 0, slotMs của S03 giảm ~4×, job từ ~14 phút xuống ~3.5 phút. Quy trình chẩn đoán đầy đủ theo bước ở BigQuery — Chẩn đoán & tinh chỉnh.
Ghi nhớ
- Mỗi stage có bốn pha theo worker: WAIT (chờ slot), READ (đọc input/shuffle), COMPUTE (xử lý), WRITE (ghi output/shuffle). Đọc tỷ lệ pha để biết nút thắt nằm ở đâu.
- Mỗi pha có cả avg và max;
max >> avg(≥ ~3–4×) = data skew — và max mới quyết định thời gian thực của stage, không phải avg. slotMslà tổng công tính toán (slot × thời gian); stage cóslotMslớn nhất là nơi cần tối ưu trước tiên.shuffleOutputBytesSpilled > 0= shuffle tràn đĩa vì thiếu bộ nhớ/shuffle quá lớn — luôn là cờ đỏ đáng điều tra.recordsWritten >> recordsReadgợi ý join làm nổ dòng;parallelInputsvscompletedParallelInputscho biết mức song song và tiến độ.- WAIT cao → thiếu tài nguyên (slot/reservation/concurrency), xử lý ở tầng tài nguyên chứ không phải viết lại SQL.
- COMPUTE max >> avg → skew khóa join/group hoặc
NULLdồn worker; tách khóa nóng, tiền tổng hợp, xử lýNULL. - Có thể đọc mọi thông số này bằng SQL qua
UNNEST(job_stages)trongINFORMATION_SCHEMA.JOBSmà không cần console.
Nguồn tham khảo
- Query plan and timeline (execution details, các pha & field của stage): https://cloud.google.com/bigquery/docs/query-plan-explanation
- INFORMATION_SCHEMA JOBS (job_stages, các cột thống kê): https://cloud.google.com/bigquery/docs/information-schema-jobs
- BigQuery documentation: https://cloud.google.com/bigquery/docs
- Best practices — performance overview (data skew, shuffle, tối ưu): https://cloud.google.com/bigquery/docs/best-practices-performance-overview
- Reservations / slots (Editions) — slot & bộ nhớ shuffle: https://cloud.google.com/bigquery/docs/reservations-intro
- "Dremel: Interactive Analysis of Web-Scale Datasets" — Melnik et al., VLDB 2010
- "Google BigQuery: The Definitive Guide" — Lakshmanan & Tigani, O'Reilly
Bài viết liên quan
So sánh các định dạng dữ liệu (CSV, JSON, XML, Avro, Parquet, ORC) và lý do lưu theo cột nhanh hơn cho phân tích. Bài đi sâu vào row vs columnar storage, nén (Snappy/gzip/zstd), schema evolution, OLTP vs OLAP, object storage và partitioning để tối ưu chi phí lẫn tốc độ truy vấn.
Data Engineering là ngành xây dựng và vận hành hệ thống biến dữ liệu thô thành dữ liệu sạch, tin cậy, sẵn sàng cho phân tích và AI. Bài giới thiệu vai trò Data Engineer trong vòng đời dữ liệu (nguồn → ingestion → storage → transformation → serving), phân biệt với Analyst/Scientist/ML Engineer, bức tranh hệ sinh thái công cụ và bài toán đưa dữ liệu core banking sang kho phân tích.
Stream processing là gì, khác biệt batch vs stream (bounded/unbounded), micro-batch (Spark) vs true streaming, Apache Flink là gì và định vị so với Spark Structured Streaming và Kafka Streams. Kiến trúc runtime JobManager/TaskManager, các tầng API, triết lý streaming-first và bối cảnh phát hiện gian lận ngân hàng.
Vì sao một máy không đủ và cần xử lý phân tán: từ MapReduce, Hadoop/HDFS đến Apache Spark in-memory. Kiến trúc driver–executor–cluster manager, các mức trừu tượng RDD/DataFrame/Dataset, cơ chế lazy evaluation với DAG, và vì sao shuffle (wide dependency) là phần tốn kém nhất. Kèm PySpark, Spark SQL và các kỹ thuật tối ưu (partition, broadcast join, cache, chống skew) cùng khi nào KHÔNG nên dùng Spark.
Cảm nhận của bạn
Bình luận
Chưa có bình luận. Hãy là người đầu tiên chia sẻ!