Chương 9 — Consistency and Consensus
🎯 Mục tiêu chương: Tìm các abstraction có guarantee mạnh (linearizability, total order broadcast, consensus) để ứng dụng có thể "quên" bớt các lỗi của hệ phân tán (mất gói, clock lệch, node pause/crash). Hiểu giới hạn của những gì làm được và không làm được, và vì sao nhiều bài toán thực tế (leader election, unique constraint, distributed commit, lock) thực chất đều quy về consensus.
Ý tưởng xuyên suốt giống Chương 7 (transactions): cài đặt một abstraction tổng quát một lần cho đúng, rồi để ứng dụng dựa vào guarantee của nó. Abstraction quan trọng nhất ở đây là consensus — làm cho mọi node đồng ý về một điều gì đó. Ví dụ kinh điển: failover trong single-leader replication; nếu hai node cùng nghĩ mình là leader (split brain) thì rất dễ mất dữ liệu.
Consistency Guarantees
- Hầu hết DB replicated chỉ cho eventual consistency: ngừng ghi và chờ "đủ lâu" thì mọi replica hội tụ về cùng giá trị. Tên hợp lý hơn là convergence.
- Đây là guarantee rất yếu: không nói khi nào hội tụ. Trước đó read có thể trả về bất cứ gì (ví dụ ghi xong đọc lại ngay mà không thấy giá trị mình vừa ghi).
- Vấn đề với lập trình viên: DB trông như một biến đọc/ghi được, nhưng semantics phức tạp hơn nhiều. Bug do eventual consistency thường chỉ lộ ra khi có network fault hoặc concurrency cao, nên rất khó test.
- Consistency model mạnh hơn thì dễ dùng đúng hơn, nhưng đổi lại performance kém hơn hoặc kém fault-tolerant hơn.
- Phân biệt: transaction isolation (Chương 7) chủ yếu là tránh race condition giữa các transaction chạy song song; distributed consistency là phối hợp trạng thái giữa các replica khi có delay và fault. Có chồng lấn nhưng phần lớn độc lập.
Lộ trình chương: linearizability → ordering (causality, total order) → distributed transactions & consensus.
Linearizability
Ý tưởng: làm cho hệ thống trông như chỉ có một bản sao dữ liệu duy nhất và mọi thao tác lên nó đều atomic. Tên khác: atomic consistency, strong consistency, immediate consistency, external consistency.
Hệ quả: ngay khi một client ghi thành công, mọi client đọc sau đó phải thấy giá trị mới. Linearizability là một recency guarantee — không được đọc từ cache/replica cũ.
Ví dụ Alice & Bob (World Cup 2014): Alice và Bob ngồi cùng phòng, cùng xem kết quả chung kết trên điện thoại. Alice refresh, thấy tỉ số chung cuộc và hét lên. Bob refresh sau khi nghe Alice, nhưng request của Bob tới một replica bị lag, nên trang hiện trận đấu vẫn đang diễn ra. Nếu hai người bấm cùng lúc thì kết quả khác nhau cũng không lạ; vi phạm nằm ở chỗ Bob biết mình bấm sau Alice, nên kỳ vọng kết quả ít nhất phải mới bằng của Alice.
What Makes a System Linearizable?
Mô hình hoá: mỗi key là một register (key trong KV store, row, document…). Thao tác gồm read(x) ⇒ v, write(x, v) ⇒ r, và cas(x, v_old, v_new) ⇒ r (compare-and-set).
Mỗi request là một khoảng thời gian từ lúc gửi tới lúc nhận response; client chỉ biết DB xử lý request vào một thời điểm nào đó trong khoảng đó.
Các quy tắc: - Read kết thúc trước khi write bắt đầu → chắc chắn trả giá trị cũ. - Read bắt đầu sau khi write kết thúc → chắc chắn trả giá trị mới. - Read concurrent (chồng lấn) với write → có thể trả cũ hoặc mới. Register chỉ có tính chất này gọi là regular register — chưa đủ. - Ràng buộc thêm để thành linearizable: phải tồn tại một thời điểm trong khoảng write mà giá trị "lật" atomically từ cũ sang mới. Khi đã có một read trả giá trị mới thì mọi read bắt đầu sau đó (dù ở client khác) cũng phải thấy giá trị mới, kể cả khi write chưa trả response.
Cách hình dung: đánh dấu mỗi thao tác bằng một điểm (thời điểm nó "có hiệu lực") nằm trong thanh request của nó; nối các điểm lại theo thứ tự. Yêu cầu: đường nối luôn tiến về phía trước theo thời gian, và chuỗi đó phải là một chuỗi đọc/ghi hợp lệ trên một register.
Một số chi tiết đáng chú ý trong ví dụ phức tạp của sách:
- Thứ tự xử lý có thể khác thứ tự gửi request, miễn là các request đó concurrent.
- Một read có thể trả giá trị mới trước khi người ghi nhận được ok — chỉ là response bị delay trên mạng.
- Linearizability không có transaction isolation: giá trị có thể bị client khác đổi bất cứ lúc nào; dùng cas để phát hiện thay đổi concurrent.
- Nếu client A đã đọc được giá trị mới, thì client B bắt đầu read sau A không được đọc giá trị cũ hơn (lại là tình huống Alice/Bob).
Có thể kiểm tra một hệ có linearizable hay không bằng cách ghi lại timing của mọi request/response và kiểm tra xem có sắp xếp được thành một thứ tự tuần tự hợp lệ không (tốn kém về tính toán nhưng làm được).
Linearizability Versus Serializability
| Tiêu chí | Serializability | Linearizability | Causal consistency |
|---|---|---|---|
| Thuộc về | Isolation của transaction | Recency guarantee trên một object (register) | Consistency model cho toàn hệ thống |
| Phạm vi | Nhiều object mỗi transaction | Từng read/write đơn lẻ | Mọi thao tác có quan hệ nhân quả |
| Guarantee | Kết quả như chạy tuần tự theo một thứ tự nào đó (có thể khác thứ tự thực tế) | Như chỉ có một bản sao; read luôn thấy write mới nhất đã hoàn thành | Cause luôn được thấy trước effect; thao tác concurrent có thể thấy theo thứ tự bất kỳ |
| Loại thứ tự | Một thứ tự serial (không cần khớp thời gian thực) | Total order khớp thời gian thực | Partial order |
| Chặn write skew? | Có | Không (trừ khi materialize conflicts) | Không |
| Chi phí / CAP | Tuỳ implementation | Chậm, bị ảnh hưởng bởi CAP | Không bị CAP ảnh hưởng, nhanh hơn nhiều |
- Có cả hai = strict serializability (strong one-copy serializability, strong-1SR).
- 2PL và actual serial execution thường là linearizable.
- SSI không linearizable: đọc từ consistent snapshot, nên cố ý không thấy các write mới hơn snapshot.
Relying on Linearizability
Xem tỉ số bóng đá chậm vài giây thì chẳng hại gì. Nhưng có vài chỗ linearizability là bắt buộc:
1. Locking và leader election. Chọn leader bằng lock: node nào lấy được lock thì làm leader. Lock này phải linearizable, nếu không các node sẽ không đồng ý ai đang giữ lock. ZooKeeper và etcd dùng consensus để cung cấp thao tác linearizable có fault tolerance; Apache Curator cung cấp các "recipe" cấp cao hơn. Vẫn phải lo fencing (Chương 8). Oracle RAC thậm chí dùng lock linearizable cho từng disk page, nên cần mạng interconnect riêng.
- Lưu ý: ZooKeeper/etcd mặc định chỉ có write linearizable; read có thể stale. Muốn read linearizable: etcd dùng quorum read, ZooKeeper gọi sync() trước khi đọc.
2. Constraints và uniqueness guarantees. Username/email duy nhất, không có hai file cùng path. Muốn enforce ngay lúc ghi (hai người đăng ký cùng tên concurrent thì một người phải nhận lỗi) thì cần linearizability — về bản chất giống lấy "lock" trên username, hay một atomic CAS. Tương tự: số dư không âm, không bán quá tồn kho, không đặt trùng ghế. Nếu chấp nhận ràng buộc lỏng (overbook rồi bồi thường) thì không cần. Foreign key và attribute constraint thì không cần linearizability.
3. Cross-channel timing dependencies. Với Alice/Bob, vi phạm chỉ bị phát hiện vì có kênh liên lạc thứ hai (giọng Alice → tai Bob). Ví dụ trong hệ thống — image resizer: - Web server nhận ảnh upload, ghi ảnh full-size vào file storage, rồi đẩy message "resize ảnh X" vào message queue (không đẩy cả ảnh vì queue dành cho message nhỏ). - Resizer đọc message, lấy ảnh từ storage, resize, ghi thumbnail. - Nếu storage không linearizable, message queue có thể nhanh hơn replication nội bộ của storage → resizer đọc thấy ảnh cũ hoặc không thấy ảnh → ảnh full-size và thumbnail lệch nhau vĩnh viễn. - Nguyên nhân: có hai kênh giao tiếp (storage và queue). Linearizability là cách đơn giản nhất để tránh race này; nếu mình kiểm soát kênh phụ thì có thể dùng cách khác kiểu read-your-writes, đổi lại phức tạp hơn.
Implementing Linearizable Systems
Cách đơn giản nhất: chỉ dùng một bản sao — nhưng không chịu lỗi được. Vậy phải xét các phương pháp replication:
| Replication method | Linearizable? | Lý do |
|---|---|---|
| Single-leader | Có thể | Đọc từ leader hoặc follower cập nhật synchronous. Nhưng có thể hỏng nếu dùng snapshot isolation, có bug concurrency, có "delusional leader" (tưởng mình còn là leader), hoặc async failover làm mất write đã commit. |
| Consensus algorithms (ZooKeeper, etcd) | Có | Giống single-leader nhưng có cơ chế chống split brain và stale replica. |
| Multi-leader | Không | Ghi concurrent trên nhiều node, replicate async → sinh conflict. |
| Leaderless (Dynamo-style) | Có lẽ không | LWW dựa trên clock (Cassandra) gần như chắc chắn không; sloppy quorum phá hỏng hoàn toàn; ngay cả strict quorum vẫn có race. |
Linearizability và quorum: Ví dụ n=3, w=3, r=2. Writer đang ghi x=1 lên 3 replica với độ trễ khác nhau. Client A đọc 2 node, thấy 1 trên một node → trả 1. Client B bắt đầu sau khi A xong, đọc 2 node khác, cả hai vẫn 0 → trả 0. Thoả w + r > n mà vẫn không linearizable (lại là Alice/Bob).
- Có thể "chữa" bằng cách: reader làm read repair synchronous trước khi trả kết quả, và writer đọc trạng thái mới nhất của một quorum trước khi ghi. Riak không làm vì tốn performance; Cassandra có đợi read repair nhưng vẫn mất linearizability vì LWW khi có write concurrent.
- Kể cả vậy cũng chỉ được read/write linearizable; CAS linearizable thì cần consensus.
- Kết luận an toàn: coi leaderless là không linearizable.
The Cost of Linearizability
Ví dụ hai datacenter bị đứt mạng giữa chúng (trong mỗi DC vẫn ổn, client vẫn tới được DC): - Multi-leader: mỗi DC vẫn hoạt động bình thường; write được xếp hàng và trao đổi khi mạng phục hồi. - Single-leader: leader nằm ở một DC. Client chỉ tới được DC follower thì không ghi được và không đọc linearizable được (chỉ đọc stale từ follower) → ứng dụng bị unavailable ở đó cho tới khi mạng được sửa.
The CAP theorem
- Nếu ứng dụng cần linearizability và một số replica bị ngắt khỏi phần còn lại → những replica đó phải chờ hoặc trả lỗi → unavailable.
- Nếu không cần linearizability → mỗi replica xử lý độc lập → vẫn available nhưng không linearizable.
- Do Eric Brewer đặt tên năm 2000, nhưng trade-off đã được biết từ những năm 1970. Công lao lớn nhất của CAP là thúc đẩy văn hoá khám phá thiết kế shared-nothing (làn sóng NoSQL).
The Unhelpful CAP Theorem
- Cách nói "chọn 2 trong 3" (C, A, P) là sai lệch: network partition là một loại fault, không phải lựa chọn — nó sẽ xảy ra dù bạn muốn hay không.
- Cách nói đúng hơn: "either Consistent or Available when Partitioned". Khi mạng ổn thì có thể có cả hai.
- Định nghĩa "availability" trong CAP rất đặc thù, không khớp với nghĩa thông thường; nhiều hệ "highly available" không thoả định nghĩa này.
- Phạm vi hẹp: chỉ xét một consistency model (linearizability) và một loại fault (network partition); không nói gì về network delay, node chết… → giá trị thực tiễn thấp, chủ yếu mang ý nghĩa lịch sử. Tránh gắn nhãn "CP/AP".
Ordering Guarantees
Ordering xuất hiện khắp cuốn sách: leader quyết định thứ tự write trong replication log (Ch5); serializability là "như chạy theo thứ tự tuần tự nào đó" (Ch7); timestamp/clock để xác định write nào sau (Ch8). Chương này cho thấy có liên hệ sâu giữa ordering, linearizability và consensus.
Ordering and Causality
Ordering quan trọng vì nó giữ causality (nhân quả). Các ví dụ từ các chương trước: - Consistent prefix reads: thấy câu trả lời trước câu hỏi → vi phạm nhân quả. - Multi-leader: update tới trước insert của cùng row. - Happens-before trong phát hiện concurrent write: A trước B, B trước A, hoặc concurrent. - Snapshot isolation: "consistent" nghĩa là consistent with causality — snapshot có câu trả lời thì phải có câu hỏi. Read skew = đọc dữ liệu vi phạm causality. - Write skew (bác sĩ on-call): hành động nghỉ trực phụ thuộc nhân quả vào việc quan sát xem ai đang trực; SSI theo dõi các dependency này. - Alice/Bob và image resizer cũng là vi phạm causality.
Hệ tuân thủ thứ tự nhân quả gọi là causally consistent (ví dụ snapshot isolation).
Causal order không phải total order
- Total order: mọi cặp phần tử đều so sánh được (như các số tự nhiên).
- Partial order: một số cặp không so sánh được (như tập hợp {a,b} và {b,c}).
- Linearizability → total order: một timeline duy nhất, không có concurrency thực sự.
- Causality → partial order: hai thao tác concurrent là không so sánh được. Giống lịch sử Git: có branch rồi merge.
Linearizability mạnh hơn causal consistency
- Linearizability bao hàm causality: hệ linearizable tự động giữ causality, kể cả khi có nhiều kênh liên lạc (queue + storage) mà không cần truyền timestamp giữa các thành phần. Đó là lý do nó dễ hiểu.
- Nhưng causal consistency có thể đạt được không cần linearizability. Causal consistency là consistency model mạnh nhất không bị chậm đi vì network delay và vẫn available khi network fail — CAP không áp dụng cho nó.
- Nhiều hệ tưởng cần linearizability nhưng thực ra chỉ cần causal consistency. Đây là hướng nghiên cứu đầy hứa hẹn nhưng chưa phổ biến trong production.
Capturing causal dependencies
- Replica muốn xử lý một thao tác thì mọi thao tác xảy ra trước nó (theo nhân quả) phải được xử lý trước; thiếu thì phải chờ.
- Cần mô tả "knowledge" của node: khi phát ra write Y, node đã thấy X chưa? (Giống câu hỏi trong điều tra: CEO có biết X lúc ra quyết định Y không?)
- Kỹ thuật: tổng quát hoá version vectors cho toàn DB (chứ không chỉ một key); DB cần biết ứng dụng đã đọc version nào (truyền version number khi ghi, như SSI kiểm tra lúc commit).
Sequence Number Ordering
Theo dõi toàn bộ causal dependency là quá tốn (client đọc nhiều thứ trước khi ghi). Thay vào đó dùng sequence number / logical clock: nhỏ gọn và cho total order. Mục tiêu: total order consistent with causality (A xảy ra trước B theo nhân quả ⇒ A có số nhỏ hơn). Total order kiểu so UUID ngẫu nhiên thì hợp lệ nhưng vô dụng.
Single-leader: leader tăng counter cho mỗi thao tác trong replication log → follower apply theo thứ tự đó luôn causally consistent (dù lag).
Noncausal sequence number generators
Khi không có single leader (multi-leader, leaderless, partitioned), các cách thường dùng: 1. Mỗi node tự sinh số riêng (node 1 số lẻ, node 2 số chẵn; hoặc nhúng node ID vào vài bit). 2. Gắn timestamp của physical clock (dùng trong LWW). 3. Cấp trước block số (node A 1–1000, node B 1001–2000).
Cả ba đều nhanh và scalable hơn đẩy mọi thứ qua một leader, nhưng không consistent with causality: node chạy nhanh/chậm khác nhau nên counter lệch nhau; clock skew; thao tác sau về nhân quả có thể rơi vào block số nhỏ hơn. (Spanner chờ hết khoảng bất định của clock để timestamp khớp causality, nhưng đa số clock không cung cấp được con số bất định đó.)
Lamport timestamps
- Timestamp = cặp (counter, node ID). So counter trước, bằng nhau thì so node ID → total order, duy nhất.
- Điểm mấu chốt: mọi node và client đều theo dõi counter lớn nhất đã thấy và gửi kèm trong mọi request. Node nhận được giá trị max lớn hơn counter của mình thì nhảy counter lên bằng max đó ngay.
- Ví dụ trong sách: client A nhận counter 5 từ node 2, rồi gửi max=5 tới node 1; node 1 đang ở 1 nên nhảy lên 5, thao tác tiếp theo mang 6.
- Nhờ vậy mọi causal dependency đều làm timestamp tăng → consistent with causality.
- So với version vectors: version vector phân biệt được concurrent hay có quan hệ nhân quả; Lamport timestamp luôn ép thành total order nên không biết được hai thao tác có concurrent không. Ưu điểm của Lamport: gọn hơn.
Timestamp ordering is not sufficient
Bài toán username duy nhất: ý tưởng "timestamp nhỏ hơn thắng" chỉ đúng sau khi đã thu thập hết các thao tác. Còn khi một node vừa nhận request và phải quyết định ngay bây giờ, nó không biết node khác có đang concurrent tạo cùng username với timestamp nhỏ hơn hay không. Muốn chắc chắn phải hỏi mọi node → một node chết/mất mạng là cả hệ đứng → không fault-tolerant.
Vấn đề cốt lõi: total order chỉ hình thành sau khi gom đủ; cần biết khi nào thứ tự đã được chốt (finalized). Đó là total order broadcast.
Total Order Broadcast
Single-leader tạo total order bằng cách tuần tự hoá mọi thứ trên CPU của leader. Thách thức: scale khi vượt khả năng một leader, và failover khi leader chết. Trong lý thuyết, bài toán này gọi là total order broadcast (hay atomic broadcast — tên gây nhầm, chẳng liên quan tới atomicity của ACID).
Lưu ý về phạm vi: DB partitioned với leader riêng cho mỗi partition thường chỉ giữ thứ tự trong từng partition; total order xuyên partition cần phối hợp thêm.
Hai safety property bắt buộc (kể cả khi node/network lỗi): - Reliable delivery: không mất message — đã tới một node thì tới mọi node. - Totally ordered delivery: mọi node nhận message theo cùng một thứ tự.
Khi mạng đứt thì message không tới được, nhưng thuật toán retry, và khi mạng phục hồi message vẫn phải tới đúng thứ tự.
Using total order broadcast
- ZooKeeper và etcd thực chất cài total order broadcast → gợi ý có liên hệ với consensus.
- State machine replication: mỗi message là một write; mọi replica xử lý cùng các write theo cùng thứ tự → các replica nhất quán.
- Serializable transaction: mỗi message là một stored procedure deterministic, mọi node chạy theo cùng thứ tự.
- Thứ tự chốt tại thời điểm deliver: không được chèn ngược một message vào vị trí trước các message đã deliver → mạnh hơn timestamp ordering.
- Có thể xem nó như một log (append-only): deliver = append.
- Dùng để làm lock service có fencing token: sequence number trong log tăng đơn điệu (ZooKeeper gọi là
zxid).
Implementing linearizable storage using total order broadcast
Total order broadcast là asynchronous (đảm bảo thứ tự, không đảm bảo khi nào tới); linearizability là recency guarantee. Khác nhau, nhưng có thể xây cái này từ cái kia.
Ví dụ username duy nhất bằng linearizable CAS xây trên log: 1. Append message vào log, "tạm" đòi username. 2. Đọc log, chờ message của mình được deliver lại cho mình. 3. Nếu message đầu tiên đòi username đó là của mình → thành công (có thể append thêm message commit). Nếu là của người khác → abort.
Vì mọi node thấy log theo cùng thứ tự, mọi node đồng ý ai thắng. Cách này cho write linearizable nhưng read có thể stale — chính xác ra là sequential consistency (timeline consistency). Để read linearizable:
- Đưa read qua log (append một message, đợi nó deliver rồi mới đọc) — giống quorum read của etcd.
- Lấy vị trí mới nhất của log một cách linearizable, đợi tới vị trí đó rồi đọc — ý tưởng của sync() trong ZooKeeper.
- Đọc từ replica được cập nhật synchronous (chain replication).
Implementing total order broadcast using linearizable storage
Ngược lại: có một register linearizable chứa số nguyên với increment-and-get (hoặc CAS). Mỗi message: increment-and-get để lấy sequence number, gắn vào message, gửi tới mọi node (resend nếu mất); bên nhận deliver theo thứ tự số.
Khác biệt then chốt với Lamport timestamps: số từ register linearizable không có lỗ hổng (no gaps). Đã deliver message 4 mà nhận được 6 thì biết phải chờ 5.
Làm counter này chịu lỗi được (node giữ nó chết, mất mạng) thì rốt cuộc sẽ ra một thuật toán consensus. Không phải trùng hợp:
Distributed Transactions and Consensus
Consensus: làm cho nhiều node đồng ý về một điều. Nghe đơn giản nhưng rất nhiều hệ thống hỏng vì tưởng nó dễ. Hai tình huống điển hình: - Leader election: mọi node phải đồng ý ai là leader, tránh split brain. - Atomic commit: transaction trải trên nhiều node/partition; hoặc tất cả commit, hoặc tất cả abort.
FLP impossibility (Fischer, Lynch, Paterson): không có thuật toán nào luôn đạt được consensus nếu có nguy cơ node crash — nhưng chỉ trong asynchronous model (deterministic, không dùng clock/timeout). Cho phép timeout (dù đôi khi nghi oan) hoặc dùng số ngẫu nhiên là đủ để giải → trong thực tế consensus làm được.
Atomic Commit and Two-Phase Commit (2PC)
From single-node to distributed atomic commit
- Single node: ghi data vào WAL, rồi ghi commit record. Thời điểm disk ghi xong commit record là điểm quyết định. Atomicity nhờ một thiết bị duy nhất (disk controller của một node).
- Nhiều node: chỉ gửi commit cho từng node rồi commit độc lập là không đủ: node thì phát hiện vi phạm constraint phải abort, request commit tới node khác bị mất, node khác crash trước khi ghi commit record…
- Commit phải irrevocable: đã commit thì dữ liệu hiển thị cho transaction khác (read committed), không thể rút lại. (Có thể "undo" bằng compensating transaction, nhưng về mặt DB đó là transaction riêng.) → Một node chỉ được commit khi chắc chắn mọi node khác cũng sẽ commit.
Introduction to two-phase commit
2PC dùng thêm một thành phần: coordinator (transaction manager; thường là library trong process ứng dụng, ví dụ trong Java EE container; ví dụ Narayana, JOTM, BTM, MSDTC). Các DB node tham gia gọi là participants.
- Phase 1 — prepare: coordinator hỏi mọi participant "có commit được không?".
- Phase 2 — commit/abort: tất cả "yes" → gửi commit; có bất kỳ "no" → gửi abort cho tất cả.
Ví von đám cưới: mục sư hỏi riêng cô dâu và chú rể "có đồng ý không"; nhận được "I do" từ cả hai thì tuyên bố thành vợ chồng — transaction đã commit và thông báo cho mọi người. Chỉ cần một người không nói "yes" là buổi lễ huỷ.
Đừng nhầm 2PC với 2PL: 2PC = atomic commit phân tán; 2PL = serializable isolation. Hoàn toàn khác nhau.
A system of promises
Chi tiết từng bước: 1. Ứng dụng xin coordinator một global transaction ID. 2. Ứng dụng mở single-node transaction trên từng participant, gắn ID đó. Lỗi ở giai đoạn này → bất kỳ ai cũng có thể abort. 3. Sẵn sàng commit → coordinator gửi prepare (kèm ID). Có request lỗi/timeout → gửi abort cho tất cả. 4. Participant nhận prepare phải đảm bảo chắc chắn commit được trong mọi tình huống: ghi hết dữ liệu xuống disk, kiểm tra conflict/constraint. Trả "yes" nghĩa là từ bỏ quyền abort, nhưng chưa commit. 5. Coordinator nhận đủ phản hồi, quyết định, và ghi quyết định xuống transaction log trên disk — đây là commit point. 6. Gửi commit/abort; lỗi hay timeout thì retry mãi mãi. Participant đã crash thì commit khi nó hồi phục (vì nó đã hứa "yes").
→ Hai point of no return: participant vote "yes"; coordinator ra quyết định. (Single-node commit gộp hai điểm này thành một: ghi commit record.)
Quay lại đám cưới: sau khi nói "I do" thì không rút lại được. Nếu bạn ngất trước khi nghe tuyên bố, bạn vẫn đã kết hôn; tỉnh dậy thì hỏi mục sư trạng thái của global transaction ID, hoặc chờ mục sư retry.
Coordinator failure
- Coordinator chết trước khi gửi prepare → participant abort an toàn.
- Participant đã vote "yes" rồi mà coordinator chết/mất mạng → participant không thể tự quyết: tự abort có thể lệch với node đã commit; tự commit có thể lệch với node đã abort. Trạng thái này gọi là in doubt / uncertain. Timeout cũng không giúp.
- Cách duy nhất: chờ coordinator hồi phục, đọc log; transaction nào không có commit record thì abort. Commit point của 2PC rốt cuộc là một single-node atomic commit trên coordinator.
Three-phase commit
2PC là blocking atomic commit. 3PC là phiên bản nonblocking nhưng giả định network delay và response time có giới hạn — không đúng với hệ thực tế (Ch8). Nonblocking atomic commit cần perfect failure detector, mà timeout thì không phải. Vì vậy 2PC vẫn được dùng.
Distributed Transactions in Practice
Danh tiếng lẫn lộn: có safety guarantee quan trọng, nhưng gây rắc rối vận hành, giết performance (distributed transaction trong MySQL được báo chậm hơn 10 lần, do thêm fsync và round-trip mạng). Hai loại hay bị lẫn:
- Database-internal: mọi participant chạy cùng phần mềm DB (VoltDB, MySQL Cluster NDB) → tối ưu riêng được, thường chạy khá tốt.
- Heterogeneous: participant là các công nghệ khác nhau (DB của hai vendor, message broker) → khó hơn nhiều.
Exactly-once message processing
Với heterogeneous transaction: ack message trong queue và các write vào DB được commit atomically trong cùng transaction. Lỗi thì cả hai abort, broker gửi lại message an toàn → message effectively được xử lý đúng một lần. Điều kiện: mọi hệ bị ảnh hưởng phải hỗ trợ cùng một atomic commit protocol — ví dụ gửi email (email server không hỗ trợ 2PC) thì retry có thể gửi email nhiều lần.
XA transactions
- X/Open XA (1991): chuẩn 2PC cho các công nghệ khác nhau; được hỗ trợ bởi PostgreSQL, MySQL, DB2, SQL Server, Oracle, ActiveMQ, HornetQ, MSMQ, IBM MQ…
- XA không phải network protocol, chỉ là C API để giao tiếp với coordinator. Trong Java: JTA, qua driver JDBC/JMS.
- Coordinator thường là library trong process ứng dụng, lưu log trên disk local. Process/máy chết → participant kẹt in doubt cho tới khi server khởi động lại và đọc log. DB server không liên lạc trực tiếp được với coordinator.
Holding locks while in doubt
Vấn đề thật sự là lock: transaction giữ row-level exclusive lock (và shared lock nếu dùng 2PL) cho tới khi commit/abort. In doubt 20 phút = giữ lock 20 phút; mất log coordinator = giữ mãi mãi cho tới khi admin xử lý tay. Các transaction khác đụng vào các row đó bị block → nhiều phần ứng dụng thành unavailable.
Recovering from coordinator failure
- Orphaned in-doubt transactions có xảy ra (log mất/hỏng do bug). Reboot DB cũng không giải quyết được vì 2PC đúng phải giữ lock qua cả restart.
- Lối thoát: admin tự quyết commit/rollback bằng tay, thường là lúc sự cố nghiêm trọng, áp lực cao.
- Heuristic decisions: participant đơn phương quyết định — thực chất là cách nói giảm của "có lẽ phá vỡ atomicity". Chỉ dành cho thảm hoạ.
Limitations of distributed transactions
- Coordinator chính là một database (lưu kết quả transaction) → không replicate thì là single point of failure; nhiều implementation mặc định không HA.
- Coordinator nằm trong application server → server không còn stateless; log của nó quan trọng ngang DB.
- XA là mẫu số chung nhỏ nhất: không phát hiện được deadlock xuyên hệ thống, không chạy được với SSI.
- Kể cả database-internal: 2PC cần tất cả participant trả lời → bất kỳ phần nào hỏng là transaction fail → khuếch đại lỗi (amplifying failures), ngược với mục tiêu fault tolerance.
Fault-Tolerant Consensus
Định nghĩa hình thức: một hoặc nhiều node propose giá trị, thuật toán decide một giá trị (ví dụ nhiều khách tranh ghế cuối cùng; mỗi node propose ID khách của mình). Bốn tính chất:
- Uniform agreement: không có hai node quyết định khác nhau.
- Integrity: không node nào quyết định hai lần.
- Validity: giá trị được quyết định phải do một node nào đó propose (loại trừ lời giải tầm thường như luôn quyết định
null). - Termination: mọi node không crash cuối cùng đều quyết định được một giá trị.
Ba cái đầu là safety, termination là liveness — chính nó hình thức hoá fault tolerance. Nếu không cần fault tolerance thì cứ chỉ định một node "độc tài" là xong — đó chính là 2PC, và 2PC không thoả termination (phải chờ coordinator). Mô hình giả định node crash là biến mất vĩnh viễn (datacenter bị lở đất vùi dưới 30 feet bùn), nên thuật toán nào phải chờ node hồi phục thì không đạt termination.
- Termination cần đa số node còn sống (quorum). Tuy vậy, hầu hết implementation luôn giữ safety kể cả khi mất đa số — outage lớn có thể làm hệ ngừng xử lý, nhưng không làm nó ra quyết định sai.
- Đa số thuật toán giả định không có Byzantine fault. Chống Byzantine được nếu dưới 1/3 node lỗi Byzantine.
Consensus algorithms and total order broadcast
Các thuật toán nổi tiếng: Viewstamped Replication (VSR), Paxos, Raft, Zab. Không nên tự cài (rất khó).
Thực tế chúng quyết định một chuỗi giá trị → chính là total order broadcast = nhiều vòng consensus lặp lại, mỗi vòng quyết định message kế tiếp: - Agreement → mọi node deliver cùng message, cùng thứ tự. - Integrity → không duplicate. - Validity → không bị hỏng hay bịa ra. - Termination → không bị mất.
VSR, Raft, Zab làm total order broadcast trực tiếp (hiệu quả hơn); với Paxos tối ưu này gọi là Multi-Paxos.
Single-leader replication and consensus
Single-leader về bản chất cũng là total order broadcast. Vì sao Ch5 không cần consensus? Vì cách chọn leader: - Người vận hành chọn tay → consensus kiểu "độc tài", không thoả termination (cần con người can thiệp). - Tự động failover → cần mọi node đồng ý ai là leader → cần consensus. Nhưng consensus lại cần leader… muốn bầu leader thì phải có leader trước? Vòng luẩn quẩn.
Epoch numbering and quorums
Lời giải: các giao thức không đảm bảo leader là duy nhất tuyệt đối, mà đảm bảo mỗi epoch có đúng một leader. Epoch number có tên khác nhau: ballot number (Paxos), view number (VSR), term number (Raft). - Mỗi lần nghi leader chết → bầu cử với epoch tăng dần. Hai leader xung đột thì epoch cao hơn thắng. - Trước khi quyết định bất cứ điều gì, leader phải kiểm tra không có leader nào với epoch cao hơn: gửi đề xuất và chờ quorum (thường là đa số) đồng ý. Node chỉ vote đồng ý nếu không biết leader nào có epoch cao hơn. - Hai vòng vote: một để bầu leader, một để duyệt đề xuất của leader. Hai quorum này phải giao nhau → nếu vote đề xuất thành công mà không thấy epoch nào cao hơn thì leader chắc chắn mình vẫn là leader.
| Tiêu chí | 2PC | Fault-tolerant consensus (Raft/Paxos/Zab) |
|---|---|---|
| Người điều phối | Coordinator cố định, không được bầu | Leader được bầu, có epoch/term |
| Số vote cần | "Yes" từ tất cả participant | Đa số (quorum) |
| Khi người điều phối chết | Block, participant in doubt tới khi coordinator hồi phục | Bầu leader mới, tiếp tục chạy |
| Termination | Không thoả | Thoả (nếu đa số còn sống) |
| Recovery | Dựa vào log coordinator; có thể cần admin/heuristic | Có quy trình để node về trạng thái nhất quán sau khi đổi leader |
| Mục đích chính | Atomic commit trên nhiều hệ (kể cả heterogeneous) | Total order broadcast, linearizable storage, leader election |
Limitations of consensus
- Vote trước khi quyết định = synchronous replication → chậm hơn; nhiều người chọn async và chấp nhận rủi ro mất dữ liệu.
- Cần đa số nghiêm ngặt: 3 node chịu được 1 lỗi, 5 node chịu được 2. Mạng bị chia thì chỉ phần đa số chạy được.
- Thường giả định tập node cố định; dynamic membership ít được hiểu rõ hơn.
- Dựa vào timeout để phát hiện lỗi → mạng có delay biến thiên (nhất là phân tán địa lý) dễ bầu leader liên tục, tốn thời gian bầu hơn làm việc.
- Nhạy với mạng: Raft có edge case khi chỉ một link chập chờn → leadership nhảy qua lại liên tục, không có tiến triển.
Membership and Coordination Services
ZooKeeper/etcd trông như KV store nhưng được thiết kế cho lượng dữ liệu nhỏ, vừa trong memory, replicate bằng total order broadcast chịu lỗi. Hiếm khi dùng trực tiếp; thường gián tiếp qua HBase, Hadoop YARN, OpenStack Nova, Kafka. ZooKeeper mô phỏng theo Google Chubby. Các tính năng:
- Linearizable atomic operations: CAS để làm lock; lock thường là lease có thời hạn hết hạn.
- Total ordering of operations:
zxid,cversiontăng đơn điệu → dùng làm fencing token. - Failure detection: session dài hạn với heartbeat; hết session timeout thì session chết, ephemeral nodes (ví dụ lock) tự động bị xoá.
- Change notifications: client watch thay đổi (node mới tham gia, node chết) thay vì poll.
Chỉ có atomic ops thực sự cần consensus, nhưng chính sự kết hợp các tính năng mới làm ZooKeeper hữu ích.
Allocating work to nodes: chọn leader/primary cho process, job scheduler; gán partition cho node và rebalance khi node thêm vào hay chết đi — dùng atomic ops + ephemeral nodes + notifications (Curator hỗ trợ). ZooKeeper chạy trên số node cố định (3 hoặc 5) nhưng phục vụ hàng nghìn client → "outsource" consensus. Dữ liệu trong đó thay đổi chậm (phút/giờ, ví dụ "node 10.1.1.23 là leader của partition 7"), không dành cho runtime state thay đổi hàng triệu lần/giây (dùng BookKeeper cho việc đó).
Service discovery: ZooKeeper/etcd/Consul dùng để tìm IP của service. Nhưng service discovery không nhất thiết cần consensus — DNS không linearizable mà vẫn ổn; điều quan trọng là availability. Một số hệ có read-only caching replicas nhận log async, không vote, phục vụ read không cần linearizable.
Membership services: xác định node nào đang sống trong cluster. Không thể phát hiện lỗi chắc chắn, nhưng kết hợp failure detection với consensus thì các node đồng ý được về membership (có thể khai tử nhầm một node còn sống, nhưng có sự đồng thuận vẫn rất hữu ích, ví dụ chọn leader = node có ID nhỏ nhất trong membership hiện tại).
Summary — các bài toán tương đương consensus
| Bài toán | Cái cần "decide" |
|---|---|
| Linearizable CAS register | Có set giá trị hay không, dựa trên giá trị hiện tại |
| Atomic transaction commit | Commit hay abort distributed transaction |
| Total order broadcast | Thứ tự deliver message |
| Locks và leases | Client nào lấy được lock |
| Membership/coordination service | Node nào còn sống, node nào coi như chết |
| Uniqueness constraint | Transaction nào được tạo record với key đó |
Với một node duy nhất hoặc một leader duy nhất thì tất cả đều dễ. Khi leader chết có ba lựa chọn: 1. Chờ leader hồi phục (nhiều XA/JTA coordinator) — không thoả termination. 2. Failover thủ công — consensus kiểu "act of God", chậm bằng tốc độ con người. 3. Tự động bầu leader bằng thuật toán consensus đã được chứng minh.
Có leader chỉ là "đá cái lon đi xa hơn": vẫn cần consensus, chỉ là ở chỗ khác và ít thường xuyên hơn. Không phải hệ nào cũng cần consensus — leaderless và multi-leader chấp nhận conflict và lịch sử phân nhánh/hợp nhất.
⚠️ Hiểu lầm & cạm bẫy thường gặp
- Nhầm linearizability với serializability: một cái là recency trên một object, một cái là isolation của transaction nhiều object. SSI serializable nhưng không linearizable.
- Nghĩ quorum
w + r > nlà linearizable: sai khi có network delay biến thiên; cộng thêm LWW hay sloppy quorum thì càng không. - "CAP: chọn 2 trong 3": partition không phải thứ để chọn. Nhãn CP/AP gây hiểu lầm; CAP chỉ nói về linearizability + partition.
- Nghĩ bỏ linearizability là vì fault tolerance: phần lớn là vì performance/latency.
- Tin ZooKeeper/etcd read mặc định là linearizable: chỉ write linearizable; read cần
sync()hoặc quorum read. - Dùng Lamport timestamp để enforce uniqueness ngay lập tức: không được, vì không biết thứ tự đã chốt chưa → cần total order broadcast.
- Nhầm Lamport timestamp với version vector: Lamport không phát hiện được concurrency.
- Nhầm 2PC với 2PL: hai thứ hoàn toàn khác nhau.
- Tưởng 2PC chịu lỗi được: coordinator chết là participant in doubt, giữ lock vô thời hạn; 2PC không thoả termination.
- Dùng heuristic decision tuỳ tiện: thực chất là phá vỡ atomicity.
- Nghĩ FLP nghĩa là consensus không làm được trong thực tế: chỉ đúng trong asynchronous model thuần tuý, không có timeout/random.
- Tự viết thuật toán consensus: lịch sử thất bại rất dày; hãy dùng ZooKeeper/etcd.
- Dùng ZooKeeper làm database đa năng hoặc lưu state thay đổi nhanh: sai mục đích.
- Nghĩ thêm node vào cluster consensus luôn tốt hơn: số node chẵn không tăng khả năng chịu lỗi (4 node vẫn chỉ chịu được 1 lỗi), và nhiều node hơn thì vote chậm hơn.
💼 Áp dụng thực tế & phỏng vấn
- Leader election / distributed lock: dùng etcd (Kubernetes lưu toàn bộ state của cluster trong etcd), ZooKeeper (Kafka đời cũ, HBase), Consul. Nhớ nói tới lease + fencing token khi thiết kế lock.
- Unique username / đặt ghế / chống bán quá tồn kho: cần linearizable CAS (single-leader DB với unique constraint, hoặc consensus store). Nếu nghiệp vụ cho phép thì dùng ràng buộc lỏng + bồi thường để tăng availability.
- Multi-region: linearizable toàn cầu thì mỗi write tốn cross-region RTT (Spanner dùng TrueTime + Paxos). Nếu chỉ cần causal hay read-your-writes thì chọn model yếu hơn cho latency thấp.
- Cross-channel race (upload file rồi đẩy event vào queue): nhắc tới đọc từ leader/primary, hoặc gắn version/ETag vào message rồi consumer chờ tới khi thấy version đó.
- Distributed transaction giữa các microservice: thường tránh 2PC/XA, dùng Saga (compensating transaction), transactional outbox + CDC, idempotent consumer để đạt "effectively once" (liên hệ Ch11–12).
- Kích thước cluster consensus: 3 hoặc 5 node (số lẻ); 5 node chịu được 2 lỗi. Giải thích được vì sao một cluster bị chia 2–3 thì chỉ phía 3 node chạy tiếp.
- Câu hỏi hay gặp: "Raft hoạt động thế nào?" → term, leader election bằng đa số, log replication, commit khi đa số đã ghi, quorum giao nhau. "Vì sao cần fencing token?" → leader/lock holder bị GC pause, lease đã hết hạn mà vẫn tưởng mình còn giữ.
- Nói về CAP khi phỏng vấn: nói "khi có partition, chọn consistency (linearizability) hay availability", rồi bổ sung PACELC (khi không có partition thì đánh đổi latency và consistency) để thể hiện hiểu sâu.