BigQuery 15 — Chẩn đoán & tinh chỉnh qua Query Plan
Từ "đọc được plan" sang "sửa được query"
Ba bài trước dạy bạn nhìn: Execution Details và cấu trúc plan, cách đọc execution graph, và ý nghĩa từng thông số theo pha của stage. Nhưng đọc được một biểu đồ chỉ có giá trị khi nó dẫn tới một hành động: query này chậm ở đâu, vì sao, và sửa thế nào cho nhanh và rẻ hơn.
Bài này là bước ghép nối. Ta lấy chính các field thật của query plan — computeMsAvg/computeMsMax, waitMsAvg/waitMsMax, shuffleOutputBytes, shuffleOutputBytesSpilled, recordsRead, parallelInputs — và biến chúng thành một bộ triệu chứng lâm sàng. Mỗi triệu chứng trỏ tới một nguyên nhân gốc, mỗi nguyên nhân có một toa thuốc cụ thể. Mục tiêu: khi mở một plan chậm, bạn không còn "đoán mò" mà chạy qua một checklist chẩn đoán như bác sĩ đọc phim X-quang.
Mô hình tinh thần: một query chậm luôn có một stage nghẽn cổ chai (bottleneck). Việc của bạn không phải tối ưu mọi stage, mà tìm đúng stage tốn nhiều nhất (
slotMslớn nhất hoặc nằm trên đường tới hạn), phân loại nó tốn ở pha nào (wait / read / compute / write), rồi áp toa tương ứng. Sửa sai stage = tốn công mà đồng hồ không nhúc nhích.
Quy trình chẩn đoán 5 bước
Trước khi vào bảng, hãy nắm một luồng cố định để không bỏ sót. Luôn bắt đầu từ đắt nhất và đi ngược lên nguyên nhân.
Điểm mấu chốt của luồng: phân biệt "tốn ở pha nào". Cùng một stage chậm, nếu chậm vì waitMs thì đó là bài toán tài nguyên (slot), còn nếu chậm vì computeMs với max >> avg thì đó là bài toán phân phối dữ liệu (skew) — hai hướng sửa hoàn toàn khác nhau. Chi tiết cách đọc các pha này nằm ở bài thông số stage.
Bảng chẩn đoán: triệu chứng → nguyên nhân → khắc phục
Đây là phần lõi của bài. Mỗi hàng là một "ca bệnh" đọc được trực tiếp từ query plan.
| # | Triệu chứng (đọc từ plan) | Nguyên nhân gốc | Cách khắc phục |
|---|---|---|---|
| 1 | Trong một stage, computeMsMax lớn hơn nhiều lần computeMsAvg (ví dụ max gấp 10–50× avg); một vài worker chạy mãi trong khi phần lớn đã xong | Data skew — một hoặc vài giá trị key chiếm phần lớn số dòng, dồn hết vào ít worker | Tách key nóng ra xử lý riêng; thêm cột phụ để salting/rebalance (băm key nóng thành nhiều nhóm); gộp trước (pre-aggregate) để giảm số dòng theo key trước khi JOIN/GROUP BY |
| 2 | shuffleOutputBytesSpilled > 0 (khác 0), thường kèm writeMs cao ở stage sinh shuffle | Shuffle quá lớn — dữ liệu trung gian vượt bộ nhớ worker, phải tràn xuống đĩa | Giảm dữ liệu trước khi shuffle: đẩy WHERE/FILTER lên sớm, chỉ SELECT cột cần, pre-aggregate trước JOIN; nếu vẫn tràn thì tăng slot (reservation) để có nhiều bộ nhớ song song hơn |
| 3 | waitMsAvg/waitMsMax chiếm phần lớn thời gian stage; parallelInputs lớn nhưng completedParallelInputs nhích chậm | Đói slot — không đủ slot để chạy hết các đơn vị song song, chúng xếp hàng chờ | Cấp thêm slot qua reservation (baseline/autoscale); giảm concurrency (giãn lịch, tách workload BI khỏi ETL); xếp query nặng ra khung giờ thấp điểm |
| 4 | Một stage READ có recordsRead / bytes đọc rất lớn so với dữ liệu query thực sự cần; thường là stage lá (input) | Thiếu prune — quét cả bảng vì không lọc được theo partition/cluster | Thêm điều kiện lọc trên cột partition để cắt phần lớn dữ liệu (partitioning); sắp xếp bảng theo cluster cho cột hay lọc/JOIN (clustering); tránh hàm bọc quanh cột partition làm mất prune |
| 5 | Stage JOIN có shuffle khổng lồ ở cả hai nhánh input; slotMs dồn vào stage join, kèm spill | JOIN bảng lớn × bảng lớn — cả hai phía đều phải repartition qua mạng | Nếu một phía nhỏ → để BigQuery broadcast (không shuffle phía lớn); denormalize trước (gộp sẵn cột hay dùng vào bảng fact); lọc cả hai phía sớm; cân nhắc bảng trung gian đã pre-join (kỹ thuật tối ưu) |
| 6 | Plan có rất nhiều stage nhỏ, mỗi stage recordsRead ít, slotMs thấp; tổng overhead điều phối lớn hơn công việc thật | Over-partition — bảng chia quá vụn (quá nhiều partition/file nhỏ), sinh vô số input tí hon | Gộp bớt độ mịn partition (ví dụ theo tháng thay vì theo giờ nếu truy vấn theo tháng); nén file nhỏ; tránh partition trên cột có cardinality quá cao |
Cách dùng bảng: mở Execution Details hoặc truy vấn
INFORMATION_SCHEMA.JOBS, tìm stage cóslotMslớn nhất, đọc bộ 4 pha (wait/read/compute/write) và cặpavg/max, rồi dò xuống hàng khớp. Một query có thể mắc nhiều bệnh cùng lúc — sửa cái đắt nhất trước rồi đo lại.
Đọc dấu hiệu skew ngay trong panel
Skew (hàng #1) là bệnh khó thấy nhất vì tổng thời gian trông "bình thường". Dấu hiệu duy nhất là khoảng cách avg–max trong một stage:
STAGE 03: JOIN + AGGREGATE slotMs: 4,210,880
parallelInputs: 800 completedParallelInputs: 800
----- pha (ms, theo worker) ----- avg max
wait 120 180
read 340 410
compute 1,900 46,700 <-- max ~24x avg = SKEW
write 210 260
records read : 1,240,000,000
shuffle out : 88.4 GB spilled: 0
Ở đây computeMsMax (46.7s) gấp ~24 lần computeMsAvg (1.9s): tuyệt đại đa số worker xong trong 2 giây, nhưng vài worker "ôm" một key nóng chạy tới 47 giây và kéo dài cả stage. Cả bảng field và ngưỡng cảnh báo chi tiết ở thông số từng pha stage.
Truy vấn INFORMATION_SCHEMA để tự động bắt bệnh
Thay vì mở panel từng job, bạn có thể quét hàng loạt job và tính sẵn tỷ số cảnh báo. Đoạn GoogleSQL sau lọc ra các job có dấu hiệu skew hoặc shuffle tràn đĩa trong 1 ngày:
-- Bắt job nghi skew (compute max >> avg) hoặc có shuffle spill
SELECT
job_id,
user_email,
total_bytes_billed / POW(1024, 3) AS gb_billed,
total_slot_ms,
stage.name AS stage_name,
stage.slot_ms AS stage_slot_ms,
stage.compute_ms_avg,
stage.compute_ms_max,
SAFE_DIVIDE(stage.compute_ms_max,
stage.compute_ms_avg) AS skew_ratio,
stage.shuffle_output_bytes_spilled AS spilled_bytes
FROM
`region-us`.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) > 8 -- nghi skew
OR stage.shuffle_output_bytes_spilled > 0 -- shuffle tràn đĩa
)
ORDER BY stage.slot_ms DESC
LIMIT 50;
Bảng kết quả cho bạn danh sách stage cần soi: skew_ratio cao → hàng #1; spilled_bytes > 0 → hàng #2. Cấu trúc job_stages (mảng các stage, UNNEST để bung ra) đã giới thiệu ở bài Execution Details. Đây là "hệ thống cảnh báo sớm" chạy hàng ngày, hợp với vận hành ở bài monitoring & cost.
Case study: dashboard rủi ro tín dụng — before → after
Bối cảnh
Phòng QLRR của NCB có một dashboard tính dư nợ theo chi nhánh trong 30 ngày. Query JOIN bảng transactions (giao dịch, ~2 tỷ dòng, không partition) với customers (khách hàng, ~8 triệu dòng) rồi GROUP BY branch_id. Dashboard mở mất ~40 giây, và một số hôm còn timeout.
Query gốc:
-- BEFORE: quét toàn bảng, join lớn×lớn
SELECT
c.branch_id,
SUM(t.amount) AS total_amount,
COUNT(*) AS n_txn
FROM `ncb.core.transactions` AS t -- 2 tỷ dòng, KHÔNG partition
JOIN `ncb.core.customers` AS c
ON t.customer_id = c.customer_id
WHERE t.txn_date >= '2026-06-15' -- lọc 30 ngày nhưng bảng không partition
GROUP BY c.branch_id;
Đọc plan phát hiện ba bệnh: (a) stage READ transactions quét gần như toàn bộ 2 tỷ dòng vì bảng không partition → không prune được txn_date (hàng #4); (b) stage JOIN shuffle lớn ở cả hai nhánh, shuffleOutputBytesSpilled > 0 (hàng #2 + #5); (c) một số branch_id (chi nhánh lớn Hà Nội) khiến computeMsMax >> avg ở stage AGGREGATE (hàng #1).
Execution graph — before
Cách sửa
Ba thay đổi, mỗi cái đánh vào một bệnh:
- Partition bảng
transactionstheoDATE(txn_date)→ stage READ chỉ quét 30 partition thay vì cả bảng (hàng #4, partitioning). - Cluster theo
customer_idđể JOIN gọn hơn, và pre-aggregate giao dịch theocustomer_idtrước khi JOIN vớicustomers→ phía lớn co lại còn ~8 triệu dòng, JOIN không còn lớn×lớn, hết spill (hàng #2, #5). - Sau pre-aggregate, phân phối theo
branch_idđều hơn vì mỗi customer chỉ còn 1 dòng → giảm skew ở AGGREGATE cuối (hàng #1).
-- AFTER: prune partition + pre-aggregate trước JOIN
WITH txn_by_cust AS (
SELECT
customer_id,
SUM(amount) AS cust_amount,
COUNT(*) AS cust_txn
FROM `ncb.core.transactions` -- ĐÃ partition theo DATE(txn_date)
WHERE txn_date >= '2026-06-15' -- prune: chỉ 30 partition
GROUP BY customer_id -- co 2 tỷ dòng -> ~vài triệu
)
SELECT
c.branch_id,
SUM(a.cust_amount) AS total_amount,
SUM(a.cust_txn) AS n_txn
FROM txn_by_cust AS a
JOIN `ncb.core.customers` AS c -- clustered theo customer_id
ON a.customer_id = c.customer_id
GROUP BY c.branch_id;
Execution graph — after
Kết quả
| Chỉ số (đọc từ plan/job) | Before | After |
|---|---|---|
| Thời gian chạy (elapsed) | ~40 s | ~6 s |
totalBytesBilled | ~1.9 TB | ~140 GB |
shuffleOutputBytesSpilled (stage JOIN) | > 0 (tràn đĩa) | 0 |
computeMsMax / computeMsAvg (AGGREGATE) | ~24× (skew nặng) | ~2× (đều) |
Các con số trên là minh hoạ định tính cho một tình huống điển hình, không phải benchmark chuẩn — hình dạng cải thiện (bytes billed và spill giảm mạnh, skew ratio về gần 1) mới là điều cần nhớ. Trên workload của bạn hãy tự đo lại bằng
INFORMATION_SCHEMA.JOBS.
Use case thực tế
Đội Data Platform NCB áp quy trình này thành một weekly tuning review. Mỗi thứ Hai, một scheduled query quét INFORMATION_SCHEMA.JOBS_BY_PROJECT 7 ngày, lọc top 20 job theo total_slot_ms và gắn cờ tự động: skew_ratio > 8, spilled_bytes > 0, hoặc gb_billed > 500.
Một quý điển hình: 20 job "đắt" nhất chiếm ~70% tổng total_slot_ms của cả project. Sau khi áp bảng chẩn đoán — chủ yếu là thêm partition cho 3 bảng fact chưa partition (hàng #4) và pre-aggregate 2 dashboard join lớn×lớn (hàng #5) — tổng slot-time hàng tuần của nhóm job này giảm khoảng 60%, và số lần dashboard timeout cuối tháng (khi concurrency cao, wait tăng — hàng #3) gần như biến mất sau khi tách reservation BI khỏi ETL. Điểm rút ra: không tối ưu dàn trải mà tập trung vào ít job đắt nhất, đúng theo nguyên tắc "sửa stage bottleneck trước".
Ghi nhớ
- Query chậm luôn có một stage bottleneck: xếp theo
slotMs, sửa cái đắt nhất trước, đừng tối ưu dàn trải. - Phân loại theo pha:
waitMscao = đói slot (tài nguyên);computeMsMax >> avg= data skew (phân phối);spilled > 0= shuffle quá lớn; READ nhiều bytes = thiếu prune. - Skew (max ≫ avg) sửa bằng tách/salting key nóng và pre-aggregate; đói slot sửa bằng reservation + giảm concurrency — hai bệnh khác nhau dù đều làm chậm.
- Thiếu prune (hàng #4) là bệnh phổ biến và dễ sửa nhất: thêm partition + cluster, tránh bọc hàm quanh cột partition.
- JOIN lớn×lớn: pre-aggregate/denormalize để một phía co nhỏ, cho BigQuery broadcast thay vì shuffle cả hai nhánh.
- Over-partition (nhiều stage tí hon) cũng hại: overhead điều phối lớn hơn công việc — gộp bớt độ mịn partition.
- Dùng
INFORMATION_SCHEMA.JOBS+UNNEST(job_stages)để tự động bắt bệnh hàng loạt (skew_ratio,spilled_bytes), không cần mở panel từng job. - Luôn đo lại
totalBytesBilled,totalSlotMs, elapsed vàspilledsau khi sửa — số liệu mới xác nhận toa thuốc, không phải cảm giác.
Nguồn tham khảo
- BigQuery — Query plan and timeline (execution details, stage phases, shuffle, spill): cloud.google.com/bigquery/docs/query-plan-explanation
- BigQuery — Best practices for performance overview (skew, join, prune): cloud.google.com/bigquery/docs/best-practices-performance-overview
- BigQuery — INFORMATION_SCHEMA JOBS (job_stages, total_slot_ms, total_bytes_billed): cloud.google.com/bigquery/docs/information-schema-jobs
- BigQuery — Introduction to partitioned tables: cloud.google.com/bigquery/docs/partitioned-tables
- BigQuery — Introduction to clustered tables: cloud.google.com/bigquery/docs/clustered-tables
- BigQuery — Introduction to reservations (slots, autoscale): cloud.google.com/bigquery/docs/reservations-intro
- BigQuery documentation (tổng quan): cloud.google.com/bigquery/docs
- "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ẻ!