SingleStore 4 — Sharding & truy vấn phân tán
SingleStore 4 — Sharding & truy vấn phân tán
Ở bài kiến trúc ta đã thấy SingleStore là hệ shared-nothing: dữ liệu chia thành partitions trải trên các leaf node, còn aggregator nhận query và điều phối. Bài này trả lời câu hỏi cốt lõi tiếp theo: một dòng cụ thể rơi vào partition nào, và khi hai bảng lớn phải JOIN với nhau thì dữ liệu di chuyển ra sao? Đây là phần quyết định hiệu năng của một cụm phân tán — chọn đúng thì JOIN chạy cục bộ trên leaf, chọn sai thì mỗi truy vấn phải kéo hàng trăm triệu dòng qua mạng.
Ý tưởng nền: trong SingleStore, bạn không đọc/ghi trực tiếp lên partition. Bạn nói chuyện với aggregator như một bảng logic duy nhất; aggregator dùng SHARD KEY để biết dòng nào thuộc partition nào và đẩy công việc xuống đúng leaf (query pushdown). Toàn bộ nghệ thuật thiết kế nằm ở việc chọn shard key sao cho dữ liệu hay dùng chung nằm chung một partition.
Lưu ý: mọi block SQL trong bài là cú pháp SingleStore (tương thích MySQL) để minh hoạ. Sandbox của app là PostgreSQL read-only nên không block nào được đánh dấu "chạy được".
SHARD KEY: từ giá trị tới partition tới leaf
Khi bạn tạo một sharded table, SingleStore băm giá trị của các cột trong SHARD KEY để chọn partition:
partition = hash(giá_trị_shard_key) % số_partition
Số partition được cố định lúc tạo cluster/database (thường là bội số của số core trên các leaf) và không đổi theo từng bảng. Mỗi partition được gán cho một leaf (master), cộng một bản replica trên leaf khác cho HA. Vì hàm băm là hàm tất định, mọi dòng có cùng giá trị shard key luôn rơi vào cùng một partition — đây chính là tính chất ta khai thác để JOIN cục bộ.
Điểm cần nắm: SHARD KEY không phải index để tăng tốc lọc (đó là việc của SORT KEY/hash index — xem Universal Storage). Nó là hàm phân phối dữ liệu. Chọn shard key là chọn "trục" mà dữ liệu được cắt ra và rải lên cụm.
Cú pháp khai báo — ví dụ bảng giao dịch ngân hàng shard theo customer_id:
CREATE TABLE transactions (
txn_id BIGINT NOT NULL,
customer_id BIGINT NOT NULL,
account_id BIGINT NOT NULL,
amount DECIMAL(18,2) NOT NULL,
txn_type VARCHAR(16) NOT NULL,
created_at DATETIME NOT NULL,
SHARD KEY (customer_id), -- băm customer_id → chọn partition
SORT KEY (created_at), -- sắp trong columnstore theo thời gian
KEY USING CLUSTERED COLUMNSTORE (created_at)
);
Nếu không khai SHARD KEY, SingleStore mặc định shard theo PRIMARY KEY; nếu cũng không có primary key thì dùng keyless sharding (băm ngẫu nhiên theo dòng) — dữ liệu vẫn phân bố đều nhưng bạn mất khả năng collocate JOIN vì không dòng nào "cùng khoá" với dòng bảng khác.
Reference table vs sharded table
Không phải bảng nào cũng nên chia nhỏ. Với các dimension table nhỏ (danh mục chi nhánh, loại sản phẩm, mã tỉnh, bảng tỷ giá…) mà hầu như mọi truy vấn đều JOIN vào, SingleStore cho phép nhân bản toàn bộ bảng lên MỌI leaf — gọi là reference table.
| Sharded table | Reference table | |
|---|---|---|
| Dữ liệu nằm ở đâu | Chia thành partition, mỗi leaf giữ một phần | Nhân bản đầy đủ trên mọi leaf (và aggregator) |
| Dùng cho | Bảng fact lớn (transactions, accounts) | Dim table nhỏ, ít đổi (branches, product_types) |
| JOIN với bảng khác | Có thể phải reshuffle/broadcast | Luôn local — bản sao có sẵn ngay trên leaf |
| Ghi (write) | Song song, phân tán trên nhiều partition | Ghi qua master aggregator rồi nhân ra các leaf → đắt hơn, không hợp bảng ghi nhiều |
| Kích thước hợp lý | Bất kỳ | Nhỏ (vài MB–vài trăm MB); lớn quá tốn RAM × số leaf |
Khai báo reference table bằng CREATE REFERENCE TABLE (không có SHARD KEY vì bảng không chia):
CREATE REFERENCE TABLE branches (
branch_id INT NOT NULL PRIMARY KEY,
branch_name VARCHAR(128) NOT NULL,
province VARCHAR(64) NOT NULL
);
Vì mỗi leaf đã có sẵn toàn bộ branches, một JOIN transactions JOIN branches không cần di chuyển dữ liệu branches qua mạng — mỗi partition tự tra bản sao cục bộ. Đây là cách rẻ nhất để JOIN một fact table lớn với nhiều dim nhỏ. Đánh đổi: mỗi lần cập nhật branches phải nhân ra tất cả leaf, và bảng chiếm RAM/đĩa nhân với số leaf — nên chỉ dùng cho bảng nhỏ, ít thay đổi.
Distributed join: dữ liệu di chuyển thế nào
Khi JOIN hai bảng, câu hỏi sống còn là: các dòng cần khớp với nhau có đang nằm cùng partition không? Nếu có, mỗi leaf tự JOIN phần của mình — nhanh và không tốn mạng. Nếu không, SingleStore phải di chuyển dữ liệu để đưa các dòng khớp về cùng chỗ, bằng một trong ba chiến lược.
Collocated join — nhanh nhất
Nếu cả hai bảng shard theo cùng khoá và JOIN đúng trên khoá đó, thì mọi dòng cần khớp đã ở sẵn cùng partition (vì cùng hàm băm, cùng giá trị → cùng partition). SingleStore đẩy JOIN xuống từng leaf, chạy local join song song, không dòng nào rời máy.
Ví dụ: shard transactions theo customer_id và cũng shard customers theo customer_id:
CREATE TABLE customers (
customer_id BIGINT NOT NULL,
full_name VARCHAR(128) NOT NULL,
segment VARCHAR(32) NOT NULL,
SHARD KEY (customer_id) -- CÙNG shard key với transactions
);
-- JOIN này là COLLOCATED: chạy local trên mỗi partition
SELECT c.segment, COUNT(*) AS n_txn, SUM(t.amount) AS total
FROM transactions t
JOIN customers c ON t.customer_id = c.customer_id
WHERE t.created_at >= '2026-01-01'
GROUP BY c.segment;
Vì hash(customer_id) giống nhau ở cả hai bảng, dòng của khách 4711 trong transactions và dòng 4711 trong customers chắc chắn nằm cùng partition. Đây là mục tiêu thiết kế: chọn shard key theo khoá JOIN phổ biến nhất để biến các JOIN đắt nhất thành local.
Reshuffle join — phân phối lại theo join key
Nếu bạn JOIN hai sharded table nhưng join key khác shard key của ít nhất một bảng, dữ liệu khớp không cùng chỗ. SingleStore phải reshuffle: đọc bảng đó, băm lại theo join key, và gửi mỗi dòng tới partition đích tương ứng, rồi mới JOIN. Ví dụ nếu transactions shard theo customer_id nhưng ta JOIN theo account_id:
-- transactions shard theo customer_id, accounts shard theo account_id
-- JOIN theo account_id ≠ shard key của transactions → RESHUFFLE
SELECT a.account_type, SUM(t.amount)
FROM transactions t
JOIN accounts a ON t.account_id = a.account_id
GROUP BY a.account_type;
Reshuffle đúng đắn nhưng đắt: chi phí tỉ lệ với lượng dữ liệu phải chuyển qua mạng nội bộ giữa các leaf. Với bảng fact hàng trăm triệu dòng, đây thường là bottleneck. EXPLAIN/PROFILE (xem thực thi truy vấn) sẽ hiện các toán tử Repartition/Reshuffle — dấu hiệu cần xem lại shard key.
Broadcast join — phát bảng nhỏ
Khi một bên JOIN nhỏ, thay vì reshuffle cả bảng lớn, SingleStore phát (broadcast) nguyên bản nhỏ tới mọi partition, rồi mỗi partition JOIN cục bộ với bản sao đó. Chi phí = kích thước bảng nhỏ × số partition. Rẻ khi bảng nhỏ (một tập lọc, một dim chưa khai là reference), đắt nếu "bảng nhỏ" thực ra lớn. Nếu một dim table luôn nhỏ và hay JOIN, hãy khai nó là reference table ngay từ đầu để có sẵn bản sao trên leaf, khỏi broadcast lặp lại mỗi truy vấn.
Tóm tắt lựa chọn của optimizer:
| Tình huống | Chiến lược | Chi phí mạng |
|---|---|---|
| Hai bảng cùng shard key = join key | Collocated | ~0 |
| JOIN với reference table | Local (bản sao có sẵn) | ~0 |
| Join key ≠ shard key, cả hai bảng lớn | Reshuffle | Cao (∝ dữ liệu) |
| Một bên nhỏ (chưa là reference) | Broadcast | Trung bình (∝ bảng nhỏ × #partition) |
Chọn shard key: tránh skew và reshuffle
Hai "kẻ thù" khi chọn shard key:
1. Data skew (lệch dữ liệu). Nếu shard key có phân bố lệch — ví dụ shard theo txn_type với đa số là 'transfer' — thì một vài partition phình to (hot partition), gánh phần lớn dữ liệu và tải, trong khi các partition khác rảnh. Quét song song mất ý nghĩa vì phải chờ partition nóng nhất. Quy tắc: chọn cột có lực phân biệt (cardinality) cao và phân bố đều, ví dụ customer_id, account_id, hoặc txn_id, tránh cột trạng thái/loại ít giá trị. Với cột có thể lệch, cân nhắc shard theo tổ hợp nhiều cột để tăng độ đều.
2. Reshuffle đắt. Shard key nên trùng với join/GROUP BY key phổ biến nhất để biến các truy vấn nặng thành collocated. Không thể tối ưu mọi JOIN cùng lúc: nếu bảng hay JOIN theo cả customer_id lẫn account_id, bạn chỉ collocate được một trục. Chiến lược thường thấy: shard bảng fact theo khoá JOIN đắt nhất/hay dùng nhất, đưa các dim nhỏ về reference table, và chấp nhận reshuffle cho các JOIN hiếm.
Vài nguyên tắc thực dụng:
- Shard key bất biến (đừng shard theo cột hay bị UPDATE — đổi giá trị nghĩa là dòng phải chuyển partition).
- Ưu tiên collocate cặp bảng lớn × lớn (nơi reshuffle đau nhất); cặp lớn × nhỏ thường xử lý ổn bằng reference/broadcast.
- Kiểm tra phân bố sau khi nạp: so số dòng giữa các partition (management views
information_schema.MV_*/SHOW PARTITIONS) để phát hiện skew sớm. SHARD KEYvàSORT KEYlà hai quyết định độc lập: shard key phân phối dữ liệu ra leaf, sort key sắp xếp bên trong columnstore để lọc/quét nhanh.
Use case thực tế
Bối cảnh (minh hoạ): hệ phân tích khách hàng của NCB có transactions ~1,2 tỷ dòng (18 tháng), customers ~8 triệu dòng, và các dim branches (~120 dòng), product_types (~40 dòng). Dashboard rủi ro chạy liên tục truy vấn "tổng chi tiêu theo phân khúc khách hàng" — JOIN transactions × customers theo customer_id.
Thiết kế:
transactionsvàcustomerscùngSHARD KEY (customer_id)→ JOIN theocustomer_idlà collocated, chạy local trên mỗi partition, không reshuffle 1,2 tỷ dòng.branches,product_typeskhai làREFERENCE TABLE→ có sẵn bản sao trên mọi leaf, JOIN thêm dim không tốn mạng.- Với báo cáo hiếm JOIN theo
account_id, chấp nhận reshuffle — không hy sinh trụccustomer_idvốn dùng 90% thời gian.
Kết quả (minh hoạ): truy vấn tổng hợp theo phân khúc giảm từ ~7 giây (khi customers bị broadcast/reshuffle do lệch shard key) xuống ~1,3 giây sau khi đưa hai bảng về cùng shard key — PROFILE không còn toán tử Reshuffle trên nhánh transactions. Kiểm tra SHOW PARTITIONS cho thấy chênh lệch số dòng giữa các partition dưới 3%, xác nhận customer_id phân bố đều, không hot partition.
-- Sau tối ưu: collocated + reference, không dòng transactions nào rời leaf
SELECT c.segment, b.province,
COUNT(*) AS n_txn, SUM(t.amount) AS total
FROM transactions t
JOIN customers c ON t.customer_id = c.customer_id -- collocated
JOIN branches b ON c.branch_id = b.branch_id -- local (reference)
WHERE t.created_at >= '2026-01-01'
GROUP BY c.segment, b.province;
Ghi nhớ
SHARD KEYlà hàm phân phối, không phải index lọc.hash(shard_key) % Nchọn partition; cùng giá trị → cùng partition (tất định). Không khai thì mặc định theo PRIMARY KEY, không có nữa thì keyless (mất khả năng collocate).- Reference table = nhân bản toàn bộ lên mọi leaf → JOIN luôn local; chỉ dùng cho dim nhỏ, ít ghi. Sharded table chia partition cho fact lớn.
- Collocated join: hai bảng cùng shard key = join key → local, nhanh nhất, ~0 mạng. Đây là đích thiết kế shard key.
- Reshuffle join: join key ≠ shard key → băm lại và chuyển dữ liệu qua mạng; đúng nhưng đắt theo kích thước bảng — bottleneck phổ biến với fact lớn.
- Broadcast join: phát nguyên bảng nhỏ tới mọi partition; rẻ khi nhỏ, đắt khi bảng "nhỏ" thực ra lớn. Dim hay JOIN nên khai reference thay vì để broadcast lặp lại.
- Chống skew: chọn cột cardinality cao, phân bố đều (
customer_id), tránh cột trạng thái/loại; dùng tổ hợp cột nếu cần. Hot partition giết hiệu năng song song. - Chống reshuffle: shard theo khoá JOIN/GROUP BY nặng nhất; không collocate được mọi trục — ưu tiên cặp lớn × lớn.
- Shard key nên bất biến (đừng shard theo cột hay UPDATE).
EXPLAIN/PROFILEhiệnReshuffle/Repartitionlà dấu hiệu xem lại thiết kế.
Nguồn tham khảo
- SingleStore Documentation — "Distributed SQL / Query Execution" (cách phân phối và pushdown truy vấn): docs.singlestore.com
- SingleStore Documentation — "SHARD KEY" và "Sharding" (Physical Database Schema Design): docs.singlestore.com
- SingleStore Documentation — "Reference Tables" (Distributed Data Modeling): docs.singlestore.com
- SingleStore Documentation — "Distributed Joins" / "Optimizing Table Data Structures" (collocated, reshuffle, broadcast): docs.singlestore.com
- SingleStore Documentation — "SORT KEY" và "CREATE TABLE" reference (cú pháp SHARD KEY / SORT KEY / CLUSTERED COLUMNSTORE): docs.singlestore.com
- SingleStore Engineering Blog — bài về distributed query execution và data skew trên cụm phân tán: singlestore.com/blog
- MySQL Documentation — cú pháp
JOIN,GROUP BY(phần tương thích cú pháp): dev.mysql.com
Bài viết liên quan
Kiến trúc shared-nothing của SingleStore: Master Aggregator giữ metadata và điều phối, Child Aggregator scale kết nối, Leaf node chứa dữ liệu chia thành partition. Bài mổ xẻ luồng một query (aggregator nhận → pushdown xuống leaf → gộp kết quả) và cơ chế High Availability master/replica, failover, redundancy level.
Điểm khác biệt lớn nhất của SingleStore: nó BIÊN DỊCH truy vấn ra mã máy (code generation) rồi chạy song song MPP trên leaf, thay vì diễn giải từng dòng. Bài dựng luồng SQL → tối ưu → sinh mã → plan biên dịch, giải thích plan cache tái dùng (biên dịch 1 lần), query pushdown xuống leaf và aggregator gộp; cách đọc EXPLAIN/PROFILE (thời gian, rows, bộ nhớ, network/reshuffle) và SHOW PLANCACHE để nhận diện reshuffle/broadcast, tối ưu truy vấn.
SingleStore (tiền thân MemSQL) là database quan hệ phân tán HTAP, tương thích giao thức MySQL, gộp OLTP và OLAP trong một hệ thống. Bài mở màn dựng mô hình tinh thần về HTAP, giải thích vấn đề nó giải quyết (tránh ETL sang warehouse riêng), chỉ rõ khi nào NÊN và KHÔNG NÊN dùng, định vị so với PostgreSQL, ClickHouse, BigQuery, TiDB/CockroachDB, và vẽ bản đồ toàn series 10 bài.
Khoá chính/ngoại/tổng hợp, ràng buộc (NOT NULL, UNIQUE, CHECK, FK) và cách mô hình hoá quan hệ 1:1, 1:n, n:n cho hệ khách hàng — tài khoản — giao dịch. Đi qua chuẩn hoá 1NF/2NF/3NF bằng ví dụ trước/sau cụ thể, rồi bàn khi nào nên cố tình phi chuẩn hoá để đọc nhanh — giúp thiết kế lược đồ đúng ngay từ đầu.
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ẻ!