Chương 6 — Partitioning
Phần II — Distributed Data

Chương 6 — Partitioning

12 phút đọc DDIA · Martin Kleppmann

🎯 Mục tiêu chương: Hiểu cách chia một dataset lớn thành nhiều partition (theo key range hoặc hash), cách secondary index tương tác với partitioning, các chiến lược rebalancing khi thêm/bớt node, và cách request được route tới đúng partition.

Tổng quan: Partitioning là gì?

Khi dataset quá lớn hoặc query throughput quá cao, replication thôi không đủ — cần chia dữ liệu thành partitions (còn gọi sharding). Mỗi record/row/document thuộc đúng một partition; mỗi partition giống như một database nhỏ độc lập (dù DB có thể hỗ trợ thao tác chạm nhiều partition).

Thuật ngữ gây nhầm lẫn — cùng một khái niệm, nhiều tên:

Hệ thốngTên gọi của "partition"
MongoDB, Elasticsearch, SolrCloudshard
HBaseregion
Bigtabletablet
Cassandra, Riakvnode
CouchbasevBucket

Lưu ý: partitioning ở đây (chủ động chia DB) không liên quan tới network partition (lỗi mạng chia cắt node — chương 8).

Lý do chính: scalability. Trong kiến trúc shared-nothing, các partition đặt trên các node khác nhau → dữ liệu trải trên nhiều disk, query load trải trên nhiều CPU. Query chạm một partition thì mỗi node tự xử lý độc lập → thêm node là tăng throughput. Query lớn phức tạp có thể parallelize nhưng khó hơn nhiều.

Lịch sử: Teradata, Tandem NonStop SQL (1980s); được NoSQL và Hadoop-based data warehouse "tái khám phá". Nguyên lý áp dụng cho cả workload transactional và analytics.

Partitioning and Replication

Partitioning thường kết hợp với replication: mỗi partition có bản sao trên nhiều node để chịu lỗi. Một node có thể chứa nhiều partition. Với leader–follower, mỗi partition có leader ở một node và followers ở các node khác → mỗi node có thể là leader của vài partition và follower của các partition khác.

Mọi điều ở chương 5 áp dụng nguyên vẹn cho replication của partition. Việc chọn partitioning scheme gần như độc lập với replication scheme, nên chương này bỏ qua replication để đơn giản.

Partitioning of Key-Value Data

Mục tiêu: chia đều dữ liệu và query load giữa các node. Lý tưởng: 10 node xử lý gấp 10 lần dữ liệu và throughput của 1 node.

  • Skew: partition không đều, có partition nhiều dữ liệu/query hơn → partitioning kém hiệu quả. Cực đoan: mọi load dồn vào 1 partition, 9/10 node ngồi chơi.
  • Hot spot: partition có load cao bất thường.
  • Gán record ngẫu nhiên vào node thì chia đều, nhưng khi đọc không biết item nằm đâu → phải query tất cả node song song. Có cách tốt hơn khi truy cập theo primary key.

Partitioning by Key Range

Gán mỗi partition một khoảng key liên tục (min → max), như các tập của bộ bách khoa toàn thư giấy (tập 1: A–B, …, tập 12: T–Z). Biết ranh giới → biết key nằm ở partition nào → gửi thẳng tới node đúng.

  • Ranh giới không cần đều nhau vì dữ liệu không phân bố đều (chia mỗi tập 2 chữ cái sẽ có tập to tập nhỏ). Ranh giới phải thích nghi với dữ liệu — do admin chọn tay hoặc DB tự chọn.
  • Dùng bởi: Bigtable, HBase, RethinkDB, MongoDB trước 2.4.
  • Trong mỗi partition giữ key đã sắp xếp (SSTable/LSM) → range scan dễ, có thể coi key như concatenated index để lấy nhiều record liên quan trong một query.

Ví dụ sensor: key = timestamp (year-month-day-hour-minute-second) → range scan "mọi reading trong tháng X" rất tiện. Nhưng hot spot: mọi write hiện tại đều rơi vào partition "hôm nay", các partition khác rảnh. - Cách sửa: prefix timestamp bằng tên sensor → partition theo sensor rồi mới theo thời gian → write trải đều (khi nhiều sensor cùng hoạt động). Đổi lại: muốn lấy nhiều sensor trong một khoảng thời gian phải chạy một range query riêng cho mỗi sensor.

Partitioning by Hash of Key

Để tránh skew và hot spot, nhiều datastore dùng hash function quyết định partition. Hash tốt biến dữ liệu lệch thành phân bố đều: hàm 32-bit trả số "ngẫu nhiên" trong 0 … 2³² − 1, ngay cả với chuỗi input gần giống nhau.

  • Không cần mạnh về mật mã: Cassandra, MongoDB dùng MD5, Voldemort dùng Fowler–Noll–Vo.
  • Cạm bẫy: hash built-in của ngôn ngữ như Java Object.hashCode() hay Ruby Object#hash có thể cho giá trị khác nhau giữa các process → không dùng được cho partitioning.
  • Mỗi partition sở hữu một khoảng hash (không phải khoảng key). Ranh giới có thể chia đều, hoặc chọn giả ngẫu nhiên (khi đó đôi khi gọi là consistent hashing).

Consistent hashing (Karger et al.): ban đầu là cách chia tải cho hệ thống cache toàn internet như CDN, dùng ranh giới ngẫu nhiên để không cần điều phối trung tâm hay consensus. Chữ "consistent" không liên quan tới replica consistency hay ACID consistency — nó mô tả một cách rebalancing. Cách này thực ra không hoạt động tốt cho database nên ít dùng thực tế; tài liệu nhiều DB dùng thuật ngữ này không chính xác → sách khuyên gọi là hash partitioning.

Cái giá: mất range query hiệu quả. Key từng kề nhau giờ rải khắp partition, mất thứ tự. - MongoDB (hash-based sharding): range query phải gửi tới mọi partition. - Riak, Couchbase, Voldemort: không hỗ trợ range query trên primary key.

Thoả hiệp của Cassandra — compound primary key: chỉ phần đầu của key được hash để chọn partition; các cột còn lại làm concatenated index để sắp xếp dữ liệu trong SSTable. Không range được trên cột đầu, nhưng nếu cố định cột đầu thì range scan hiệu quả trên các cột sau. - Ví dụ mạng xã hội: key (user_id, update_timestamp) → lấy mọi update của một user trong khoảng thời gian, sắp theo thời gian, từ một partition duy nhất. User khác nhau nằm ở partition khác nhau. Mô hình đẹp cho quan hệ one-to-many.

Tiêu chíKey range partitioningHash partitioningCompound key (Cassandra)
Phân bố tảiDễ lệch, dễ hot spot nếu key tuần tự (timestamp)Đều hơn nhiềuĐều theo phần hash
Range queryHiệu quả (key sắp xếp)Không hiệu quả — scatter tới mọi partition hoặc không hỗ trợHiệu quả trên cột sau khi cố định cột đầu
Ranh giớiThích nghi theo dữ liệu (tay hoặc tự động)Khoảng hash, chia đều hoặc ngẫu nhiênTheo hash của partition key
Rebalancing điển hìnhDynamic partitioning (split/merge)Fixed number of partitions (hoặc dynamic)Theo cách của hash
Ví dụBigtable, HBase, RethinkDB, MongoDB trước 2.4Cassandra, MongoDB (hash mode), Riak, Couchbase, VoldemortCassandra

Skewed Workloads and Relieving Hot Spots

Hash giảm hot spot nhưng không loại bỏ hoàn toàn: nếu mọi read/write cùng nhắm một key, chúng vẫn vào cùng partition (hash của hai ID giống nhau vẫn giống nhau).

Ví dụ: celebrity với hàng triệu follower làm gì đó → bão write vào cùng key (user ID của celebrity, hoặc ID của hành động đang được comment).

Hầu hết hệ thống không tự bù được skew kiểu này → trách nhiệm của application: - Kỹ thuật đơn giản: thêm số ngẫu nhiên vào đầu hoặc cuối key hot. Chỉ 2 chữ số thập phân đã chia write của một key ra 100 key khác nhau, trải trên nhiều partition. - Trade-off: read phải làm thêm việc — đọc cả 100 key rồi gộp lại. - Cần bookkeeping: chỉ nên làm với số ít key hot (làm với mọi key là overhead vô ích) → phải theo dõi key nào đang bị split.

Partitioning and Secondary Indexes

Các scheme trên dựa vào mô hình key-value: truy cập qua primary key → biết partition. Secondary index phức tạp hơn vì nó thường không định danh duy nhất record mà là để tìm các record có một giá trị: mọi hành động của user 123, mọi bài chứa từ "hogwash", mọi xe màu đỏ…

  • Secondary index là "cơm ăn hằng ngày" của RDBMS, phổ biến trong document DB.
  • Nhiều key-value store (HBase, Voldemort) tránh secondary index vì phức tạp; Riak bắt đầu thêm vì hữu ích.
  • Là lý do tồn tại của search server như Solr, Elasticsearch.

Vấn đề: secondary index không map gọn vào partition. Hai cách: document-based và term-based.

Partitioning Secondary Indexes by Document (local index)

Ví dụ: website bán xe cũ. Mỗi listing có document ID; partition theo ID (0–499 ở partition 0, 500–999 ở partition 1…). Muốn lọc theo color và make → secondary index. Khi thêm xe đỏ, partition tự thêm ID vào index entry color:red.

  • Mỗi partition hoàn toàn tách biệt, tự duy trì secondary index cho chỉ những document của nó. → gọi là local index.
  • Write đơn giản: chỉ đụng tới partition chứa document ID đang ghi.
  • Read tốn kém: xe đỏ nằm rải ở cả partition 0 và 1 → phải gửi query tới mọi partition rồi gộp kết quả → scatter/gather.
    • Dù query song song, dễ bị tail latency amplification (chậm bằng partition chậm nhất).
  • Dù vậy dùng rất rộng rãi: MongoDB, Riak, Cassandra, Elasticsearch, SolrCloud, VoltDB.
  • Vendor khuyên thiết kế partitioning sao cho query secondary index phục vụ được từ một partition — nhưng không phải lúc nào cũng làm được, nhất là khi một query dùng nhiều secondary index cùng lúc (lọc cả color lẫn make).
  • Lưu ý: nếu DB chỉ có key-value và bạn tự tạo index (map value → document IDs) ở tầng application, rất dễ bị lệch do race condition và write thất bại giữa chừng → cần multi-object transaction.

Partitioning Secondary Indexes by Term (global index)

Xây global index bao phủ dữ liệu mọi partition. Nhưng không thể để index trên một node (thành bottleneck) → global index cũng phải được partition, nhưng theo cách khác với primary key.

  • Ví dụ: color:red chứa xe đỏ từ mọi partition; index được chia sao cho màu bắt đầu bằng a–r ở partition 0, s–z ở partition 1. Index theo make chia tương tự (ranh giới giữa f và h).
  • Gọi là term-partitioned vì term (ví dụ color:red) quyết định partition của index. Tên "term" từ full-text index: term là các từ xuất hiện trong document.
  • Có thể partition theo term trực tiếp (hỗ trợ range scan, ví dụ giá xe) hoặc theo hash của term (phân bố tải đều hơn).

Trade-off: - ✅ Read hiệu quả: chỉ hỏi partition chứa term cần tìm, không scatter/gather. - ❌ Write chậm và phức tạp: ghi một document có thể ảnh hưởng nhiều partition của index (mỗi term có thể ở partition/node khác). - Giữ index luôn up-to-date đòi hỏi distributed transaction qua mọi partition bị ảnh hưởng — không phải DB nào cũng hỗ trợ. - Thực tế cập nhật global index thường asynchronous: đọc index ngay sau write có thể chưa thấy thay đổi. Ví dụ Amazon DynamoDB: global secondary index cập nhật trong phần nhỏ của giây bình thường, nhưng có thể trễ lâu hơn khi hạ tầng gặp sự cố. - Dùng bởi: DynamoDB GSI, Riak search, Oracle data warehouse (cho chọn local hoặc global).

Tiêu chíDocument-partitioned (local index)Term-partitioned (global index)
Index nằm ở đâuCùng partition với documentPartition riêng, chia theo term (hoặc hash term)
WriteNhanh, chỉ 1 partitionChậm, có thể chạm nhiều partition; thường cập nhật async
Read theo secondary keyScatter/gather mọi partition → tail latencyChỉ 1 partition (cho mỗi term)
Consistency của indexĐồng bộ với document (cùng partition)Có thể trễ (eventual) nếu không có distributed transaction
Ví dụMongoDB, Riak, Cassandra, Elasticsearch, SolrCloud, VoltDBDynamoDB GSI, Riak search, Oracle DW

Rebalancing Partitions

Theo thời gian mọi thứ thay đổi: throughput tăng (cần thêm CPU), dữ liệu tăng (cần thêm disk/RAM), máy chết (máy khác phải gánh). Quá trình di chuyển load từ node này sang node khác gọi là rebalancing.

Yêu cầu tối thiểu: - Sau rebalancing, load (dữ liệu, read, write) chia công bằng giữa các node. - Trong khi rebalancing, DB vẫn nhận read và write. - Không di chuyển nhiều dữ liệu hơn cần thiết — để nhanh và giảm tải network/disk I/O.

Strategies for Rebalancing

How not to do it: hash mod N

Tại sao không dùng hash(key) mod N (N = số node)? Vì khi N thay đổi, hầu hết key phải di chuyển. Ví dụ hash(key) = 123456: - 10 node → node 6 (123456 mod 10) - 11 node → node 3 (123456 mod 11) - 12 node → node 0 (123456 mod 12)

Di chuyển liên tục như vậy khiến rebalancing cực đắt. Cần cách chỉ di chuyển dữ liệu khi thật cần.

Fixed number of partitions

Tạo nhiều partition hơn số node rất nhiều ngay từ đầu, mỗi node giữ nhiều partition. Ví dụ: 10 node, 1.000 partition → ~100 partition/node. - Thêm node: node mới "trộm" vài partition từ mỗi node cũ cho đến khi cân bằng. Bớt node: ngược lại. - Chỉ di chuyển nguyên partition; số partition và việc gán key → partition không đổi, chỉ đổi gán partition → node. - Chuyển dữ liệu tốn thời gian → trong lúc đó vẫn dùng assignment cũ cho read/write. - Có thể tính tới phần cứng không đồng đều: node mạnh nhận nhiều partition hơn. - Dùng bởi: Riak, Elasticsearch, Couchbase, Voldemort. - Số partition thường cố định từ lúc setup; nhiều hệ không cài đặt split/merge vì phức tạp → số partition ban đầu = số node tối đa có thể có. Phải chọn đủ lớn cho tăng trưởng, nhưng mỗi partition có overhead quản lý nên không chọn quá lớn. - Khó khi kích thước dataset biến động lớn: partition quá to → rebalancing và recovery đắt; quá nhỏ → overhead nhiều. Hiệu năng tốt nhất khi partition "vừa phải".

Dynamic partitioning

Với key range partitioning, số partition cố định với ranh giới cố định rất bất tiện (chọn sai ranh giới → mọi dữ liệu dồn vào một partition). Nên HBase, RethinkDB tạo partition động: - Partition vượt ngưỡng kích thước (HBase mặc định 10 GB) → split làm đôi. Dữ liệu bị xoá nhiều, partition nhỏ dưới ngưỡng → merge với partition kề. Giống cách B-tree tách node. - Sau khi split, một nửa có thể chuyển sang node khác để cân tải (HBase chuyển file qua HDFS). - ✅ Số partition thích nghi với khối lượng dữ liệu: ít dữ liệu → ít partition, overhead nhỏ; nhiều dữ liệu → mỗi partition vẫn bị giới hạn kích thước. - ❌ DB rỗng bắt đầu với một partition duy nhất → khi còn nhỏ, mọi write do một node xử lý, các node khác rảnh. Giảm thiểu bằng pre-splitting (HBase, MongoDB) — nhưng với key range cần biết trước phân bố key. - Dùng được cho cả hash partitioning: MongoDB từ 2.4 hỗ trợ cả key-range và hash, đều split động.

Partitioning proportionally to nodes

Với dynamic partitioning, số partition ∝ kích thước dữ liệu; với fixed, kích thước partition ∝ kích thước dữ liệu. Cả hai: số partition độc lập với số node.

Cách thứ ba (Cassandra, Ketama): số partition tỉ lệ với số node — cố định số partition mỗi node. - Node không đổi → partition lớn dần theo dữ liệu; thêm node → partition nhỏ lại. Vì dữ liệu nhiều thường đi kèm nhiều node → kích thước partition khá ổn định. - Node mới join: chọn ngẫu nhiên một số partition có sẵn để split, lấy một nửa mỗi partition. Ngẫu nhiên có thể không công bằng, nhưng trung bình trên nhiều partition (Cassandra mặc định 256 partition/node) thì node mới nhận phần tải công bằng. Cassandra 3.0 có thuật toán tránh split không công bằng. - Ranh giới ngẫu nhiên đòi hỏi hash-based partitioning. Đây là cách gần nhất với định nghĩa gốc của consistent hashing.

Chiến lượcSố partitionKích thước partitionƯu điểmNhược điểmVí dụ
hash mod N= N node—Đơn giảnĐổi N là gần như mọi key phải di chuyển(Không nên dùng)
Fixed number of partitionsCố định từ đầu (≫ số node)∝ kích thước datasetĐơn giản vận hành; chỉ di chuyển nguyên partitionKhó chọn số ban đầu khi dataset biến động; giới hạn số node tối đaRiak, Elasticsearch, Couchbase, Voldemort
Dynamic partitioning∝ kích thước datasetGiữ trong khoảng min–max (HBase 10 GB)Thích nghi với khối lượng dữ liệuDB rỗng chỉ có 1 partition → cần pre-splittingHBase, RethinkDB, MongoDB
Proportional to nodes∝ số node (cố định/node)Khá ổn địnhTự điều chỉnh khi thêm nodeSplit ngẫu nhiên có thể không đều; cần hash partitioningCassandra (256/node), Ketama

Operations: Automatic or Manual Rebalancing

Có một dải từ hoàn toàn tự động tới hoàn toàn thủ công. Ở giữa: Couchbase, Riak, Voldemort tự sinh đề xuất assignment nhưng cần admin commit.

  • Tự động: tiện, ít việc vận hành; nhưng khó đoán. Rebalancing đắt (reroute request, chuyển nhiều dữ liệu) → có thể quá tải mạng/node, ảnh hưởng request khác.
  • Nguy hiểm khi kết hợp automatic failure detection: một node quá tải phản hồi chậm → các node khác tưởng nó chết → tự rebalance dời tải khỏi nó → thêm tải cho chính node đó, các node khác và mạng → tình hình tệ hơn, có thể cascading failure.
  • Vì vậy có con người trong vòng lặp là điều tốt: chậm hơn nhưng tránh bất ngờ vận hành.

Request Routing

Dữ liệu đã chia trên nhiều node — client muốn đọc/ghi key "foo" thì kết nối IP/port nào? Assignment thay đổi khi rebalance, nên cần ai đó theo dõi. Đây là trường hợp của bài toán tổng quát service discovery.

Ba cách tiếp cận: 1. Client gửi tới node bất kỳ (qua round-robin load balancer). Node sở hữu partition thì xử lý; không thì forward tới node đúng rồi trả kết quả. 2. Mọi request qua routing tier trước — tầng này không xử lý request, chỉ là partition-aware load balancer. 3. Client tự biết partitioning và assignment → kết nối thẳng node đúng.

Vấn đề cốt lõi: thành phần ra quyết định routing học về thay đổi assignment thế nào? Mọi bên phải đồng thuận, nếu không request đi sai node. Consensus khó cài đúng (chương 9).

CáchMô tảVí dụ
Coordination service (ZooKeeper)Node đăng ký vào ZooKeeper; ZK giữ mapping partition → node có thẩm quyền; routing tier/client subscribe và được thông báo khi thay đổiHBase, SolrCloud, Kafka; LinkedIn Espresso (Helix trên ZooKeeper)
Config server riêng + routing daemonKiến trúc tương tự nhưng tự cài đặtMongoDB (config server + mongos)
Gossip protocolCác node lan truyền thay đổi trạng thái cluster cho nhau; request gửi tới node bất kỳ rồi forward (cách 1). Phức tạp hơn ở node nhưng không phụ thuộc dịch vụ ngoàiCassandra, Riak
Routing tier đơn giảnKhông tự rebalance nên thiết kế đơn giản; routing tier học thay đổi từ nodeCouchbase (moxi)

Với routing tier hoặc gửi tới node ngẫu nhiên, client vẫn cần biết IP để kết nối — nhưng IP ít thay đổi hơn assignment nên thường dùng DNS là đủ.

Parallel Query Execution

Phần lớn NoSQL chỉ hỗ trợ query đơn giản: đọc/ghi một key (cộng scatter/gather cho document-partitioned index). Ngược lại, MPP (massively parallel processing) relational DB cho analytics hỗ trợ query phức tạp: join, filter, group, aggregate. MPP query optimizer chia query thành nhiều stage và partition, chạy song song trên các node — đặc biệt lợi cho query quét phần lớn dataset. Chi tiết ở chương 10.

Tổng kết

  • Partitioning cần thiết khi dữ liệu quá lớn cho một máy. Mục tiêu: chia đều dữ liệu và load, tránh hot spot.
  • Hai cách chính: key range (range query hiệu quả, rủi ro hot spot, thường rebalance bằng split động) và hash (phá thứ tự key, tải đều hơn, thường fixed number of partitions; cũng dùng dynamic được). Hybrid: compound key.
  • Secondary index: local (document-partitioned) — ghi 1 partition, đọc scatter/gather; global (term-partitioned) — ghi nhiều partition, đọc 1 partition.
  • Routing: từ load balancer partition-aware đơn giản đến parallel query engine phức tạp.
  • Các partition vận hành phần lớn độc lập — đó là lý do scale được; nhưng thao tác ghi nhiều partition (một thành công, một thất bại?) rất khó reasoning → chương sau.

⚠️ Hiểu lầm & cạm bẫy thường gặp

  • Nhầm partitioning với network partition — hai khái niệm hoàn toàn khác.
  • "Consistent hashing" = consistency — không; nó chỉ là một cách rebalancing, và thực tế ít dùng đúng nghĩa trong DB. Nên gọi là hash partitioning.
  • Dùng hash(key) mod N → thêm/bớt một node là gần như toàn bộ dữ liệu phải di chuyển.
  • Dùng hash built-in của ngôn ngữ (Java hashCode(), Ruby Object#hash) → khác nhau giữa các process.
  • Dùng timestamp làm leading key trong key range partitioning → mọi write vào partition "hiện tại" → hot spot.
  • Nghĩ hash giải quyết mọi hot spot — một key cực nóng (celebrity) vẫn dồn vào một partition; phải salt/split key ở application.
  • Salt mọi key — overhead read vô ích; chỉ nên làm cho số ít hot key và phải theo dõi.
  • Quên chi phí scatter/gather của local secondary index → tail latency amplification khi số partition tăng.
  • Kỳ vọng global secondary index cập nhật đồng bộ — thường async (DynamoDB GSI), đọc ngay sau write có thể chưa thấy.
  • Chọn số fixed partitions quá nhỏ → giới hạn số node tối đa; quá lớn → overhead quản lý.
  • Quên rằng dynamic partitioning bắt đầu với 1 partition → cần pre-splitting.
  • Bật rebalancing tự động + failure detection tự động không cẩn thận → cascading failure.

💼 Áp dụng thực tế & phỏng vấn

  • Chọn shard key là câu hỏi trung tâm trong system design: phân bố đều không? query phổ biến có chạm một shard không? có hot key không? Ví dụ: chat app shard theo conversation_id, timeline shard theo user_id với sort key là timestamp (mô hình (user_id, timestamp) như Cassandra/DynamoDB partition key + sort key).
  • Time-series: tránh timestamp làm leading key; dùng (sensor_id/device_id, timestamp) hoặc bucket theo thời gian + salt.
  • Hot key / celebrity problem (Twitter, Instagram): key salting (thêm suffix 0–99), cache phía trước, hoặc xử lý riêng tài khoản lớn.
  • Secondary index trong DynamoDB: LSI = local (cùng partition key), GSI = global (async, eventual). Elasticsearch/Cassandra secondary index = local → query theo field không phải partition key sẽ scatter/gather.
  • Rebalancing: Elasticsearch phải chọn số primary shard lúc tạo index (fixed partitions); HBase/MongoDB split chunk tự động; Cassandra dùng vnodes (256 token/node). Redis Cluster dùng 16384 hash slots — một ví dụ fixed number of partitions.
  • Request routing: Kafka/HBase dựa ZooKeeper (Kafka mới dùng KRaft), MongoDB dùng mongos + config server, Cassandra dùng gossip + client driver token-aware.
  • Tránh mod N khi thiết kế cache phân tán — nhắc tới consistent hashing / hash slots để giảm lượng dữ liệu di chuyển khi thêm node.
  • Nhắc trade-off cross-shard query và cross-shard transaction — thiết kế data model để phần lớn query nằm trong một partition.

❓ Câu hỏi ôn tập

Bấm vào câu hỏi để xem đáp án
1So sánh key range partitioning và hash partitioning về phân bố tải và range query.
Key range giữ key sắp xếp nên range scan hiệu quả nhưng dễ hot spot nếu key tuần tự (timestamp). Hash phân bố đều nhưng phá thứ tự key, range query phải gửi tới mọi partition hoặc không được hỗ trợ.
2Trong database sensor dùng timestamp làm key, vấn đề là gì và sửa thế nào? Cái giá của cách sửa?
Mọi write rơi vào partition của thời điểm hiện tại → hot spot. Prefix key bằng tên sensor để write trải đều; cái giá là lấy nhiều sensor trong một khoảng thời gian phải chạy một range query cho mỗi sensor.
3Compound primary key của Cassandra giải quyết trade-off gì? Cho ví dụ.
Hash phần đầu key để chọn partition (tải đều), phần sau dùng để sắp xếp trong partition (range scan). Ví dụ (user_id, update_timestamp): lấy mọi update của một user trong một khoảng thời gian từ một partition.
4Làm sao xử lý hot key như tài khoản celebrity? Trade-off?
Thêm số ngẫu nhiên (ví dụ 2 chữ số) vào key hot để chia write ra 100 key trên nhiều partition. Đổi lại read phải đọc cả 100 key rồi gộp, và cần bookkeeping để biết key nào đang bị split.
5So sánh local (document-partitioned) và global (term-partitioned) secondary index.
Local: index nằm cùng partition với document, write chỉ chạm 1 partition, read phải scatter/gather mọi partition (tail latency). Global: index partition theo term, read chỉ hỏi 1 partition, nhưng write chạm nhiều partition và thường cập nhật async.
6Vì sao không dùng hash mod N để gán key cho node? Giải pháp fixed number of partitions hoạt động thế nào?
Khi N đổi, gần như mọi key đổi node (123456 → node 6, 3, 0 với N=10, 11, 12). Fixed partitions: tạo nhiều partition hơn số node (ví dụ 1.000 cho 10 node), node mới "trộm" nguyên partition từ node cũ; gán key → partition không đổi.
7Ưu nhược của dynamic partitioning và partitioning proportionally to nodes?
Dynamic: split khi quá ngưỡng (HBase 10 GB), merge khi quá nhỏ, số partition thích nghi dữ liệu; nhưng DB rỗng chỉ có 1 partition → cần pre-splitting. Proportional: cố định số partition mỗi node (Cassandra 256), node mới split ngẫu nhiên partition có sẵn, kích thước partition ổn định; cần hash partitioning, split có thể không đều.
8Ba cách route request tới đúng partition, và làm sao hệ thống biết được assignment hiện tại?
Gửi tới node bất kỳ rồi forward; qua routing tier partition-aware; client tự biết partition. Assignment được theo dõi bằng coordination service như ZooKeeper (HBase, Kafka, SolrCloud, Espresso), config server (MongoDB), hoặc gossip protocol giữa các node (Cassandra, Riak).

Đây là bản tóm tắt và ghi chú, không thay thế sách gốc. Hãy ủng hộ tác giả bằng cách đọc bản gốc.