Chương 8 — The Trouble with Distributed Systems
🎯 Mục tiêu chương: Nhìn thẳng vào mọi thứ có thể hỏng trong hệ phân tán — mạng không tin cậy, đồng hồ không tin cậy, process bị pause bất kỳ lúc nào — và học cách suy luận về "sự thật" khi một node không thể biết chắc điều gì (quorum, fencing token, system model, safety vs liveness).
Mở đầu
Các chương trước vẫn còn quá lạc quan. Chương này giả định cái gì có thể hỏng thì sẽ hỏng (ngoại lệ: tạm giả định fault là non-Byzantine). Khác biệt cốt lõi giữa lập trình trên một máy và hệ phân tán: có rất nhiều cách mới để mọi thứ sai. Nhiệm vụ của kỹ sư là xây hệ thống vẫn đáp ứng đảm bảo cho người dùng dù mọi thứ bên dưới hỏng. Chương 8 chỉ mô tả vấn đề; Chương 9 mới bàn về thuật toán giải quyết.
Faults and Partial Failures
- Trên một máy, phần mềm tốt thường deterministic: hoặc chạy đúng, hoặc hỏng hẳn (kernel panic, blue screen). Đây là lựa chọn thiết kế có chủ đích: thà crash còn hơn trả kết quả sai; máy tính che giấu thực tại vật lý lộn xộn và trình ra một mô hình lý tưởng.
- Trong hệ phân tán, ta phải đối mặt thực tế vật lý. Coda Hale kể: network partition kéo dài trong một DC, hỏng PDU, hỏng switch, vô tình tắt nguồn cả rack, hỏng backbone cả DC, mất điện cả DC, thậm chí một tài xế hạ đường huyết đâm xe bán tải vào hệ thống HVAC của DC.
- Partial failure: một phần hệ thống hỏng theo cách khó lường trong khi phần khác vẫn chạy. Điều khó là partial failure nondeterministic — thao tác liên quan nhiều node có lúc chạy, có lúc hỏng, và bạn thậm chí không biết nó thành công hay không.
Cloud Computing and Supercomputing
Phổ triết lý xây hệ lớn:
- HPC / supercomputer: hàng nghìn CPU cho tính toán khoa học (dự báo thời tiết, mô phỏng phân tử). Job checkpoint định kỳ; một node hỏng thì dừng cả cluster, sửa xong chạy lại từ checkpoint → biến partial failure thành total failure, giống một máy đơn.
- Cloud computing: datacenter multi-tenant, máy commodity, mạng IP/Ethernet, cấp phát tài nguyên co giãn, tính tiền theo dùng.
- Datacenter doanh nghiệp truyền thống nằm giữa.
Hệ internet service khác supercomputer ở chỗ:
- Online: phải phục vụ low latency mọi lúc, không được dừng cả cluster để sửa (batch job thì dừng được).
- Supercomputer dùng phần cứng chuyên dụng, tin cậy, giao tiếp qua shared memory/RDMA; cloud dùng máy commodity rẻ hơn nhưng tỷ lệ hỏng cao hơn.
- Mạng datacenter dùng IP/Ethernet theo topology Clos; supercomputer dùng mesh/torus chuyên biệt.
- Hệ càng lớn, càng chắc chắn luôn có thứ gì đó đang hỏng; nếu chiến lược là bỏ cuộc thì hệ lớn sẽ dành phần lớn thời gian để phục hồi.
- Chịu được node hỏng giúp vận hành: rolling upgrade, giết VM chậm để xin VM mới.
- Triển khai phân tán địa lý đi qua internet — chậm và kém tin cậy hơn mạng nội bộ.
→ Phải chấp nhận partial failure và xây cơ chế fault tolerance trong phần mềm: xây hệ tin cậy từ các thành phần không tin cậy. Ngay cả hệ nhỏ vài node cũng cần; hãy chủ động tạo lỗi trong môi trường test. Trong hệ phân tán, nghi ngờ, bi quan và paranoia đều đáng giá.
Building a Reliable System from Unreliable Components
Trực giác nói hệ chỉ tin cậy bằng mắt xích yếu nhất — sai. Ví dụ:
- Error-correcting code truyền dữ liệu chính xác qua kênh đôi khi lật bit.
- TCP trên IP: IP có thể drop, delay, duplicate, reorder packet; TCP truyền lại packet mất, loại bỏ trùng, sắp xếp lại thứ tự.
Nhưng luôn có giới hạn: ECC không cứu được kênh ngập nhiễu; TCP che được mất/trùng/đảo packet nhưng không xóa được delay. Lớp cao hơn vẫn hữu ích vì xử lý các lỗi mức thấp khó nhằn, phần còn lại dễ suy luận hơn.
Unreliable Networks
Sách tập trung vào shared-nothing system: các máy chỉ giao tiếp qua mạng, mỗi máy có memory/disk riêng. Lý do phổ biến: rẻ, không cần phần cứng đặc biệt, dùng được cloud, đạt độ tin cậy qua redundancy nhiều DC.
Internet và mạng nội bộ datacenter là asynchronous packet network: gửi message đi nhưng không có đảm bảo khi nào tới, hay có tới không. Gửi request và chờ response, có thể:
- Request bị mất (ai đó rút cáp).
- Request nằm trong hàng đợi, tới sau (mạng hoặc bên nhận quá tải).
- Node đích đã chết (crash, tắt nguồn).
- Node đích tạm ngừng phản hồi (VD GC pause dài) nhưng sẽ trả lời lại sau.
- Node đích đã xử lý nhưng response bị mất (switch cấu hình sai).
- Node đích đã xử lý nhưng response bị trễ (mạng hoặc chính máy bạn quá tải).
Người gửi không thể phân biệt các trường hợp này — chỉ biết là chưa nhận được response. Cách xử lý thông thường là timeout, nhưng khi timeout bạn vẫn không biết request có tới nơi hay không (và nó có thể vẫn đang nằm trong queue và tới nơi sau).
Network Faults in Practice
- Nghiên cứu trong một datacenter cỡ vừa: khoảng 12 network fault mỗi tháng, một nửa làm mất kết nối một máy, một nửa mất kết nối cả rack.
- Nghiên cứu khác về top-of-rack switch, aggregation switch, load balancer: thêm thiết bị mạng dư thừa không giảm fault nhiều như mong đợi vì không chống được lỗi con người (cấu hình sai switch) — nguyên nhân chính gây outage.
- Public cloud như EC2 nổi tiếng có glitch mạng tạm thời thường xuyên.
- Nâng cấp phần mềm switch gây tái cấu hình topology → packet bị trễ hơn một phút.
- Cá mập cắn cáp ngầm dưới biển.
- Network interface có thể drop mọi packet chiều vào nhưng gửi chiều ra bình thường → link chạy một chiều không đảm bảo chiều kia.
Network partitions
Khi một phần mạng bị cắt khỏi phần còn lại: network partition / netsplit (sách dùng từ network fault để tránh nhầm với partition/shard của Ch6).
- Dù hiếm, việc nó có thể xảy ra nghĩa là phần mềm phải xử lý được. Nếu không định nghĩa và test, điều tồi tệ tùy ý có thể xảy ra: cluster deadlock vĩnh viễn kể cả khi mạng hồi phục, hoặc xóa sạch dữ liệu.
- Xử lý không nhất thiết là chịu đựng (tolerate): nếu mạng thường ổn, hiển thị lỗi cho user cũng hợp lệ — miễn là bạn biết phần mềm phản ứng thế nào và hệ phục hồi được. Chủ động gây lỗi mạng để test (tinh thần Chaos Monkey).
Detecting Faults
Nhiều hệ cần tự phát hiện node hỏng: load balancer loại node chết khỏi rotation; single-leader replication cần promote follower khi leader chết. Một số trường hợp có phản hồi rõ ràng:
- Máy còn sống nhưng không process nào lắng nghe port → OS gửi
RST/FIN. Nhưng nếu node crash giữa lúc xử lý request thì không biết đã xử lý bao nhiêu. - Process crash nhưng OS còn chạy → script báo cho node khác để takeover nhanh, không cần chờ timeout (HBase làm vậy).
- Có quyền truy cập management interface của switch → phát hiện link failure ở mức phần cứng (không khả dụng qua internet hay DC dùng chung).
- Router biết IP không tới được → trả ICMP Destination Unreachable (nhưng router cũng không có phép màu).
Không thể dựa vào phản hồi nhanh. Ngay cả TCP ack packet cũng không có nghĩa app đã xử lý (app có thể crash trước). Muốn chắc chắn thành công cần positive response từ chính application. Nói chung phải giả định không nhận được response gì: retry vài lần, chờ timeout, rồi tuyên bố node chết.
Timeouts and Unbounded Delays
Timeout dài bao nhiêu? Không có đáp án đơn giản:
- Timeout dài: chờ lâu mới phát hiện node chết; user phải chờ hoặc thấy lỗi.
- Timeout ngắn: phát hiện nhanh nhưng dễ tuyên bố nhầm node chết khi nó chỉ chậm tạm thời. Hậu quả: node thật ra còn sống và đang làm một hành động (VD gửi email), node khác takeover → hành động bị làm 2 lần. Chuyển trách nhiệm sang node khác cũng tăng tải; nếu node chỉ chậm vì overload, dồn tải sang node khác có thể gây cascading failure (cực đoan: mọi node tuyên bố nhau chết và cả hệ dừng).
Giả sử mạng đảm bảo delay tối đa d và node xử lý request trong tối đa r → mọi request thành công có response trong 2d + r, đó là timeout hợp lý. Nhưng thực tế không có đảm bảo nào: asynchronous network có unbounded delay, server cũng không đảm bảo thời gian xử lý tối đa. Với failure detection, "nhanh hầu hết thời gian" là chưa đủ: timeout thấp thì một spike RTT tạm thời cũng làm hệ mất cân bằng.
Network congestion and queueing
Giống kẹt xe, độ biến thiên delay chủ yếu do queueing:
- Nhiều node cùng gửi tới một đích → switch phải xếp hàng đưa vào link đích từng packet (network congestion). Queue đầy → packet bị drop, phải gửi lại dù mạng vẫn "chạy tốt".
- Tới máy đích mà mọi CPU core bận → OS xếp hàng request tới khi app sẵn sàng, có thể lâu tùy ý.
- Môi trường ảo hóa: OS của VM bị pause hàng chục ms khi VM khác dùng CPU; dữ liệu đến bị buffer bởi hypervisor.
- TCP flow control / congestion avoidance / backpressure giới hạn tốc độ gửi → xếp hàng ngay tại bên gửi.
- TCP coi packet mất nếu không được ack trong timeout (tính từ RTT quan sát) và tự gửi lại; app không thấy mất packet nhưng thấy delay.
TCP Versus UDP
Ứng dụng nhạy latency (video call, VoIP) dùng UDP: không flow control, không retransmit → tránh một phần nguồn gây delay biến thiên (vẫn chịu switch queue và scheduling). UDP hợp khi dữ liệu đến trễ là vô giá trị: trong cuộc gọi VoIP, không kịp truyền lại packet trước khi phát ra loa → lấp bằng im lặng; việc "retry" chuyển lên tầng con người ("Anh nói lại được không?").
Thêm nữa:
- Queueing delay dao động rất rộng khi hệ gần full capacity; hệ dư capacity dễ xả queue.
- Public cloud/multi-tenant: link, switch, NIC, CPU đều dùng chung; batch job như MapReduce dễ bão hòa link; noisy neighbor gây delay biến thiên mà bạn không kiểm soát được.
- Chỉ có thể chọn timeout bằng thực nghiệm: đo phân phối RTT trong thời gian dài trên nhiều máy, cân nhắc trade-off giữa độ trễ phát hiện và rủi ro timeout sớm.
- Tốt hơn: đo liên tục response time và jitter, tự điều chỉnh timeout — Phi Accrual failure detector (dùng trong Akka, Cassandra); TCP retransmission timeout cũng tương tự.
Synchronous Versus Asynchronous Networks
So với mạng điện thoại cố định truyền thống (rất tin cậy): mỗi cuộc gọi thiết lập một circuit — băng thông cố định được đặt trước dọc suốt tuyến. VD mạng ISDN chạy 4.000 frame/giây; mỗi cuộc gọi được cấp 16 bit mỗi frame mỗi chiều → đảm bảo gửi đúng 16 bit audio mỗi 250 micro giây. Mạng như vậy là synchronous: không queueing vì chỗ đã được đặt trước ở hop kế → latency tối đa cố định = bounded delay.
Can we not simply make network delays predictable?
- Circuit là băng thông dành riêng, không ai khác dùng; TCP thì cơ hội sử dụng băng thông sẵn có, idle thì không dùng gì.
- Ethernet và IP là packet-switched → queueing → unbounded delay; không có khái niệm circuit.
- Vì sao dùng packet switching? Vì tối ưu cho bursty traffic: web page, email, file không có yêu cầu băng thông cố định, chỉ cần xong càng nhanh càng tốt. Dùng circuit cho file phải đoán băng thông: đoán thấp → chậm, lãng phí; đoán cao → không thiết lập được circuit. TCP tự thích nghi với capacity sẵn có.
- Có mạng lai (ATM — đối thủ của Ethernet thập niên 80, không liên quan máy rút tiền), InfiniBand có flow control end-to-end ở link layer. Với QoS và admission control có thể giả lập circuit switching hoặc bounded delay theo thống kê.
Latency and Resource Utilization
Delay biến thiên là hệ quả của dynamic resource partitioning:
- Dây giữa hai tổng đài chở tối đa 10.000 cuộc gọi: chia tĩnh — dù bạn là cuộc gọi duy nhất, bạn vẫn chỉ được đúng 1 slot.
- Internet chia băng thông động: các bên chen nhau, switch quyết định từng khoảnh khắc → có queueing nhưng tối đa hóa utilization → mỗi byte rẻ hơn.
- CPU cũng vậy: chia core động giữa các thread → thread có lúc chờ trong run queue nhưng dùng phần cứng tốt hơn. VM cũng vì lý do utilization.
→ Latency guarantee khả thi nếu chia tài nguyên tĩnh (phần cứng riêng, băng thông độc quyền) nhưng đắt hơn. Multi-tenancy chia động thì rẻ nhưng delay biến thiên. Delay biến thiên không phải quy luật tự nhiên mà là kết quả của trade-off cost/benefit. Hiện QoS không được bật trong cloud multi-tenant hay internet → phải giả định congestion, queueing, unbounded delay; không có giá trị timeout "đúng".
Unreliable Clocks
Ứng dụng dùng đồng hồ để trả lời hai loại câu hỏi:
- Duration: request timeout chưa? p99 response time? QPS trung bình 5 phút qua? user ở trên site bao lâu?
- Point in time: bài viết đăng lúc nào? gửi email nhắc lúc nào? cache entry hết hạn khi nào? timestamp của log lỗi?
Trong hệ phân tán, thời gian rất rắc rối: message cần thời gian di chuyển và delay biến thiên, nên khó xác định thứ tự sự kiện giữa nhiều máy. Mỗi máy có đồng hồ phần cứng riêng (thường là quartz crystal oscillator), không hoàn toàn chính xác, chạy nhanh/chậm khác nhau. Đồng bộ phổ biến bằng NTP — chỉnh theo thời gian từ nhóm server, các server lại lấy từ nguồn chính xác hơn như GPS.
Monotonic Versus Time-of-Day Clocks
Time-of-day clocks
- Trả về ngày giờ hiện tại theo lịch (wall-clock time):
clock_gettime(CLOCK_REALTIME)trên Linux,System.currentTimeMillis()trong Java — số giây/ms kể từ epoch (0h UTC 1/1/1970), không tính leap second. - Được đồng bộ bằng NTP nên timestamp ở máy này (lý tưởng) có nghĩa giống ở máy khác.
- Nhưng nếu lệch quá xa NTP server có thể bị reset cưỡng bức và nhảy lùi về quá khứ; cộng với việc bỏ qua leap second → không phù hợp để đo thời gian trôi qua.
- Lịch sử có độ phân giải thô (bước 10 ms trên Windows cũ).
Monotonic clocks
- Dùng để đo duration (timeout, response time):
clock_gettime(CLOCK_MONOTONIC),System.nanoTime(). Được đảm bảo luôn tiến lên. - Giá trị tuyệt đối vô nghĩa (có thể là số ns kể từ khi máy khởi động); chỉ hiệu giữa hai lần đọc có nghĩa. So sánh monotonic clock giữa hai máy là vô nghĩa.
- Máy nhiều socket CPU có timer riêng mỗi CPU; OS cố gắng che giấu nhưng nên dè dặt với đảm bảo monotonic.
- NTP có thể slew (điều chỉnh tốc độ) monotonic clock tối đa ±0,05% nhưng không làm nó nhảy tới/lùi. Độ phân giải tốt (micro giây hoặc hơn).
- Trong hệ phân tán, dùng monotonic clock đo timeout thường ổn vì không cần đồng bộ giữa các node.
| Tiêu chí | Time-of-day clock | Monotonic clock |
|---|---|---|
| Trả về | Ngày giờ theo lịch (từ epoch) | Giá trị tùy ý, chỉ hiệu số có nghĩa |
| API ví dụ | CLOCK_REALTIME, System.currentTimeMillis() | CLOCK_MONOTONIC, System.nanoTime() |
| Có thể nhảy lùi? | Có (NTP reset, leap second) | Không (chỉ bị slew tối đa 0,05%) |
| So sánh giữa các máy | Được (nếu đồng bộ tốt — nhưng không chắc) | Vô nghĩa |
| Dùng cho | Point in time: timestamp, lịch hẹn, log | Duration: timeout, đo latency |
| Độ phân giải | Có thể thô trên hệ cũ | Thường micro giây hoặc tốt hơn |
Clock Synchronization and Accuracy
Time-of-day clock cần đồng bộ, và việc này kém tin cậy hơn ta nghĩ:
- Drift của quartz: phụ thuộc nhiệt độ. Google giả định drift 200 ppm → lệch 6 ms nếu resync mỗi 30 giây, 17 giây nếu resync mỗi ngày. Đây là giới hạn chính xác ngay cả khi mọi thứ hoạt động đúng.
- Lệch quá nhiều với NTP server → từ chối đồng bộ hoặc bị reset cưỡng bức → app thấy thời gian đi lùi hoặc nhảy tới.
- Node vô tình bị firewall chặn NTP → có thể không ai phát hiện trong thời gian dài.
- NTP chỉ tốt bằng network delay: đồng bộ qua internet đạt sai số tối thiểu khoảng 35 ms, spike mạng có thể gây sai số ~1 giây; delay lớn có thể làm NTP client bỏ cuộc.
- Một số NTP server sai/cấu hình sai, lệch hàng giờ. Client query nhiều server và loại outlier, nhưng vẫn là đặt cược vào "người lạ trên internet".
- Leap second tạo ra phút 59 hoặc 61 giây, từng làm crash nhiều hệ lớn. Cách xử lý tốt: smearing — NTP server "nói dối", chia điều chỉnh dần trong cả ngày.
- VM: đồng hồ ảo hóa; khi VM bị pause hàng chục ms, từ góc nhìn app đồng hồ nhảy tới đột ngột.
- Thiết bị không kiểm soát (mobile, embedded): không tin được; user cố ý chỉnh sai giờ để gian lận game.
Có thể đạt độ chính xác cao nếu đầu tư: quy định MiFID II yêu cầu quỹ giao dịch cao tần đồng bộ trong 100 micro giây so với UTC (để debug flash crash, phát hiện thao túng thị trường) — cần GPS receiver, PTP (Precision Time Protocol), triển khai và giám sát cẩn thận.
Relying on Synchronized Clocks
Đồng hồ trông đơn giản nhưng đầy bẫy: một ngày có thể không đúng 86.400 giây, time-of-day clock có thể đi lùi, giờ ở node này có thể khác xa node kia. Giống mạng, phần mềm phải được thiết kế để chịu được đồng hồ sai.
Nguy hiểm đặc biệt: đồng hồ sai rất khó bị phát hiện. CPU hỏng hay mạng sai cấu hình thì hệ ngừng chạy, sẽ sớm được sửa; còn quartz hỏng hay NTP sai cấu hình thì mọi thứ trông vẫn ổn trong khi đồng hồ trôi dần — hậu quả thường là mất dữ liệu âm thầm chứ không phải crash. → Nếu dùng phần mềm cần đồng hồ đồng bộ, phải giám sát clock offset giữa mọi máy; node lệch quá xa phải bị coi là chết và loại khỏi cluster.
Timestamps for ordering events
Ví dụ multi-leader replication: client A ghi x = 1 ở node 1, replicate sang node 3; client B tăng x ở node 3 (x = 2); cả hai được replicate tới node 2. Mỗi ghi gắn timestamp theo time-of-day clock của node gốc. Dù skew giữa node 1 và 3 dưới 3 ms (tốt hơn thực tế thường gặp), ghi x = 1 có timestamp 42.004 s còn x = 2 có 42.003 s, dù x = 2 rõ ràng xảy ra sau. Node 2 kết luận x = 1 mới hơn và drop ghi của B → phép tăng bị mất.
Đây là last write wins (LWW), dùng rộng rãi trong multi-leader và leaderless DB như Cassandra, Riak. Vấn đề cơ bản (kể cả khi timestamp sinh ở client):
- Ghi biến mất bí ẩn: node có đồng hồ chậm không ghi đè được giá trị do node đồng hồ nhanh ghi cho tới khi hết khoảng skew → mất dữ liệu tùy ý mà không báo lỗi.
- LWW không phân biệt được ghi tuần tự nhanh với ghi thực sự đồng thời → cần cơ chế causality tracking như version vector.
- Hai node có thể sinh cùng timestamp (nhất là độ phân giải ms) → cần tiebreaker (số ngẫu nhiên lớn), nhưng cũng có thể vi phạm causality.
Ngay cả với NTP tốt: gửi packet lúc 100 ms (đồng hồ bên gửi) mà tới lúc 99 ms (đồng hồ bên nhận) → như thể tới trước khi gửi. NTP không thể đủ chính xác vì độ chính xác của nó bị giới hạn bởi chính RTT mạng — cái mà ta cần đo.
Logical clock (dựa trên counter tăng dần, không dựa trên quartz) an toàn hơn để sắp thứ tự sự kiện: chỉ đo thứ tự tương đối (trước/sau), không đo giờ. Đối lập với physical clock (time-of-day, monotonic).
Clock readings have a confidence interval
Đọc được ns không có nghĩa chính xác tới ns. Với NTP server trong LAN sync mỗi phút, drift vẫn có thể vài ms; qua internet thì tốt nhất vài chục ms, spike > 100 ms khi nghẽn. Nên coi mỗi lần đọc đồng hồ là một khoảng với confidence interval (VD 95% chắc thời gian hiện tại nằm giữa 10,3 và 10,5 giây). Nếu chỉ biết ±100 ms thì các chữ số micro giây vô nghĩa.
Uncertainty tính được: GPS/đồng hồ nguyên tử gắn trực tiếp thì hãng cung cấp sai số; lấy từ server thì = drift từ lần sync cuối + sai số của server + RTT. Nhưng hầu hết hệ không lộ uncertainty: clock_gettime() không cho biết sai số là 5 ms hay 5 năm.
Ngoại lệ: Google TrueTime (Spanner) trả về [earliest, latest] — thời gian thực chắc chắn nằm trong khoảng này; độ rộng phụ thuộc thời gian kể từ lần sync cuối.
Synchronized clocks for global snapshots
- Snapshot isolation cần transaction ID tăng đơn điệu phản ánh causality. Một node thì dùng counter đơn giản; nhưng phân tán nhiều partition/nhiều DC thì sinh ID toàn cục cần phối hợp → nút thắt cổ chai. (Twitter Snowflake sinh ID xấp xỉ tăng dần bằng cách chia block ID cho các node, nhưng không đảm bảo thứ tự nhất quán với causality.)
- Dùng timestamp đồng hồ đồng bộ làm transaction ID? Được nếu đồng bộ đủ tốt — vấn đề là uncertainty.
- Spanner làm vậy qua TrueTime: nếu hai khoảng A = [A_earliest, A_latest] và B = [B_earliest, B_latest] không chồng lấn (A_latest < B_earliest) thì chắc chắn B sau A. Spanner cố ý chờ bằng độ dài confidence interval trước khi commit read-write transaction, để transaction đọc dữ liệu sau đó có khoảng không chồng lấn. Để chờ ngắn, Google đặt GPS receiver hoặc đồng hồ nguyên tử ở mỗi DC, đồng bộ trong khoảng 7 ms.
- Đây là lĩnh vực nghiên cứu; chưa phổ biến ngoài Google (tại thời điểm sách viết).
Process Pauses
Ví dụ nguy hiểm khác: DB single-leader per partition. Leader làm sao biết mình vẫn là leader? Dùng lease — như lock có timeout; chỉ một node giữ lease tại một thời điểm; leader phải renew định kỳ, chết thì lease hết hạn và node khác takeover.
while (true) {
request = getIncomingRequest();
// Đảm bảo lease còn ít nhất 10 giây
if (lease.expiryTimeMillis - System.currentTimeMillis() < 10000) {
lease = lease.renew();
}
if (lease.isValid()) {
process(request);
}
}
Sai ở đâu?
- Dựa vào đồng hồ đồng bộ: expiry do máy khác đặt (VD giờ hiện tại + 30 s) nhưng so với đồng hồ local. Lệch vài giây là hỏng.
- Kể cả chuyển sang monotonic clock local: code giả định rất ít thời gian trôi giữa lúc kiểm tra thời gian và lúc
process(request). Nếu thread bị dừng 15 giây quanh dònglease.isValid(), lease đã hết hạn và node khác đã thành leader, nhưng thread không hề biết → xử lý request một cách không an toàn trước khi vòng lặp sau phát hiện.
Pause dài như vậy có thực không? Có:
- GC stop-the-world (JVM) — từng kéo dài vài phút. Ngay cả GC "concurrent" như CMS cũng phải stop-the-world đôi khi.
- VM bị suspend/resume (lưu memory xuống disk), dùng cho live migration; pause dài tùy tốc độ ghi memory.
- Laptop gập nắp.
- Context switch của OS hoặc hypervisor chuyển sang VM khác (steal time); máy tải cao thì chờ lâu.
- Synchronous disk I/O — kể cả ngầm (Java classloader lazy-load class). I/O pause và GC pause có thể cộng dồn. Disk là network block device (Amazon EBS) thì chịu thêm delay mạng.
- Swapping/paging: truy cập memory gây page fault → chờ disk; cực đoan là thrashing. Server thường tắt swap (thà kill process còn hơn thrashing).
SIGSTOP(Ctrl-Z) dừng process cho tới khiSIGCONT; có thể do ops vô tình gửi.
Tất cả có thể preempt thread ở bất kỳ điểm nào mà thread không hề biết. Giống viết code thread-safe trên một máy — nhưng trong hệ phân tán không có shared memory, nên mutex, semaphore, atomic counter... không dùng trực tiếp được. Node phải giả định có thể bị pause lâu giữa chừng một hàm; trong lúc đó thế giới vẫn chạy và có thể tuyên bố nó đã chết; khi tỉnh lại, nó không biết mình đã ngủ.
Response time guarantees
Các nguyên nhân pause có thể loại bỏ nếu cố gắng. Hệ điều khiển máy bay, tên lửa, robot, ô tô cần phản hồi trước deadline — hard real-time system. (Ví dụ: không ai muốn túi khí bung chậm vì GC pause.)
Is real-time really real?
- "Real-time" trong embedded = được thiết kế và kiểm thử để đáp ứng timing guarantee trong mọi trường hợp; khác nghĩa "real-time" mơ hồ trên web (push dữ liệu, stream processing).
- Cần hỗ trợ từ mọi tầng: RTOS cấp CPU time đảm bảo; thư viện ghi rõ worst-case execution time; hạn chế/cấm cấp phát động (có real-time GC nhưng app không được bắt nó làm quá nhiều); test và đo đạc cực kỳ nhiều.
- Rất đắt, giới hạn ngôn ngữ/tool → chủ yếu cho thiết bị nhúng safety-critical. Real-time ≠ high-performance: thậm chí throughput thấp hơn vì ưu tiên đúng hạn.
- Với hầu hết hệ server-side, không kinh tế → phải sống chung với pause và đồng hồ không ổn định.
Limiting the impact of garbage collection
- Runtime biết tốc độ cấp phát và memory còn trống nên có thể linh hoạt thời điểm GC.
- Coi GC như planned outage ngắn: runtime cảnh báo sắp GC → app ngừng gửi request mới tới node đó, chờ xử lý xong request đang chạy, rồi mới GC → che GC khỏi client, giảm p99 (một số hệ giao dịch tài chính dùng cách này).
- Biến thể: chỉ GC object sống ngắn (nhanh), và restart process định kỳ trước khi tích lũy đủ object sống lâu cần full GC; restart từng node, chuyển traffic đi trước như rolling upgrade.
- Không ngăn hoàn toàn GC pause nhưng giảm tác động.
Knowledge, Truth, and Lies
Một node không thể biết chắc điều gì — chỉ đoán dựa trên message nhận được (hoặc không nhận được). Không có response thì không phân biệt được lỗi mạng hay lỗi node. Thay vì triết lý, ta phát biểu system model (giả định về hành vi) và thiết kế hệ đáp ứng các giả định đó; thuật toán có thể được chứng minh đúng trong một system model.
The Truth Is Defined by the Majority
Ba kịch bản:
- Asymmetric fault: node nhận được mọi message nhưng message gửi đi bị drop. Nó vẫn chạy tốt nhưng node khác không nghe được, sau timeout tuyên bố nó chết — "bị kéo ra nghĩa địa dù vẫn gào 'tôi chưa chết!'".
- Node nửa-mất-kết-nối nhận ra message không được ack, biết mạng có vấn đề, nhưng vẫn bị tuyên bố chết và không làm gì được.
- GC pause một phút: mọi thread dừng, không phản hồi; node khác tuyên bố nó chết. GC xong, node tiếp tục như không có gì — "ngóc đầu dậy khỏi quan tài" — mà không hề biết một phút đã trôi qua.
Bài học: node không thể tin phán đoán của chính mình. Hệ phân tán không thể dựa vào một node duy nhất → dùng quorum (bỏ phiếu). Kể cả quyết định tuyên bố node chết: nếu quorum nói nó chết thì nó phải coi như chết và step down, dù nó cảm thấy mình còn sống.
Quorum phổ biến nhất là đa số tuyệt đối (> 1/2): 3 node chịu được 1 lỗi, 5 node chịu được 2. An toàn vì chỉ có thể có một đa số tại một thời điểm — không thể có hai đa số với quyết định mâu thuẫn. (Chi tiết ở consensus, Ch9.)
The leader and the lock
Thường hệ cần chỉ có một thứ gì đó: một leader cho mỗi partition (tránh split brain), một client giữ lock cho một resource, một user sở hữu một username. Một node tin mình là "the chosen one" không có nghĩa quorum đồng ý — có thể nó đã bị giáng chức do mạng hoặc GC pause và leader mới đã được bầu. Nếu nó vẫn hành động như chosen one và node khác tin nó → hệ làm sai.
Ví dụ lỗi thật (HBase từng dính): storage service chỉ cho một client ghi file tại một thời điểm, client phải lấy lease từ lock service. Client 1 lấy lease rồi bị GC pause lâu, lease hết hạn; client 2 lấy lease và ghi file; client 1 tỉnh lại, vẫn tin lease hợp lệ, cũng ghi → file bị corrupt.
Fencing tokens
Kỹ thuật fencing: mỗi lần lock service cấp lock/lease, nó trả kèm một fencing token — số tăng dần mỗi lần cấp. Mỗi request ghi tới storage phải kèm token.
- Client 1 lấy lease với token 33, rồi pause lâu, lease hết hạn.
- Client 2 lấy lease với token 34, ghi với token 34.
- Client 1 tỉnh lại, gửi ghi với token 33 → storage nhớ đã xử lý token 34 lớn hơn → từ chối.
Với ZooKeeper, dùng zxid hoặc cversion làm fencing token (đảm bảo tăng đơn điệu).
Điểm mấu chốt: chính resource phải chủ động kiểm tra token — không thể dựa vào client tự kiểm tra lock status. Resource không hỗ trợ sẵn thì có thể lách (VD đưa token vào tên file). Kiểm tra ở server là điều tốt: service không nên giả định client luôn cư xử đúng, vì client thường do người khác vận hành với ưu tiên khác.
Byzantine Faults
Fencing token chặn được node vô tình sai, nhưng node cố ý phá có thể gửi token giả. Sách giả định node unreliable but honest: có thể chậm, không trả lời, state cũ, nhưng nếu trả lời thì theo đúng protocol.
Nếu node có thể "nói dối" (gửi phản hồi sai/hỏng tùy ý, VD nói đã nhận message khi chưa nhận) → Byzantine fault; bài toán đạt đồng thuận trong môi trường thiếu tin cậy là Byzantine Generals Problem.
The Byzantine Generals Problem
- Tổng quát hóa Two Generals Problem: hai tướng ở hai trại cần thống nhất kế hoạch, chỉ liên lạc qua người đưa tin có thể bị trễ/mất (như packet).
- Phiên bản Byzantine: n tướng, một số là kẻ phản bội gửi thông tin giả, không biết trước là ai.
- Tên gọi: "Byzantine" theo nghĩa quá phức tạp, quan liêu, xảo quyệt; Lamport chọn để tránh xúc phạm quốc gia nào (tên "Albanian Generals" bị khuyên không nên dùng).
Hệ Byzantine fault-tolerant vẫn đúng dù một số node không theo protocol hoặc attacker can thiệp mạng. Liên quan khi:
- Hàng không vũ trụ: bức xạ làm hỏng memory/register → node phản hồi tùy ý; hệ điều khiển bay phải chịu Byzantine fault.
- Nhiều tổ chức không tin nhau: Bitcoin/blockchain giúp các bên không tin nhau thống nhất giao dịch có xảy ra không, không cần bên trung gian.
Trong hệ của sách (datacenter một tổ chức, bức xạ thấp), thường giả định không có Byzantine fault; protocol BFT phức tạp và đắt, không thực tế. Web app vẫn phải coi client (browser) có thể độc hại → input validation, sanitization, output escaping (chống SQL injection, XSS) — nhưng giải bằng cách để server làm authority, không dùng protocol BFT.
- Bug phần mềm có thể coi là Byzantine fault, nhưng nếu mọi node chạy cùng phần mềm thì BFT không cứu được. BFT thường cần hơn 2/3 node hoạt động đúng (4 node chịu tối đa 1 lỗi) → muốn chống bug phải có 4 bản cài đặt độc lập.
- Chống tấn công: attacker chiếm được một node thường chiếm được tất cả (cùng phần mềm) → cơ chế truyền thống (authentication, access control, encryption, firewall) vẫn là chính.
Weak forms of lying
Các biện pháp thực dụng (không phải BFT đầy đủ) chống "nói dối yếu" do lỗi phần cứng, bug, cấu hình sai:
- Packet bị hỏng đôi khi lọt qua checksum TCP/UDP → thêm checksum ở tầng application.
- App public phải sanitize input (kiểm tra khoảng giá trị, giới hạn độ dài chuỗi chống DoS); service nội bộ cũng nên sanity-check cơ bản.
- NTP client cấu hình nhiều server, ước lượng sai số từng cái, kiểm tra đa số đồng ý một khoảng thời gian → server sai bị loại như outlier.
System Model and Reality
Thuật toán phân tán cần không phụ thuộc quá nhiều vào chi tiết phần cứng/phần mềm → formalize các loại fault bằng system model.
Giả định về timing:
| Timing model | Giả định | Thực tế? |
|---|---|---|
| Synchronous | Network delay, process pause, clock error đều có giới hạn trên cố định đã biết (không phải bằng 0) | Không thực tế cho hầu hết hệ — unbounded delay và pause có xảy ra |
| Partially synchronous | Hầu hết thời gian hành xử như synchronous, nhưng đôi khi vượt giới hạn tùy ý | Mô hình thực tế nhất cho nhiều hệ |
| Asynchronous | Không được giả định gì về timing, không có đồng hồ (không dùng timeout) | Rất hạn chế; chỉ một số thuật toán thiết kế được |
Giả định về node failure:
| Node model | Giả định | Ghi chú |
|---|---|---|
| Crash-stop | Node chỉ hỏng bằng cách crash và không bao giờ quay lại | Đơn giản nhất để suy luận |
| Crash-recovery | Node crash bất kỳ lúc nào và có thể quay lại sau thời gian không biết trước; stable storage (disk) còn, in-memory state mất | Kết hợp với partially synchronous là mô hình hữu ích nhất |
| Byzantine (arbitrary) | Node có thể làm bất cứ điều gì, kể cả lừa dối | Cần BFT, đắt; dùng cho aerospace, blockchain |
→ Mô hình hữu ích nhất cho hệ thực tế: partially synchronous + crash-recovery.
Correctness of an algorithm
Định nghĩa "đúng" bằng các property. Ví dụ thuật toán sinh fencing token:
- Uniqueness: không có hai request nào nhận cùng token.
- Monotonic sequence: nếu request x hoàn thành trước khi y bắt đầu thì t_x < t_y.
- Availability: node yêu cầu token và không crash thì cuối cùng sẽ nhận được response.
Thuật toán đúng trong một system model nếu luôn thỏa các property trong mọi tình huống mà model cho là có thể xảy ra. Nhưng nếu mọi node crash hoặc delay vô hạn thì không thuật toán nào làm được gì → cần phân biệt safety và liveness.
Safety and liveness
- Uniqueness và monotonic sequence là safety; availability là liveness. Dấu hiệu: liveness thường có chữ "eventually" (và đúng vậy, eventual consistency là liveness property).
- Không chính thức: safety = không có gì tồi tệ xảy ra; liveness = điều tốt cuối cùng sẽ xảy ra.
- Chính xác:
- Safety bị vi phạm thì chỉ ra được thời điểm cụ thể nó bị phá (VD thao tác trả token trùng), và không thể undo — thiệt hại đã xảy ra.
- Liveness có thể chưa thỏa tại một thời điểm (chưa nhận response) nhưng luôn còn hy vọng thỏa trong tương lai.
- Thường yêu cầu safety luôn đúng trong mọi tình huống của model (kể cả mọi node crash, cả mạng hỏng → không được trả kết quả sai). Liveness được phép kèm điều kiện: chỉ cần phản hồi nếu đa số node không crash và mạng cuối cùng hồi phục. Định nghĩa partially synchronous yêu cầu hệ cuối cùng quay về trạng thái synchronous (gián đoạn chỉ kéo dài hữu hạn).
Mapping system models to the real world
System model là trừu tượng hóa đơn giản của thực tế:
- Crash-recovery giả định stable storage sống sót qua crash — nhưng nếu disk hỏng, dữ liệu bị xóa do lỗi phần cứng/cấu hình, hay firmware bug khiến server không nhận ra ổ cứng khi reboot?
- Quorum algorithm dựa vào việc node nhớ dữ liệu đã tuyên bố lưu; node bị "mất trí nhớ" phá vỡ điều kiện quorum → phá tính đúng. Có thể cần model mới (storage hầu như sống sót) nhưng khó suy luận hơn.
- Implementation thực tế vẫn phải có code xử lý trường hợp "không thể xảy ra", dù chỉ là
printf("Sucks to be you")+exit(666)để con người dọn dẹp — đây có lẽ là khác biệt giữa computer science và software engineering.
Model trừu tượng vẫn cực kỳ giá trị: thu gọn độ phức tạp thành tập fault có thể suy luận, chứng minh thuật toán đúng. Chứng minh đúng không có nghĩa implementation luôn đúng, nhưng là bước đầu tốt — phân tích lý thuyết phát hiện lỗi có thể ẩn lâu trong hệ thật. Phân tích lý thuyết và kiểm thử thực nghiệm quan trọng như nhau.
Tổng kết chương
- Gửi packet qua mạng có thể mất hoặc trễ tùy ý; không có reply thì không biết message có tới không.
- Đồng hồ node có thể lệch nhiều, nhảy tới/lùi; dựa vào nó nguy hiểm vì hiếm khi biết khoảng sai số.
- Process có thể pause lâu bất kỳ lúc nào (GC), bị tuyên bố chết rồi sống lại mà không biết.
- Partial failure là đặc trưng định nghĩa hệ phân tán. Phát hiện lỗi đã khó (timeout không phân biệt lỗi mạng và lỗi node; node "limping" — VD NIC Gigabit tụt xuống 1 Kb/s do bug driver — còn khó xử lý hơn node chết hẳn). Chịu lỗi càng khó vì không có shared state; quyết định quan trọng cần quorum.
- Nếu bài toán giải được trên một máy thì thường nên làm vậy — đừng mở hộp Pandora. Nhưng hệ phân tán cần cho fault tolerance và low latency (đặt dữ liệu gần user), không chỉ scalability.
- Mạng/đồng hồ/process không tin cậy không phải quy luật tự nhiên: có thể có bounded delay và hard real-time, nhưng rất đắt và giảm utilization → hầu hết hệ chọn rẻ và không tin cậy.
⚠️ Hiểu lầm & cạm bẫy thường gặp
- "Timeout = node đã chết" → sai. Timeout chỉ nói không nhận được phản hồi; request có thể đã được xử lý (→ cần idempotency khi retry).
- "TCP ack = request đã được xử lý" → không; chỉ có response ở tầng application mới chắc chắn.
- Đặt timeout cố định quá ngắn → false positive, failover thừa, cascading failure. Nên đo phân phối RTT hoặc dùng adaptive detector (Phi Accrual).
- Dùng
System.currentTimeMillis()/ wall clock để đo duration hoặc timeout → có thể âm hoặc nhảy; dùng monotonic clock. - So sánh monotonic clock giữa hai máy → vô nghĩa.
- Tin rằng NTP làm đồng hồ các máy "đủ khớp" để sắp thứ tự sự kiện → LWW có thể âm thầm mất ghi; dùng logical clock / version vector.
- Tin timestamp độ phân giải micro giây là chính xác tới micro giây — sai số thực có thể hàng chục–trăm ms.
- Nghĩ đồng hồ sai sẽ gây crash dễ thấy → thực tế gây mất dữ liệu âm thầm; phải giám sát clock offset.
- Nghĩ lease/distributed lock (VD Redis lock đơn giản) là đủ an toàn → GC pause/VM pause làm hai client cùng tin mình giữ lock; cần fencing token kiểm tra ở resource.
- Nghĩ "code chạy nhanh nên không bị ngắt giữa check và act" → thread có thể bị pause bất kỳ đâu.
- Nghĩ thêm thiết bị mạng dự phòng sẽ loại bỏ network fault → lỗi con người (cấu hình) vẫn là nguyên nhân chính.
- Nghĩ link hoạt động một chiều thì chiều kia cũng ổn → asymmetric fault có thật.
- Nhầm "real-time" trên web với hard real-time; nhầm real-time với high-performance.
- Nghĩ BFT chống được bug phần mềm hay hacker → không, nếu mọi node chạy cùng code.
💼 Áp dụng thực tế & phỏng vấn
- Health check & failover (load balancer, leader election): thảo luận trade-off timeout ngắn/dài, nguy cơ split brain, dùng quorum (ZooKeeper, etcd, Raft) để quyết định leader thay vì một node tự quyết.
- Distributed lock là câu hỏi kinh điển ("thiết kế job scheduler đảm bảo một job chỉ chạy một lần"): trả lời lease + fencing token (ZooKeeper
zxid, etcd revision), storage từ chối token cũ; nêu vấn đề GC pause. Tham chiếu tranh luận Redlock: lock không có fencing thì không an toàn cho correctness, chỉ ổn cho efficiency. - Retry & idempotency: vì không phân biệt được "request mất" và "response mất", mọi API có side effect (thanh toán, gửi email) nên có idempotency key; retry với exponential backoff + jitter để không làm overload tệ hơn.
- ID & ordering: Snowflake ID tiện nhưng không đảm bảo thứ tự causal; muốn thứ tự đúng cần logical clock (Lamport, vector clock) hoặc sequencer/consensus. Nhắc đến Spanner TrueTime (commit wait) khi thảo luận global consistent snapshot; CockroachDB dùng hybrid logical clock với giới hạn max clock offset và tự loại node lệch quá.
- Cassandra/DynamoDB LWW: biết rằng clock skew có thể làm mất ghi; dùng conditional write/lightweight transaction hoặc CRDT khi cần.
- Tuning JVM service: GC pause ảnh hưởng p99 và có thể gây bị tuyên bố chết (Kafka broker/ZooKeeper session timeout, Elasticsearch node rời cluster). Chọn GC low-pause (G1, ZGC, Shenandoah), tắt swap, theo dõi steal time trên cloud.
- Chaos engineering: Chaos Monkey, Jepsen (Kyle Kingsbury) test hành vi khi partition — đáng nhắc khi được hỏi "làm sao bạn biết hệ chịu được lỗi mạng?".
- Câu trả lời mẫu cho "tại sao hệ phân tán khó?": partial failure nondeterministic + mạng async unbounded delay + đồng hồ không tin cậy + process pause → không node nào biết chắc sự thật; phải dựa vào quorum, safety luôn giữ, liveness có điều kiện.
- Khi thiết kế, nêu rõ system model giả định (thường partially synchronous + crash-recovery) và property cần giữ (safety vs liveness).