BigQuery 7 — Join, Window & CTE

15 thg 7, 2026 3 lượt xem
#join
#data-engineering
#bigquery
#shuffle
#window-function
#cte

Vì sao một câu JOIN có thể rẻ hoặc đắt gấp trăm lần

Hai câu truy vấn cho ra cùng một kết quả nhưng một câu chạy 2 giây, câu kia chạy 4 phút và ngốn gấp 50 lần slot-time. Khác biệt hầu như luôn nằm ở một chỗ: cách engine thực thi JOIN — cụ thể là nó có phải shuffle (xáo trộn dữ liệu qua mạng) hay không, và nếu có thì shuffle bao nhiêu.

BigQuery là hệ thống phân tán: dữ liệu nằm rải trên hàng trăm worker (slot). Khi bạn JOIN hai bảng, engine phải đưa các dòng "khớp key" về cùng một worker để so khớp. Có hai chiến lược để làm việc đó, và chọn sai (hoặc để engine buộc phải chọn chiến lược đắt) là nguyên nhân số một khiến truy vấn phân tích chậm.

Bài này xây mô hình tinh thần về broadcast join vs hash join, giải thích vì sao thứ tự và kích thước bảng ảnh hưởng tới lượng shuffle, rồi bổ sung ba công cụ viết truy vấn gọn: window function, CTE (WITH)subquery. Cuối cùng là data skew — kẻ thù âm thầm khi bạn JOIN theo một key lệch.

Mô hình tinh thần: JOIN là bài toán "gom cùng key về một chỗ"

Hãy tưởng tượng bảng transactions (giao dịch, hàng tỉ dòng) nằm rải trên 200 worker, và bảng customers (khách hàng) cũng vậy. Để tính transactions JOIN customers ON customer_id, mỗi dòng giao dịch cần được ghép với đúng dòng khách hàng có cùng customer_id. Nhưng dòng giao dịch của khách A đang nằm ở worker 7, còn dòng khách A trong customers lại nằm ở worker 42. Chúng phải "gặp nhau".

Có hai cách để chúng gặp nhau:

  1. Phát bản sao bảng nhỏ cho mọi worker — nếu customers đủ nhỏ, engine gửi toàn bộ nó tới từng worker. Mỗi worker giữ nguyên phần transactions của mình (không di chuyển) và tra cứu cục bộ. Đây là broadcast join.
  2. Xáo cả hai bảng theo key — nếu cả hai bảng đều lớn, engine băm (hash) mỗi dòng theo customer_id và gửi tới worker phụ trách khoảng hash đó. Sau khi xáo, mọi dòng cùng customer_id (từ cả hai bảng) đều nằm chung một worker. Đây là hash join (còn gọi shuffle join).

Khác biệt cốt lõi: broadcast không shuffle bảng lớn, còn hash join shuffle cả hai. Shuffle là thao tác ghi dữ liệu ra tầng lưu trữ trung gian rồi đọc lại qua mạng Jupiter — đó là phần đắt nhất của hầu hết truy vấn.

Broadcast join vs Hash (shuffle) join

Broadcast join (BigQuery gọi nội bộ là broadcast JOIN) chỉ khả thi khi một vế đủ nhỏ để nhét vừa bộ nhớ mỗi worker. Ưu điểm: bảng lớn không nhúc nhích — không shuffle, không ghi trung gian, cực nhanh. Trong query plan bạn sẽ thấy bảng lớn được đọc và join ngay trong một stage, shuffleOutputBytes thấp.

Hash join dùng khi cả hai vế đều lớn. Engine chèn một pha shuffle: mỗi bảng được đọc, băm theo key, và ghi ra shuffle output. Stage sau đọc lại các bucket đã băm và thực hiện join cục bộ. Ở đây shuffleOutputBytes lớn, và nếu bộ nhớ worker không đủ, một phần bị tràn ra đĩa — bạn thấy shuffleOutputBytesSpilled > 0, dấu hiệu cần chú ý. Cách đọc các thông số này chi tiết ở bài đọc shuffle & data skew trong execution graph.

Điểm mấu chốt cần nhớ:

Broadcast tránh được shuffle bảng lớn. Hash join phải shuffle cả hai. BigQuery tự chọn dựa trên ước lượng kích thước — nhưng ước lượng có thể sai, và cách bạn viết truy vấn ảnh hưởng tới lựa chọn đó.

Vì sao thứ tự và kích thước bảng ảnh hưởng shuffle

BigQuery tối ưu hoá thứ tự join dựa trên thống kê, nhưng có một best practice kinh điển vẫn còn giá trị: đặt bảng lớn nhất trước, các bảng nhỏ dần về sau trong mệnh đề FROM/JOIN. Lý do:

  • Engine có xu hướng chọn vế nhỏ hơn làm phía được băm/broadcast. Nếu bạn viết small JOIN big, tối ưu hoá có thể vẫn sửa lại, nhưng khi thống kê thiếu chính xác (bảng mới nạp, không có metadata tốt), đặt bảng lớn trước giúp engine ra quyết định đúng ngay từ đầu.
  • Với JOIN nhiều bảng, thứ tự quyết định kích thước kết quả trung gian. Nếu bạn join hai bảng lớn trước rồi mới lọc bằng bảng nhỏ, kết quả trung gian khổng lồ phải shuffle. Nếu lọc sớm bằng bảng nhỏ (giảm số dòng trước), lượng dữ liệu cần shuffle ở bước sau nhỏ hẳn.

Nguyên tắc gộp lại: giảm dữ liệu càng sớm càng tốt trước khi shuffle.

  • Đẩy WHERE (đặc biệt là lọc theo cột phân vùng/cụm) xuống trước JOIN để mỗi vế bé đi.
  • Chỉ SELECT các cột cần dùng — cột thừa cũng bị shuffle theo, làm phình shuffleOutputBytes.
  • Nếu một vế đủ nhỏ, giữ nó nhỏ (đừng vô tình JOIN nó với bảng khác trước) để engine chọn broadcast.

Một hệ quả kiến trúc quan trọng: nếu tránh được JOIN thì tránh. Với dữ liệu có quan hệ 1-nhiều ổn định (một giao dịch có nhiều dòng chi tiết), mô hình nested & repeated cho phép lưu dữ liệu con lồng trong dòng cha, biến một JOIN tốn shuffle thành thao tác đọc cục bộ. Xem nested & repeated để tránh join.

SQL thật: ba cách join và ảnh hưởng tới shuffle

-- Ví dụ 1: JOIN có khả năng dùng broadcast (bảng chi nhánh rất nhỏ)
-- customers lớn, branches ~vài trăm dòng => engine thường broadcast branches
SELECT
  c.customer_id,
  c.full_name,
  b.branch_name,
  b.region
FROM `ncb-dwh.core.customers`      AS c          -- bảng lớn: đặt trước
JOIN `ncb-dwh.core.branches`       AS b          -- bảng nhỏ: phía được phát
  ON c.branch_id = b.branch_id
WHERE c.status = 'ACTIVE';

-- Ví dụ 2: JOIN hai bảng lớn => hash join (shuffle theo customer_id)
-- Lọc sớm bằng phân vùng ngày để giảm dữ liệu TRƯỚC khi shuffle
SELECT
  t.customer_id,
  c.segment,
  COUNT(*)                AS so_giao_dich,
  SUM(t.amount)           AS tong_tien
FROM `ncb-dwh.core.transactions` AS t
JOIN `ncb-dwh.core.customers`    AS c
  ON t.customer_id = c.customer_id
WHERE t.txn_date BETWEEN DATE '2026-06-01' AND DATE '2026-06-30'  -- cắt phân vùng, giảm shuffle
  AND c.status = 'ACTIVE'
GROUP BY t.customer_id, c.segment;

Ở Ví dụ 2, mệnh đề WHERE t.txn_date BETWEEN ... cực kỳ quan trọng: nó cắt bớt phân vùng ngày trước pha shuffle, nên lượng dữ liệu băm theo customer_id nhỏ đi tương ứng. Bỏ điều kiện này, bạn shuffle cả năm dữ liệu — cùng kết quả logic nhưng chi phí gấp bội.

Window functions: tính toán "theo nhóm" mà không cần JOIN lại

Window function (hàm cửa sổ) cho phép tính giá trị tổng hợp trên một nhóm dòng liên quan nhưng vẫn giữ nguyên từng dòng — khác GROUP BY (gộp nhiều dòng thành một). Đây là cách gọn để trả lời "giao dịch này là giao dịch thứ mấy của khách?", "số dư luỹ kế đến thời điểm này?", "giao dịch này chiếm bao nhiêu % tổng chi tiêu tháng của khách?".

Cấu trúc: hàm() OVER (PARTITION BY ... ORDER BY ... [khung dòng]).

  • PARTITION BY chia dữ liệu thành các cửa sổ (ví dụ mỗi khách một cửa sổ) — về mặt thực thi giống một pha phân nhóm, cần shuffle theo cột partition.
  • ORDER BY sắp thứ tự bên trong mỗi cửa sổ (ví dụ theo thời gian) để các hàm như ROW_NUMBER, LAG, SUM(...) OVER (... ORDER BY ...) (luỹ kế) hoạt động.

Trong query plan, window function xuất hiện dưới step có kind = ANALYTIC_FUNCTION, thường kèm một pha SORT (cho ORDER BY trong cửa sổ) và shuffle theo cột PARTITION BY. Vì vậy một cột partition lệch (skew) cũng làm chậm window function y như làm chậm JOIN.

-- Với mỗi khách: đánh số thứ tự giao dịch, số dư chi tiêu luỹ kế,
-- và tỉ trọng từng giao dịch trên tổng chi tiêu tháng của khách đó.
SELECT
  customer_id,
  txn_id,
  txn_ts,
  amount,
  ROW_NUMBER() OVER (
    PARTITION BY customer_id ORDER BY txn_ts
  )                                                   AS thu_tu_giao_dich,
  SUM(amount) OVER (
    PARTITION BY customer_id ORDER BY txn_ts
    ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
  )                                                   AS chi_tieu_luy_ke,
  ROUND(
    amount / SUM(amount) OVER (PARTITION BY customer_id) * 100, 2
  )                                                   AS ti_trong_pct
FROM `ncb-dwh.core.transactions`
WHERE txn_date BETWEEN DATE '2026-06-01' AND DATE '2026-06-30';

CTE (WITH) và subquery: chia nhỏ truy vấn cho dễ đọc

CTE (Common Table Expression, mệnh đề WITH) đặt tên cho một truy vấn con để tái sử dụng và đọc từ trên xuống thay vì lồng nhiều tầng subquery. Đây thuần tuý là cú pháp cho con người — BigQuery vẫn nội tuyến (inline) CTE vào kế hoạch thực thi, nên CTE không tự động cache kết quả.

Lưu ý thực thi quan trọng:

  • Nếu một CTE được tham chiếu nhiều lần, BigQuery có thể tính lại nó mỗi lần dùng (không mặc định vật chất hoá). Nếu CTE đó nặng và dùng lại 3-4 chỗ, cân nhắc ghi ra bảng tạm hoặc dùng CREATE TEMP TABLE để tính một lần.
  • CTE giúp bạn sắp xếp logic để lọc và thu gọn dữ liệu sớm — gián tiếp giảm shuffle ở các bước join phía sau.
-- CTE làm rõ pipeline: lọc giao dịch tháng -> gộp theo khách -> join thông tin khách
WITH giao_dich_thang AS (
  SELECT customer_id, amount
  FROM `ncb-dwh.core.transactions`
  WHERE txn_date BETWEEN DATE '2026-06-01' AND DATE '2026-06-30'
),
tong_hop_khach AS (
  SELECT customer_id,
         COUNT(*)      AS so_gd,
         SUM(amount)   AS tong_chi
  FROM giao_dich_thang
  GROUP BY customer_id
)
SELECT
  c.customer_id,
  c.segment,
  th.so_gd,
  th.tong_chi
FROM tong_hop_khach AS th
JOIN `ncb-dwh.core.customers` AS c   -- join sau khi đã gộp => vế trái đã nhỏ
  ON th.customer_id = c.customer_id
WHERE th.tong_chi > 50000000
ORDER BY th.tong_chi DESC;

Subquery (truy vấn con lồng trong WHERE/FROM/SELECT) làm được nhiều việc tương tự nhưng khó đọc khi lồng sâu. Một dạng đáng chú ý là subquery tương quan (correlated) — subquery tham chiếu cột của truy vấn ngoài; BigQuery thường viết lại nó thành JOIN, nhưng đôi khi tạo pha thực thi tốn kém. Khi có thể, ưu tiên viết bằng JOIN hoặc window function tường minh để engine tối ưu tốt hơn.

Data skew: khi một key "gánh" cả cluster

Cả hash join lẫn window function đều dựa vào giả định: sau khi băm theo key, dữ liệu phân bố đều trên các worker. Thực tế ngân hàng thường vi phạm giả định này. Ví dụ:

  • Một tài khoản tổng (nostro/suspense) nhận hàng chục triệu bút toán trong khi tài khoản khách thường chỉ vài chục.
  • JOIN theo merchant_id mà một siêu thị lớn chiếm 30% tổng giao dịch.
  • Giá trị NULL ở cột key: mọi dòng NULL băm về cùng một bucket.

Khi đó, một worker phải xử lý lượng dữ liệu khổng lồ trong khi các worker khác ngồi chơi. Toàn bộ stage phải chờ worker chậm nhất (straggler).

Dấu hiệu skew trong execution details rất rõ: ở stage bị lệch, computeMsMax (hoặc readMs/writeMs) lớn hơn hẳn giá trị ...Avg tương ứng. Chênh lệch avg vs max lớn = có worker phải làm việc nặng bất thường. Cách đọc và định lượng các pha này ở bài các pha stage & chỉ số skew.

Cách giảm skew khi join/window:

  • Lọc bỏ hoặc tách riêng key nóng: xử lý NULL và các key khổng lồ bằng nhánh truy vấn riêng rồi UNION ALL.
  • Pre-aggregate: gộp bớt trước khi join (như CTE tong_hop_khach ở trên) để giảm số dòng của key nóng.
  • Salting (thêm hậu tố ngẫu nhiên vào key rồi gộp lại) trong trường hợp cực đoan — đánh đổi độ phức tạp lấy phân bố đều hơn.
  • Nếu quan hệ ổn định, nested & repeated loại bỏ hẳn JOIN gây skew.

Use case thực tế

Đội phân tích rủi ro NCB cần báo cáo tháng: với mỗi khách hàng active, tính tổng chi tiêu, số giao dịch, thứ hạng chi tiêu trong phân khúc (segment), và tỉ trọng so với trung bình phân khúc. Bảng transactions ~1,8 tỉ dòng/năm, customers ~8 triệu dòng.

Bản viết đầu tiên JOIN thẳng transactions với customers rồi mới lọc tháng và tính window trên toàn bộ, không cắt phân vùng: truy vấn quét ~420 GB, shuffleOutputBytes rất lớn, có shuffleOutputBytesSpilled > 0, và một stage có computeMsMax gấp ~15 lần computeMsAvg (skew ở nhóm tài khoản doanh nghiệp lớn). Thời gian ~3,5 phút.

Bản tối ưu: (1) lọc txn_date theo tháng ngay trong CTE đầu để cắt phân vùng — chỉ còn ~34 GB được đọc; (2) pre-aggregate theo customer_id trước, rồi mới JOIN customers (vế trái đã thu nhỏ còn vài triệu dòng); (3) tính RANK() OVER (PARTITION BY segment ORDER BY tong_chi DESC) trên kết quả đã gộp. Kết quả: bytes billed giảm ~12 lần, không còn spill, thời gian còn ~18 giây. Cùng một con số nghiệp vụ, chi phí giảm hơn một bậc độ lớn — khác biệt hoàn toàn nằm ở giảm dữ liệu trước khi shufflejoin sau khi đã gộp.

Ghi nhớ

  • JOIN = "gom cùng key về một worker". Broadcast phát bảng nhỏ đi khắp nơi, không shuffle bảng lớn; hash join shuffle cả hai bảng theo hash(key).
  • BigQuery tự chọn chiến lược theo ước lượng kích thước, nhưng cách bạn viết truy vấn (thứ tự bảng, lọc sớm, chỉ chọn cột cần) định hình lượng shuffle.
  • Giảm dữ liệu trước khi shuffle: đẩy WHERE/cắt phân vùng xuống trước JOIN; pre-aggregate rồi mới join; đặt bảng lớn trước.
  • Window function tính theo nhóm mà giữ nguyên dòng (OVER/PARTITION BY/ORDER BY); về thực thi vẫn shuffle theo cột PARTITION BY, hiện dưới kind = ANALYTIC_FUNCTION.
  • CTE (WITH) là cú pháp cho con người, BigQuery inline chứ không tự cache; CTE nặng dùng lại nhiều lần nên vật chất hoá ra bảng tạm.
  • Data skew (key lệch, NULL, tài khoản tổng) làm một worker gánh cả stage; nhận biết qua computeMsMax >> computeMsAvg và xử lý bằng tách key nóng / pre-aggregate / salting.
  • Nếu quan hệ 1-nhiều ổn định, cân nhắc nested & repeated để tránh JOIN và shuffle hẳn.

Nguồn tham khảo

  • BigQuery documentation — Query syntax (JOIN, WITH, subquery): cloud.google.com/bigquery/docs
  • Best practices — performance overview (join order, giảm dữ liệu trước shuffle): cloud.google.com/bigquery/docs/best-practices-performance-overview
  • Query plan and timeline / execution details (shuffleOutputBytes, các pha waitMs/readMs/computeMs/writeMs): cloud.google.com/bigquery/docs/query-plan-explanation
  • INFORMATION_SCHEMA JOBS (job_stages, total_bytes_billed, total_slot_ms): cloud.google.com/bigquery/docs/information-schema-jobs
  • Nested & repeated fields (tránh join): cloud.google.com/bigquery/docs/nested-repeated
  • "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.

13 thg 7, 2026 10

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.

13 thg 7, 2026 9

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.

13 thg 7, 2026 8

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.

13 thg 7, 2026 7

Cảm nhận của bạn

Bình luận

Bạn cần để viết bình luận.

Chưa có bình luận. Hãy là người đầu tiên chia sẻ!