System Design · طراحی سیستم سنیورSenior ~49 دقیقه مطالعه~43 min read

نظریهٔ سیستم‌های توزیع‌شدهDistributed Systems Theory

نظریهٔ سیستم‌های توزیع‌شده به تو یاد می‌دهد که چرا «شکست جزئی»، ساعت‌های نامطمئن و شبکهٔ غیرقابل‌اعتماد همه‌چیز را سخت می‌کنند و چطور با ترتیب علّی، اجماع (Raft/Paxos)، کوئوروم، مدل‌های سازگاری و الگوهایی مثل saga و fencing token سیستمی بسازی که در دنیای واقعی نمی‌شکند.Distributed systems theory teaches you why partial failure, unreliable clocks, and an unreliable network make everything hard — and how causal ordering, consensus (Raft/Paxos), quorums, consistency models, and patterns like saga and fencing tokens let you build systems that survive the real world.

پیش‌نیاز:Prerequisites: مبانیِ طراحیِ سیستمSystem Design Fundamentals


بذار با یک حقیقت تلخ شروع کنم که مرز بین یک مهندس متوسط و یک سنیور واقعی است: در یک برنامهٔ تک‌ماشینه، یا همه‌چیز کار می‌کند یا همه‌چیز خراب است. در یک سیستم توزیع‌شده، همیشه یک تکه‌ای خراب است و تو نمی‌دانی کدام تکه. به این می‌گویند «شکست جزئی» (partial failure) و تمام سختیِ این حوزه از همین یک جمله بیرون می‌آید.

وقتی متد save() را روی یک دیتابیس لوکال صدا می‌زنی، دو حالت داری: موفق یا خطا. اما وقتی همان درخواست از روی شبکه به یک سرویس دیگر می‌رود، یک حالت سومِ شیطانی هم ظاهر می‌شود: نمی‌دانی. شاید درخواست اصلاً نرسید. شاید رسید و پردازش شد ولی جوابش گم شد. شاید هنوز دارد پردازش می‌شود و سه ثانیهٔ دیگر جواب می‌دهد. این «نمی‌دانم» قلبِ نظریهٔ سیستم‌های توزیع‌شده است.

این فصل قرار نیست فرمول‌های آکادمیک را طوطی‌وار ردیف کند. قرار است بعد از خواندنش، وقتی در مصاحبه یا در جلسهٔ طراحی کسی گفت «خب اینجا اگر network partition بشه چی؟»، تو نه با ترس بلکه با یک نقشهٔ ذهنیِ دقیق جواب بدهی: بدانی چه چیزی تئوریماً ممکن نیست، چه trade-offای داری، و کدام الگو دردت را دوا می‌کند.

نقشهٔ راه این فصل

۱) چرا توزیع‌شدن سخت است: شکست جزئی، شبکهٔ نامطمئن، ساعت‌های دروغگو، و «هشت دروغ». ۲) مدل‌های شکست: crash، omission، timing، و Byzantine. ۳) زمان و ترتیب: happens-before، ساعت Lamport و vector clock. ۴) اجماع (consensus): قضیهٔ FLP، Paxos در یک نگاه، و Raft با جزئیات (leader election و log replication). ۵) تکثیر و کوئوروم: قانون R+W>N و انواع سازگاری خواندن/نوشتن. ۶) CAP و PACELC به‌صورت عمیق. ۷) طیف سازگاری: linearizable تا eventual و read-your-writes. ۸) تحویل پیام: at-least-once در برابر exactly-once (و چرا واقعاً effectively-once است). ۹) تراکنش توزیع‌شده: 2PC/3PC در برابر saga. ۱۰) gossip، anti-entropy، CRDT، و در آخر split-brain و fencing token.


۱) چرا توزیع‌شدن این‌قدر سخت است؟

شبکه به تو دروغ می‌گوید

در سال ۱۹۹۴، Peter Deutsch و همکارانش در Sun فهرستی درست کردند به اسم «هشت دروغِ محاسبات توزیع‌شده» (The Eight Fallacies of Distributed Computing). این‌ها فرض‌های غلطی هستند که هر مهندسِ تازه‌کار ناخودآگاه می‌کند و بعد در تولید (production) با پوست و استخوان تاوانش را می‌دهد:

۱) شبکه قابل‌اعتماد است. ۲) تأخیر (latency) صفر است. ۳) پهنای‌باند بی‌نهایت است. ۴) شبکه امن است. ۵) توپولوژی ثابت است. ۶) یک مدیرِ شبکه وجود دارد. ۷) هزینهٔ انتقال صفر است. ۸) شبکه همگن (homogeneous) است.

هر هشت فرض غلط‌اند — و هر کدام یک incident است

هر باگِ عجیبِ «فقط گاهی اتفاق می‌افتد» در سیستم‌های توزیع‌شده، معمولاً یکی از این هشت فرض است که یک نفر جایی در کد ناخودآگاه کرده. مثلاً restTemplate.getForObject(url) بدون timeout یعنی «فرض کردم شبکه قابل‌اعتماد است و تأخیر صفر است». نتیجه؟ یک سرویسِ کُند، تمام thread‌های Tomcat تو را قفل می‌کند و کل برنامه‌ات می‌خوابد. به این می‌گویند cascading failure.

مسئلهٔ دو ژنرال

دو ژنرال در دو طرف یک دره هستند و باید هم‌زمان به دشمن حمله کنند. تنها راه ارتباط، فرستادنِ پیک از میان دره است، جایی که پیک ممکن است اسیر شود. ژنرال اول پیام می‌فرستد «فردا سحر حمله». ولی از کجا بداند پیام رسیده؟ منتظر تأییدیه می‌ماند. ژنرال دوم تأییدیه می‌فرستد، ولی از کجا بداند تأییدیه‌اش رسیده؟ او هم باید منتظر تأییدِ تأییدیه بماند… این حلقه هیچ‌وقت تمام نمی‌شود. این «مسئلهٔ دو ژنرال» ثابت می‌کند که روی یک کانالِ نامطمئن، توافقِ قطعی (agreement) به‌صورت تئوری غیرممکن است. هر بار که در کد catch (TimeoutException) می‌نویسی، در واقع داری با روحِ همین دو ژنرال دست‌وپنجه نرم می‌کنی: نمی‌دانی طرف کارش را انجام داد یا نه.

چرا timeout جواب قطعی نیست

نکتهٔ ظریفی که خیلی‌ها جا می‌اندازند: timeout هیچ‌وقت به تو نمی‌گوید طرف مقابل مُرده یا فقط کُند است. اگر ۵ ثانیه صبر کردی و جواب نیامد، سه احتمال داری که از هم قابل‌تفکیک نیستند: (الف) سرور مُرده، (ب) سرور زنده ولی کُند است و ۶ ثانیهٔ دیگر جواب می‌دهد، (ج) جواب آماده بود ولی در راهِ برگشت گم شد. این عدمِ‌قطعیت، تصمیم‌گیری را دشوار می‌کند: اگر فرض کنی مُرده و کار را دوباره بفرستی، ممکن است کار دو بار انجام شود (duplicate). اگر فرض کنی زنده و صبر کنی، ممکن است تا ابد منتظر بمانی.


۲) مدل‌های شکست: دشمنت را دقیق بشناس

قبل از اینکه بخواهی سیستمی «مقاوم» بسازی، باید دقیق تعریف کنی مقاوم در برابر چه نوع خرابی. مدل‌های شکست از «مهربان» به «شرور» مرتب می‌شوند:

مدل شکست چه اتفاقی می‌افتد مثال واقعی سختی مقابله
Crash-stop (fail-stop) نود می‌ایستد و دیگر برنمی‌گردد pod با OOMKilled کشته می‌شود ساده‌ترین
Crash-recovery نود می‌افتد ولی بعداً برمی‌گردد (شاید با state قدیمی) ری‌استارت سرویس، از دست رفتن حافظهٔ RAM متوسط
Omission نود بعضی پیام‌ها را می‌اندازد ولی زنده است صف پرِ TCP، بسته‌های گم‌شده متوسط
Timing (performance) نود جواب می‌دهد ولی خیلی دیر GC pause طولانی، دیسکِ کُند زیاد
Byzantine (arbitrary) نود رفتار دلخواه/بدخواهانه دارد، پیامِ متناقض می‌فرستد باگِ خرابی حافظه، نودِ هک‌شده شدیدترین
قضاوت سنیور: تقریباً همیشه crash-recovery را هدف بگیر

اکثر سیستم‌های صنعتی (Kafka، etcd، PostgreSQL، اکثر دیتابیس‌ها) برای مدلِ crash-recovery با omission طراحی شده‌اند، نه Byzantine. مقابله با Byzantine (که در بلاک‌چین لازم است) هزینهٔ هنگفتی دارد: به‌جای اکثریتِ ساده ($n/2+1$)، به $3f+1$ نود برای تحملِ $f$ خیانتکار نیاز داری. در یک سیستم سازمانیِ درونی که به نودهای خودت اعتماد داری، صرفِ آن هزینه احمقانه است. در مصاحبه اگر بگویی «سیستم ما Byzantine fault tolerant است» بلافاصله می‌پرسند «چرا؟ نودهایت که خودی‌اند» — و اگر جواب خوبی نداشتی، ضعف نشان دادی.

نکتهٔ حیاتی: timing failure از crash failure خطرناک‌تر است، چون یک نودِ کُند بدتر از یک نودِ مُرده است. نودِ مُرده را از حلقه خارج می‌کنی و تمام. اما نودِ کُند در health check سبز نشان می‌دهد، هنوز ترافیک می‌گیرد، و همه را کُند می‌کند. به این پدیده «gray failure» می‌گویند و علتِ خیلی از incident‌های واقعی است.


۳) زمان و ترتیب: چرا نمی‌توانی به ساعت اعتماد کنی

ساعتِ دیوار دروغ می‌گوید

بزرگ‌ترین اشتباهِ مبتدی در سیستم توزیع‌شده این است که فکر کند System.currentTimeMillis() روی دو ماشین قابل‌مقایسه است. نیست.

هرگز برای اندازه‌گیری مدت‌زمان از currentTimeMillis استفاده نکن

System.currentTimeMillis() ساعتِ دیوار (wall clock) است و با NTP هماهنگ می‌شود؛ یعنی می‌تواند ناگهان به عقب بپرد. اگر دو تایم‌استمپ از آن بگیری و کم کنی، ممکن است عددِ منفی بگیری! برای اندازه‌گیریِ مدت‌زمان همیشه از System.nanoTime() (ساعتِ monotonic که فقط جلو می‌رود) استفاده کن. برای منطقِ کسب‌وکار که «الان چه ساعتی است» را می‌خواهد از Instant.now() استفاده کن، ولی هیچ‌وقت دو Instant از دو ماشینِ مختلف را برای تصمیمِ ترتیبی مقایسه نکن.

// اشتباه: با NTP ممکن است منفی شود
long start = System.currentTimeMillis();
doWork();
long elapsed = System.currentTimeMillis() - start; // می‌تواند غلط باشد!

// درست: ساعت monotonic فقط جلو می‌رود
long startNanos = System.nanoTime();
doWork();
long elapsedNanos = System.nanoTime() - startNanos; // همیشه معتبر

مشکل عمیق‌تر است: حتی با NTP، ساعتِ دو سرور معمولاً چند میلی‌ثانیه (و در بدترین حالت چند ثانیه) با هم اختلاف دارند (clock skew). پس این کد که «رکورد با تایم‌استمپ جدیدتر برنده است» (last-write-wins) ذاتاً باگ‌دار است: نوشته‌ای که واقعاً بعداً اتفاق افتاده ممکن است تایم‌استمپِ کوچک‌تری داشته باشد و اشتباهاً دور ریخته شود.

happens-before: ترتیبِ علّی به‌جای ترتیبِ ساعت

Leslie Lamport در مقالهٔ افسانه‌ایِ ۱۹۷۸ گفت: فراموش کن که «چه ساعتی» اتفاق افتاد. فقط بپرس «آیا A می‌توانست بر B اثر بگذارد؟». این رابطه را «happens-before» می‌نامیم و با $\rightarrow$ نشان می‌دهیم:

  • اگر A و B در یک process باشند و A زودتر اجرا شود: $A \rightarrow B$.
  • اگر A فرستادنِ یک پیام باشد و B دریافتِ همان پیام: $A \rightarrow B$.
  • ترانزیتیو است: اگر $A \rightarrow B$ و $B \rightarrow C$ آنگاه $A \rightarrow C$.

اگر نه $A \rightarrow B$ و نه $B \rightarrow A$، آن‌ها همزمان (concurrent) هستند — یعنی هیچ رابطهٔ علّی‌ای ندارند و ترتیبشان بی‌معنی است.

نمودار زیر رابطهٔ happens-before را نشان می‌دهد: فرستادن پیام یک لبهٔ علّی می‌سازد. | This diagram shows the happens-before relation: sending a message creates a causal edge.

sequenceDiagram
    participant P1 as Process 1
    participant P2 as Process 2
    P1->>P1: event a
    P1->>P2: message m (event b = send)
    Note over P2: event c = receive m
    P2->>P2: event d
    Note over P1,P2: a -> b -> c -> d (causal chain)

ساعت Lamport: یک عددِ منطقی به هر رویداد

ساعتِ Lamport یک شمارندهٔ ساده روی هر process است که این رابطهٔ علّی را به یک عدد ترجمه می‌کند. قانونش سه خط است:

public class LamportClock {
    private long counter = 0;

    // قبل از هر رویداد محلی: یکی زیاد کن
    public synchronized long tick() {
        return ++counter;
    }

    // موقع فرستادن پیام: tick بزن و عدد را همراه پیام بفرست
    public synchronized long onSend() {
        return ++counter;
    }

    // موقع دریافت: max خودت و فرستنده، سپس یکی زیاد کن
    public synchronized long onReceive(long senderClock) {
        counter = Math.max(counter, senderClock) + 1;
        return counter;
    }
}

تضمینِ Lamport: اگر $A \rightarrow B$ آنگاه $L(A) < L(B)$. اما عکسش درست نیست: اگر $L(A) < L(B)$ لزوماً یعنی A علّتِ B نیست؛ شاید فقط همزمان بوده‌اند. این محدودیتِ اصلیِ ساعت Lamport است.

ساعت Lamport فقط یک‌طرفه است

$A \rightarrow B \implies L(A) < L(B)$ درست است. ولی $L(A) < L(B) \implies A \rightarrow B$ غلط است. یعنی با ساعت Lamport نمی‌توانی تشخیص بدهی دو رویداد علّی‌اند یا صرفاً همزمان. برای این تشخیص به vector clock نیاز داری.

vector clock: تشخیصِ همزمانی

vector clock این ضعف را برطرف می‌کند. به‌جای یک عدد، هر نود یک آرایه نگه می‌دارد: یک خانه به‌ازای هر نود. خانهٔ خودش را برای هر رویداد زیاد می‌کند و موقع دریافت، خانه‌به‌خانه maximum می‌گیرد.

public class VectorClock {
    private final int nodeId;
    private final long[] vector;

    public VectorClock(int nodeId, int numNodes) {
        this.nodeId = nodeId;
        this.vector = new long[numNodes];
    }

    public synchronized long[] tick() {
        vector[nodeId]++;
        return vector.clone();
    }

    public synchronized void onReceive(long[] other) {
        for (int i = 0; i < vector.length; i++) {
            vector[i] = Math.max(vector[i], other[i]);
        }
        vector[nodeId]++;
    }

    // A قبل از B است اگر همهٔ خانه‌ها <= باشند و حداقل یکی <
    public static boolean happensBefore(long[] a, long[] b) {
        boolean strictlyLess = false;
        for (int i = 0; i < a.length; i++) {
            if (a[i] > b[i]) return false;
            if (a[i] < b[i]) strictlyLess = true;
        }
        return strictlyLess;
    }
}

حالا اگر نه $A \rightarrow B$ و نه $B \rightarrow A$، با قطعیت می‌فهمی که concurrent بوده‌اند — و این دقیقاً همان چیزی است که یک دیتابیسِ eventually consistent مثل نسخه‌های اولیهٔ Riak/DynamoDB برای تشخیصِ conflict به آن نیاز دارد. وقتی دو نوشتهٔ همزمان به یک کلید می‌خورند، سیستم هر دو را به‌عنوان «sibling» نگه می‌دارد و از application می‌خواهد conflict را حل کند.

فرق ساعت Lamport و vector clock چیست و کِی کدام را استفاده می‌کنی؟

هر دو ترتیبِ علّی (happens-before) را ثبت می‌کنند، ولی با قدرتِ متفاوت. ساعت Lamport یک عددِ اسکالر است و فقط تضمین می‌کند اگر A علّتِ B باشد، عددش کوچک‌تر است — ولی نمی‌تواند «همزمانی» را تشخیص دهد. سبک است و برای ساختنِ یک total order قراردادی (مثلاً برای شکستنِ تساوی) عالی است. vector clock یک آرایه به طول تعدادِ نودها است و می‌تواند دقیقاً بگوید دو رویداد علّی‌اند یا concurrent — که برای تشخیصِ conflict در سیستم‌های multi-master ضروری است. هزینه‌اش این است که اندازه‌اش با تعداد نودها رشد می‌کند (مشکلِ scale). در عمل: اگر فقط ترتیبِ کلی می‌خواهی Lamport؛ اگر باید conflict را تشخیص و حل کنی vector clock (یا نسخهٔ فشرده‌اش، dotted version vector).


۴) اجماع (Consensus): سخت‌ترین مسئلهٔ این حوزه

اجماع یعنی چند نود روی یک مقدار به توافق برسند، حتی اگر بعضی نودها بیفتند. این بنیانِ همه‌چیز است: انتخابِ leader، commit کردنِ تراکنش، ثبتِ ترتیبِ لاگ. سه ویژگی می‌خواهیم: (۱) Agreement: همه روی یک مقدار توافق کنند. (۲) Validity: مقدارِ توافق‌شده باید یکی از مقادیرِ پیشنهادی باشد. (۳) Termination: در نهایت تصمیم گرفته شود.

قضیهٔ FLP: اجماعِ قطعی در شبکهٔ ناهمگام غیرممکن است

در ۱۹۸۵، Fischer، Lynch و Paterson ثابت کردند که در یک سیستمِ کاملاً ناهمگام (asynchronous) — جایی که کرانی برای تأخیرِ پیام نداری — حتی با فقط یک نودِ خراب، هیچ الگوریتمِ قطعی‌ای نمی‌تواند همیشه به اجماع برسد. پس چطور etcd و Raft کار می‌کنند؟ چون در عمل تقلب می‌کنند: با timeout (فرضِ همگامیِ نسبی) و تصادفی‌بودن (randomized timeouts) دورِ FLP می‌زنند. FLP نمی‌گوید اجماع غیرممکن است؛ می‌گوید نمی‌توانی هم‌زمان همیشه ایمن و همیشه زنده باشی. الگوریتم‌های واقعی safety را قربانی نمی‌کنند ولی گاهی liveness را (ممکن است یک دور بیشتر طول بکشد).

قضیهٔ FLP را ساده توضیح بده — یعنی اجماع غیرممکن است؟

نه، غیرممکن نیست. FLP می‌گوید در یک سیستمِ کاملاً ناهمگام (بدونِ هیچ کرانی برای تأخیرِ پیام) هیچ الگوریتمِ قطعی‌ای نمی‌تواند هم‌زمان هر سه چیز را تضمین کند: safety، liveness، و تحملِ حتی یک خرابی. شهودش: وقتی نودی جواب نمی‌دهد، نمی‌توانی تشخیص بدهی «مُرده» یا «فقط کُند» است؛ اگر برای ایمنی صبر کنی ممکن است تا ابد بلاک شوی (نقضِ liveness)، و اگر جلو بروی ممکن است آن نود بعداً با تصمیمی متفاوت برگردد (نقضِ safety). راهِ فرارِ سیستم‌های واقعی: فرضِ همگامیِ نسبی با timeout و کمی تصادفی‌بودن. آن‌ها safety را هرگز قربانی نمی‌کنند و فقط در بدترین حالت liveness را موقتاً از دست می‌دهند (یک دورِ انتخابات بیشتر طول می‌کشد). پس FLP یک محدودیتِ تئوریک است، نه یک بن‌بستِ عملی.

Paxos در یک نگاه

Paxos (کارِ همان Lamport) اولین الگوریتمِ اثبات‌شدهٔ اجماع بود و در قلبِ سیستم‌هایی مثل Google Chubby و Spanner است. سه نقش دارد: Proposer، Acceptor، Learner. در دو فاز کار می‌کند: فاز ۱ (Prepare/Promise) که در آن یک proposer با یک شمارهٔ یکتا اجازه می‌گیرد، و فاز ۲ (Accept/Accepted) که مقدار را جا می‌اندازد. اکثریتِ acceptorها کافی است. مشکلِ Paxos شهرتِ بدش در فهمیدنی‌نبودن است؛ خودِ Lamport مجبور شد مقالهٔ «Paxos Made Simple» را بنویسد و باز هم کسی راحت پیاده‌اش نمی‌کند. Multi-Paxos نسخهٔ عملی‌اش برای یک دنبالهٔ تصمیم است.

Raft: اجماعِ فهمیدنی

Raft در ۲۰۱۴ توسط Ongaro و Ousterhout ساخته شد با یک هدفِ صریح: قابل‌فهم بودن. امروز موتورِ اجماع در etcd (و در نتیجه Kubernetes)، Consul، RabbitMQ (quorum queues)، ClickHouse، و حالت KRaft کافکا است. Raft مسئله را به سه زیرمسئلهٔ مستقل می‌شکند: انتخابِ leader، تکثیرِ لاگ، و ایمنی.

هر نود یکی از سه حالت را دارد و بر اساس term (یک شمارندهٔ منطقیِ زمان که نقشِ «دورهٔ ریاست‌جمهوری» را بازی می‌کند) بین آن‌ها جابه‌جا می‌شود.

نمودار حالت زیر چرخهٔ عمرِ یک نودِ Raft را نشان می‌دهد. | This state diagram shows the lifecycle of a Raft node.

stateDiagram-v2
    [*] --> Follower
    Follower --> Candidate: election timeout\n(no heartbeat)
    Candidate --> Candidate: split vote\n(new term, retry)
    Candidate --> Leader: wins majority
    Candidate --> Follower: sees higher term\nor new leader
    Leader --> Follower: discovers higher term
    Leader --> Leader: send heartbeats

انتخابِ leader: هر follower یک election timeout تصادفی (معمولاً ۱۵۰ تا ۳۰۰ میلی‌ثانیه) دارد. اگر در این مدت heartbeat از leader نگیرد، فرض می‌کند leader مُرده: term خودش را یکی زیاد می‌کند، به candidate تبدیل می‌شود، به خودش رأی می‌دهد و از بقیه رأی (RequestVote) می‌خواهد. هر نود در هر term فقط یک بار رأی می‌دهد. اگر candidate اکثریت را گرفت، leader می‌شود و شروع به فرستادنِ heartbeat می‌کند. تصادفی‌بودنِ timeout کلید است: باعث می‌شود دو نود به‌ندرت هم‌زمان candidate شوند و split vote کمتر پیش بیاید — این دقیقاً همان جایی است که Raft دورِ FLP می‌زند.

انتخابِ leader مثل انتخابِ رئیسِ جلسه

تصور کن در یک جلسه رئیس (leader) هر چند ثانیه یک بار می‌گوید «من هستم، ادامه بدید» (heartbeat). اگر برای مدتی صدایی نشنیدی (election timeout)، فرض می‌کنی رئیس رفته. اما به‌جای اینکه همه هم‌زمان داد بزنند «من رئیس می‌شوم» (که هرج‌ومرج است)، هر کس یک تایمرِ تصادفی دارد؛ اولی که تایمرش تمام شود دستش را بلند می‌کند و می‌گوید «دورِ جدید، من نامزدم، به من رأی بدهید». اگر اکثریتِ اتاق موافقت کردند، رئیسِ جدید است. term همان «شمارهٔ دورِ جلسه» است تا کسی با یک دستورِ کهنه از دورِ قبل، جلسه را به‌هم نریزد.

تکثیرِ لاگ (log replication): فقط leader از client نوشته می‌گیرد. هر دستور را به‌عنوان یک entry به لاگش اضافه می‌کند و در AppendEntries به followerها می‌فرستد. وقتی اکثریت نود entry را روی دیسک نوشتند، leader آن را «committed» علامت می‌زند، به state machine اعمال می‌کند و به client جواب می‌دهد. این «اکثریتِ نوشتن پیش از commit» تضمین می‌کند حتی اگر leader بلافاصله بیفتد، هر leaderِ جدید (که باید رأیِ اکثریت را داشته باشد) قطعاً آن entry را دارد.

چرا Raft به تعدادِ فردِ نود (۳ یا ۵) نیاز دارد؟

چون اجماع بر پایهٔ اکثریت (quorum = $n/2+1$) است و تعدادِ فرد بهترین نسبتِ «تحملِ خرابی به هزینه» را می‌دهد. با ۳ نود، quorum ۲ است و می‌توانی ۱ خرابی را تحمل کنی. با ۴ نود، quorum ۳ است — یعنی باز هم فقط ۱ خرابی، ولی با یک ماشینِ اضافه که هیچ سودی نمی‌رساند و فقط احتمالِ خرابی را بالا می‌برد! با ۵ نود، quorum ۳ است و ۲ خرابی را تحمل می‌کنی. پس همیشه فرد: ۳ برای اکثر کارها، ۵ برای سیستم‌های حیاتی. عددِ زوج نه‌تنها کمکی نمی‌کند بلکه ریسکِ split-brain را در پارتیشنِ ۵۰-۵۰ زیاد می‌کند.

در تولید، به latencyِ اجماع حساس باش

هر نوشته در یک سیستمِ Raft-based باید منتظرِ fsync روی دیسکِ اکثریتِ نودها بماند. یعنی throughputِ نوشتنت را دیسکِ کُندترین نودِ داخلِ quorum و round-trip شبکه محدود می‌کند. به همین دلیل etcd را در یک datacenter (یا AZهای نزدیک) می‌گذارند، نه پخش در قاره‌ها. اگر cluster اعضایش را روی قاره‌های مختلف پخش کنی، هر نوشته چند صد میلی‌ثانیه طول می‌کشد. سنیورها این را از قبل حساب می‌کنند؛ جونیورها بعد از incident می‌فهمند.


۵) تکثیر و کوئوروم: قانونِ طلاییِ R+W>N

وقتی داده را روی $N$ نسخه (replica) کپی می‌کنی، دو عدد را خودت انتخاب می‌کنی: برای هر نوشتن، منتظرِ تأییدِ $W$ نسخه می‌مانی؛ برای هر خواندن، از $R$ نسخه می‌پرسی و جدیدترین را برمی‌داری. این مدلِ Dynamo-style (که Cassandra، Riak و اسکایلا از آن استفاده می‌کنند) به تو اجازه می‌دهد سازگاری را per-request تنظیم کنی.

قانونِ R + W > N

اگر $R + W > N$ باشد، مجموعهٔ نودهایی که آخرین بار نوشتی و مجموعه‌ای که الان می‌خوانی حتماً حداقل یک نودِ مشترک دارند (اصلِ لانهٔ کبوتری). آن نودِ مشترک جدیدترین مقدار را دارد، پس خواندنت قطعاً آخرین نوشته را می‌بیند: به این «strong consistency» در مدلِ quorum می‌گویند. اگر $R + W \le N$، ممکن است داده‌ای بخوانی که هنوز به آن نودها نرسیده — یعنی eventual consistency.

مثلاً با $N=3$: انتخابِ $W=2, R=2$ متعادل است ($2+2>3$). اگر خواندنِ سریع می‌خواهی $W=3, R=1$ (نوشتنِ کُند، خواندنِ فوری). اگر نوشتنِ سریع می‌خواهی $W=1, R=3$. و $W=1, R=1$ یعنی حداکثر سرعت ولی هیچ تضمینی نیست.

پیکربندی (N=3) معنی مناسبِ ریسک
W=3, R=1 نوشتن روی همه، خواندن از یکی داده‌ای که زیاد خوانده می‌شود نوشتن اگر یک نود بیفتد بلاک می‌شود
W=2, R=2 quorum متعادل حالتِ پیش‌فرضِ اکثر سیستم‌ها متعادل
W=1, R=1 سریع‌ترین، بی‌تضمین متریک، لاگ، cache خواندنِ کهنه (stale)
quorumِ لغزان (sloppy quorum) امنیت را می‌شکند

بعضی سیستم‌ها برای بالا نگه‌داشتنِ availability از «sloppy quorum» استفاده می‌کنند: اگر نودهای اصلی در دسترس نباشند، نوشته را روی نودهای جایگزین (hinted handoff) می‌گذارند. این availability را بالا می‌برد ولی تضمینِ $R+W>N$ را می‌شکند: ممکن است خواندن و نوشتن هیچ نودِ مشترکی نداشته باشند و داده‌ٔ کهنه بخوانی. این trade-off را باید بدانی؛ Cassandra با consistency levelهایی مثل QUORUM در برابر LOCAL_QUORUM دقیقاً همین را کنترل می‌کند.

به‌صورتِ شهودی توضیح بده چرا R+W>N سازگاریِ قوی می‌دهد.

تصور کن N=۳ نسخه داری با W=۲ و R=۲. وقتی می‌نویسی، داده روی حداقل ۲ نود از ۳ نود می‌نشیند. وقتی می‌خوانی، از ۲ نود از ۳ نود می‌پرسی. سؤال: آیا ممکن است این دو مجموعهٔ ۲تایی هیچ نودِ مشترکی نداشته باشند؟ نه — چون ۲+۲=۴ بزرگ‌تر از ۳ است، طبقِ اصلِ لانهٔ کبوتری حتماً حداقل یک نود در هر دو مجموعه هست. آن نودِ مشترک آخرین نوشته را دارد، پس خواندنت قطعاً جدیدترین مقدار را می‌بیند (به‌شرطِ اینکه با نسخه‌گذاری، جدیدترین را از بین جواب‌ها تشخیص بدهی). این کلِ جادوی R+W>N است: تضمینِ هم‌پوشانی. اگر R+W≤N بود، این هم‌پوشانی تضمین نمی‌شد و ممکن بود از نودهایی بخوانی که هنوز نوشتهٔ جدید را ندیده‌اند.

برای جبرانِ ناسازگاری، این سیستم‌ها دو مکانیزمِ پس‌زمینه دارند: read repair (وقتی یک خواندن نسخه‌های مختلف دید، جدیدترین را به نودهای عقب‌مانده پس می‌فرستد) و anti-entropy با درختِ Merkle که در بخش gossip توضیح می‌دهم.


۶) CAP و PACELC: تصمیمِ بنیادین

CAP: انتخابِ اجباریِ زمانِ پارتیشن

قضیهٔ CAP (که Eric Brewer مطرح و Gilbert و Lynch اثبات کردند) سه خاصیت را نام می‌برد: Consistency (همه یک مقدار می‌بینند)، Availability (هر درخواست جواب می‌گیرد)، Partition tolerance (سیستم با وجودِ قطعِ شبکه کار می‌کند).

بزرگ‌ترین سوءتفاهمِ صنعت این است که فکر می‌کنند «۲ تا از ۳ تا را انتخاب کن». این غلط است. حقیقت ظریف‌تر است:

P را انتخاب نمی‌کنی — پارتیشن به تو تحمیل می‌شود

در یک سیستمِ توزیع‌شدهٔ واقعی، پارتیشنِ شبکه یک واقعیتِ فیزیکی است نه یک گزینه؛ کابل قطع می‌شود، سوییچ می‌افتد. پس P را همیشه باید تحمل کنی. انتخابِ واقعیِ تو فقط این است: هنگامِ پارتیشن، C را قربانی کنی یا A را؟ یعنی CAP در عمل یعنی «CP یا AP». اگر گفتی «سیستمِ ما CA است» یعنی گفته‌ای «سیستمِ ما در برابر قطعِ شبکه هیچ تضمینی ندارد» — که تقریباً همیشه اشتباه است. یک دیتابیسِ تک‌نودی CA است، ولی به‌محضِ توزیع‌شدن باید بین CP و AP انتخاب کنی.

  • CP (مثل etcd، ZooKeeper، HBase، MongoDB با majority): هنگامِ پارتیشن، طرفِ اقلیت از پاسخ‌دادن سر باز می‌زند تا داده‌ٔ متناقض ندهد. سازگاری را حفظ می‌کند، دسترس‌پذیری را قربانی.
  • AP (مثل Cassandra، DynamoDB، Riak): هر دو طرفِ پارتیشن جواب می‌دهند، حتی اگر داده کهنه باشد؛ بعداً هماهنگ می‌شوند (eventual). دسترس‌پذیری را حفظ، سازگاری را موقتاً قربانی.
چرا یک سیستمِ توزیع‌شده نمی‌تواند CA باشد؟

چون در یک سیستمِ واقعاً توزیع‌شده، پارتیشنِ شبکه اجتناب‌ناپذیر است؛ دیر یا زود کابل قطع می‌شود یا یک نود از بقیه بریده می‌شود. «CA بودن» یعنی «فرض می‌کنیم هیچ‌وقت پارتیشن رخ نمی‌دهد» — که فرضِ باطلی است. وقتی پارتیشن رخ داد، مجبوری یکی را انتخاب کنی: یا به هر دو طرف اجازهٔ پاسخ می‌دهی (A را نگه می‌داری، C را از دست می‌دهی → AP)، یا طرفِ اقلیت را خاموش می‌کنی (C را نگه می‌داری، A را از دست می‌دهی → CP). تنها سیستمی که واقعاً CA است یک دیتابیسِ تک‌نودی است، چون اصلاً شبکه‌ای برای پارتیشن‌شدن ندارد. پس در طراحیِ توزیع‌شده، CA یک گزینه نیست؛ فقط CP یا AP.

PACELC: چیزی که CAP جا انداخت

قضیهٔ CAP یک نقصِ بزرگ دارد: فقط دربارهٔ زمانِ پارتیشن حرف می‌زند، که به‌ندرت اتفاق می‌افتد. ولی سیستم تو ۹۹.۹٪ وقت در حالتِ عادی است — و همان‌جا هم یک trade-off داری. Daniel Abadi در ۲۰۱۲ با PACELC این را کامل کرد:

فرمولِ PACELC

اگر پارتیشن (P) رخ دهد، بین Availability و Consistency انتخاب کن؛ در غیر این صورت (E) — یعنی حالتِ عادی — بین Latency و Consistency انتخاب کن. چون برای اینکه یک خواندن قطعاً جدیدترین داده را بدهد، باید با نسخه‌های دیگر هماهنگ کنی و این latency می‌آورد. پس همیشه، حتی بدون پارتیشن، بینِ «سریع» و «سازگار» در حال معامله‌ای.

سیستم دستهٔ PACELC یعنی
DynamoDB, Cassandra PA/EL زمانِ پارتیشن availability؛ در عادی latency (سریع ولی eventual)
CockroachDB, Spanner PC/EC همیشه سازگاری را مقدم می‌دارد، حتی به قیمتِ کندی
MongoDB (majority) PC/EC سازگاری‌محور
PostgreSQL (async replica) PC/EL روی primary سازگار؛ خواندن از replica سریع ولی ممکن است کهنه
فرق CAP و PACELC چیست و چرا PACELC کامل‌تر است؟

CAP فقط می‌گوید هنگامِ پارتیشن بین C و A انتخاب کن. مشکلش این است که پارتیشن نادر است، پس CAP دربارهٔ رفتارِ سیستم در ۹۹٪ زمان — یعنی حالتِ عادی — ساکت است. PACELC این خلأ را پر می‌کند: می‌گوید حتی بدونِ پارتیشن هم یک trade-offِ همیشگی بین latency و consistency داری، چون سازگاریِ قوی نیازمندِ هماهنگیِ همگام بین نسخه‌هاست که ذاتاً کُند است. مثالِ کلیدی: DynamoDB در دستهٔ PA/EL است — یعنی هم زمانِ پارتیشن و هم در حالتِ عادی، سرعت را به سازگاری ترجیح می‌دهد؛ در مقابل، Spanner در دستهٔ PC/EC است و همیشه سازگاری را انتخاب می‌کند (به کمکِ ساعتِ اتمیِ TrueTime برای کم‌کردنِ هزینه‌اش). در مصاحبه، اشاره به PACELC نشان می‌دهد که فهمیده‌ای trade-off فقط زمانِ بحران نیست، همیشگی است.


۷) طیفِ سازگاری: از قوی تا ضعیف

«سازگاری» یک کلمهٔ مبهم است؛ در واقع یک طیف است. از قوی‌ترین (گران‌ترین) به ضعیف‌ترین:

  • Linearizable (atomic): قوی‌ترین. سیستم طوری رفتار می‌کند که انگار یک نسخهٔ واحد از داده وجود دارد و هر عملیات در یک لحظهٔ اتمی بین شروع و پایانش رخ می‌دهد. اگر نوشتنِ من تمام شد، هر خواندنِ بعدی (به هر ساعتی) قطعاً آن را می‌بیند. این چیزی است که یک متغیرِ لوکال به تو می‌دهد. گران‌ترین است چون نیاز به هماهنگیِ سراسری دارد.
  • Sequential: همه عملیات‌ها را در یک ترتیبِ کلیِ واحد می‌بینند، و ترتیبِ عملیات‌های هر client حفظ می‌شود — ولی این ترتیب لازم نیست با زمانِ واقعیِ ساعتِ دیوار بخواند.
  • Causal: فقط عملیات‌هایی که رابطهٔ علّی دارند (happens-before) به ترتیب دیده می‌شوند؛ عملیات‌های concurrent می‌توانند به هر ترتیبی دیده شوند. این «شیرین‌ترین نقطه» است: قوی‌ترین سازگاری‌ای که بدونِ قربانی‌کردنِ availability قابل‌دستیابی است.
  • Read-your-writes: تضمین می‌کند بعد از اینکه من چیزی نوشتم، خودم آن را می‌بینم (هرچند دیگران شاید هنوز نه). خیلی از باگ‌های UX از نبودِ همین می‌آید.
  • Monotonic reads: اگر یک بار مقداری را دیدم، خواندن‌های بعدی‌ام عقب‌تر نمی‌روند.
  • Eventual: ضعیف‌ترین. فقط قول می‌دهد «اگر نوشتن‌ها متوقف شوند، بالاخره همه نسخه‌ها یکی می‌شوند». دربارهٔ «کِی» هیچ نمی‌گوید.

نمودار زیر طیفِ سازگاری را از قوی به ضعیف مرتب می‌کند. | This diagram orders the consistency spectrum from strong to weak.

flowchart TD
    L[Linearizable\nstrongest, costliest] --> S[Sequential]
    S --> C[Causal\nbest without losing availability]
    C --> RYW[Read-your-writes]
    RYW --> MR[Monotonic reads]
    MR --> E[Eventual\nweakest, cheapest]
read-your-writes و تلهٔ read replica

کلاسیک‌ترین باگِ تولید: کاربر پروفایلش را ذخیره می‌کند (نوشتن روی primary)، صفحه رفرش می‌شود (خواندن از یک read replica که هنوز sync نشده)، و کاربر تغییرش را نمی‌بیند و فکر می‌کند سیستم خراب است. راه‌حل‌ها: بعد از نوشتن، برای مدتِ کوتاهی خواندن‌های همان کاربر را به primary بفرست (sticky/read-from-primary)؛ یا از خواندنِ causal با نگه‌داشتنِ آخرین versionِ دیده‌شده استفاده کن. این دقیقاً همان جایی است که فهمِ عمیقِ سازگاری به یک تجربهٔ کاربریِ بهتر ترجمه می‌شود.

سازگاریِ قوی را فقط جایی بخر که لازم است

سازگاریِ linearizable گران است و throughput را پایین می‌آورد. سنیورها آن را per-use-case انتخاب می‌کنند: موجودیِ انبار و کیفِ پول → قوی؛ تعدادِ لایکِ یک پست، شمارِ بازدید، feed → eventual کاملاً کافی است. یک اشتباهِ رایج، خرجِ سازگاریِ قوی برای داده‌ای است که هیچ‌کس متوجهِ چند ثانیه تأخیرش نمی‌شود.


سازگاریِ causal چیست و چرا به آن «شیرین‌ترین نقطه» می‌گویند؟

سازگاریِ causal تضمین می‌کند عملیات‌هایی که رابطهٔ علّی دارند (happens-before) برای همه به همان ترتیب دیده شوند؛ مثلاً اگر پستی گذاشتی و بعد کامنتی روی آن، هیچ‌کس نباید کامنت را قبل از پست ببیند. اما عملیات‌های concurrent (بدونِ رابطهٔ علّی) می‌توانند به هر ترتیبی دیده شوند. چرا شیرین‌ترین نقطه؟ چون قضیهٔ CAP ثابت می‌کند سازگاریِ قوی (linearizable) با availability هنگامِ پارتیشن ناسازگار است — ولی causal قوی‌ترین مدلی است که می‌توان هم‌زمان با availability و بدونِ بلاک‌شدن هنگامِ پارتیشن به آن رسید. برای اکثرِ اپلیکیشن‌ها (چت، شبکهٔ اجتماعی، کامنت) causal دقیقاً همان چیزی است که کاربر انتظارِ منطقی‌اش را دارد، بدونِ پرداختِ هزینهٔ گرانِ linearizable.


۸) تحویلِ پیام: چرا «exactly-once» یک افسانه است

وقتی A پیامی به B می‌فرستد و ممکن است پیام گم شود، سه گارانتیِ ممکن داری:

  • At-most-once: بفرست و فراموش کن؛ اگر گم شد، گم شد. هیچ duplicate نداری ولی ممکن است پیام از دست برود. مناسبِ متریک و telemetry.
  • At-least-once: تا وقتی ack نگرفتی retry کن. هیچ‌وقت پیام گم نمی‌شود، ولی قطعاً روزی duplicate خواهی داشت (چون ممکن است ackِ برگشتی گم شود و تو دوباره بفرستی). این پیش‌فرضِ واقع‌بینانهٔ اکثر سیستم‌هاست.
  • Exactly-once: آرزوی همه؛ هر پیام دقیقاً یک بار اثر بگذارد.
exactly-once در واقع effectively-once است

تحویلِ فیزیکیِ exactly-once روی یک شبکهٔ نامطمئن غیرممکن است (به‌خاطرِ همان مسئلهٔ دو ژنرال). کاری که سیستم‌های واقعی می‌کنند این است: تحویل را at-least-once نگه می‌دارند (پس ممکن است پیام چند بار برسد)، ولی پردازش را idempotent می‌کنند یا duplicateها را dedup می‌کنند، طوری که اثرِ نهایی انگار یک بار بوده. به این می‌گویند effectively-once. کلید همیشه idempotency است، نه جادو.

idempotency: سلاحِ اصلی

عملیاتِ idempotent یعنی اجرای دوباره‌اش همان نتیجهٔ اجرای یک‌بار را می‌دهد. balance = 100 (set) idempotent است؛ balance += 10 (increment) نیست. الگوی استاندارد این است که به هر پیام یک کلیدِ یکتا (idempotency key) بدهی و قبل از پردازش چک کنی که آیا قبلاً دیده‌ای:

@Service
public class PaymentConsumer {

    private final ProcessedMessageRepository processed;
    private final PaymentService payments;

    @KafkaListener(topics = "payments")
    @Transactional
    public void onMessage(PaymentEvent event) {
        // dedup: اگر این messageId را قبلاً دیده‌ایم، رد شو
        if (processed.existsById(event.getMessageId())) {
            return; // duplicate — بی‌خطر نادیده بگیر
        }
        payments.charge(event.getAccountId(), event.getAmount());
        // ثبتِ dedup و کار در یک تراکنشِ اتمیک
        processed.save(new ProcessedMessage(event.getMessageId()));
    }
}

نکتهٔ حیاتی: ثبتِ کلیدِ dedup و اثرِ کسب‌وکار باید در یک تراکنشِ اتمیک باشند، وگرنه پنجره‌ای می‌ماند که کار انجام شده ولی dedup ثبت نشده (یا برعکس). این همان الگوی inbox است. جدولِ dedup چنین است — و upsertِ آن بین Oracle و PostgreSQL فرق دارد:

-- PostgreSQL: ON CONFLICT برای درجِ idempotent
INSERT INTO processed_message (message_id, processed_at)
VALUES (:id, now())
ON CONFLICT (message_id) DO NOTHING;
-- Oracle (19c/23ai): MERGE برای همان منطق
MERGE INTO processed_message t
USING (SELECT :id AS message_id FROM dual) s
ON (t.message_id = s.message_id)
WHEN NOT MATCHED THEN
  INSERT (message_id, processed_at) VALUES (s.message_id, SYSTIMESTAMP);
تفاوتِ گویش: ON CONFLICT در برابر MERGE

PostgreSQL از INSERT ... ON CONFLICT ... DO NOTHING/DO UPDATE استفاده می‌کند که خوانا و اتمیک است. Oracle تا مدت‌ها فقط MERGE داشت (که پرگوتر است و نیاز به FROM dual دارد)؛ از Oracle 23ai دستورِ INSERT ... ON CONFLICT مشابهِ PostgreSQL هم اضافه شده اما برای پرتابل بودن، MERGE امن‌ترین انتخابِ مشترک است. یادت باشد در Oracle رشتهٔ خالی '' معادلِ NULL است — پس روی ستون‌های idempotency key هرگز رشتهٔ خالی نگذار.

exactly-once در کافکا

کافکا از نسخهٔ 0.11 «exactly-once semantics» را با دو مکانیزم می‌سازد. اول idempotent producer: هر producer یک PID و به هر رکورد یک sequence number می‌دهد؛ broker با این دو، retryهای تکراری را روی لاگ حذف می‌کند. دوم transactions: نوشتن روی چند partition و commitِ offsetِ مصرف را اتمیک می‌کند.

پیش‌فرض‌های کافکا از 3.0

از کافکا 3.0 به بعد، producer به‌صورتِ پیش‌فرض enable.idempotence=true و acks=all است. برای transaction باید transactional.id ست کنی و در سمتِ consumer isolation.level=read_committed بگذاری تا پیام‌های تراکنشِ ناتمام را نخوانی. نکتهٔ مهم: exactly-once کافکا فقط درونِ خودِ کافکا (Kafka-to-Kafka، مثلِ Kafka Streams) واقعاً end-to-end است. به‌محضِ اینکه یک side-effectِ خارجی داری (نوشتن در دیتابیسِ دیگر، صداکردنِ یک API)، باز به idempotency در آن سمت نیاز داری. exactly-once کافکا تو را از idempotency بی‌نیاز نمی‌کند.

چطور در یک microservice پیام‌محور، exactly-once را تضمین می‌کنی؟

جوابِ صادقانه: exactly-once تحویل را تضمین نمی‌کنی، effectively-once پردازش را تضمین می‌کنی. سه لایه: (۱) producer را idempotent کن (در کافکا پیش‌فرض روشن است) تا retryها duplicate نسازند. (۲) در سمتِ مصرف، هر پیام یک کلیدِ یکتا داشته باشد و با یک جدولِ dedup (الگوی inbox) قبل از پردازش چک کنی؛ ثبتِ dedup و اثرِ کسب‌وکار در یک تراکنشِ دیتابیسیِ اتمیک باشند. (۳) هرجا ممکن است، خودِ عملیات را ذاتاً idempotent طراحی کن (set به‌جای increment، upsert به‌جای insert). اگر مصاحبه‌گر اصرار کرد «exactly-once واقعی می‌خواهم»، توضیح بده که این فیزیکاً روی شبکهٔ نامطمئن ممکن نیست و بهترین کاری که می‌شود کرد effectively-once است — و همین را همهٔ سیستم‌های صنعتی می‌کنند.


۹) تراکنشِ توزیع‌شده: 2PC در برابر saga

وقتی یک عملیات باید روی چند سرویس/دیتابیس اتمیک باشد (یا همه یا هیچ)، دو مکتب داری.

Two-Phase Commit (2PC)

یک coordinator دو فاز اجرا می‌کند. فاز ۱ (prepare): از همه می‌پرسد «آماده‌ای commit کنی؟» و هر participant قفل می‌گیرد و «yes/no» می‌گوید. فاز ۲ (commit): اگر همه yes گفتند، coordinator به همه commit می‌فرستد؛ وگرنه abort.

نمودار زیر یک 2PCِ موفق را نشان می‌دهد. | This diagram shows a successful two-phase commit.

sequenceDiagram
    participant C as Coordinator
    participant A as Service A
    participant B as Service B
    C->>A: prepare
    C->>B: prepare
    A-->>C: yes (locked)
    B-->>C: yes (locked)
    C->>A: commit
    C->>B: commit
    A-->>C: ack
    B-->>C: ack
2PC یک blocking protocol است — پاشنهٔ آشیلِ آن

اگر coordinator بعد از فاز prepare اما قبل از فرستادنِ commit بیفتد، participantها در حالتِ «prepared» گیر می‌کنند: قفل‌ها را گرفته‌اند، نمی‌دانند commit کنند یا abort، و باید تا برگشتِ coordinator صبر کنند. در این مدت آن ردیف‌های قفل‌شده برای بقیه بلاک است. به همین دلیل 2PC در معماریِ microservice تقریباً منسوخ است: throughput را می‌کُشد و availability را پایین می‌آورد. 3PC یک فاز (pre-commit) اضافه می‌کند تا blocking را کم کند ولی در برابرِ پارتیشنِ شبکه ناامن می‌شود و در عمل کسی استفاده‌اش نمی‌کند.

Saga: تراکنشِ توزیع‌شدهٔ عملی

الگوی saga تراکنشِ بزرگ را به دنباله‌ای از تراکنش‌های محلیِ کوچک می‌شکند که هرکدام روی یک سرویس commit می‌شوند. اگر مرحله‌ای شکست خورد، به‌جای rollback (که در سیستم توزیع‌شده ممکن نیست)، تراکنش‌های جبرانی (compensating transactions) را به‌ترتیبِ معکوس اجرا می‌کنی تا اثرِ مراحلِ قبلی را خنثی کنی.

دو سبک دارد: orchestration (یک هماهنگ‌کنندهٔ مرکزی مراحل را دستور می‌دهد؛ خوانا و قابل‌ردیابی) و choreography (هر سرویس با eventها به رویدادها واکنش می‌دهد؛ decoupled ولی دنبال‌کردنِ جریان سخت‌تر).

نمودار زیر یک saga سفارش با یک شکست و جبران را نشان می‌دهد. | This diagram shows an order saga with a failure and compensation.

sequenceDiagram
    participant O as Order
    participant P as Payment
    participant I as Inventory
    O->>P: reserve payment
    P-->>O: paid
    O->>I: reserve stock
    I-->>O: OUT OF STOCK (fail)
    Note over O,P: compensate in reverse
    O->>P: refund payment
    P-->>O: refunded
saga اتمیک نیست — isolation ندارد

تلهٔ بزرگِ saga: چون هر مرحله جداگانه commit می‌شود، در میانهٔ saga سیستم در حالتی است که دنیای بیرون می‌تواند آن را ببیند (پول کم شده ولی سفارش هنوز کامل نیست). این نبودِ isolation باعثِ آنومالی‌هایی مثل dirty read می‌شود. راه‌حل‌ها: استفاده از یک وضعیتِ «pending/reserved» به‌جای نهایی، قفل‌های semanticِ سطحِ کاربرد، و طراحیِ compensationها به‌صورت idempotent (چون ممکن است چند بار retry شوند). و همیشه compensationها را از قبل طراحی کن — «refund» گاهی از خودِ «charge» سخت‌تر است (پول رفته، کاربر رفته).

کِی 2PC و کِی saga؟

2PC سازگاریِ قوی و اتمیک می‌دهد ولی blocking است، قفل‌ها را طولانی نگه می‌دارد و availability را پایین می‌آورد؛ برای سیستم‌های با throughputِ بالا و توزیع‌شده نامناسب است. عملاً فقط جایی می‌ماند که participantها یک XA resource manager دارند و تعدادشان کم و شبکه‌شان قابل‌اعتماد است (مثلاً چند دیتابیس در یک datacenter). saga برای microserviceها استانداردِ عملی است: هر سرویس تراکنشِ محلیِ خودش را commit می‌کند و با compensation به شکست واکنش می‌دهد — availability بالا، بدونِ قفلِ توزیع‌شده، به قیمتِ از دست دادنِ atomicity و isolation. اگر مصاحبه‌گر پرسید «چطور atomicity را در saga جبران می‌کنی؟» بگو نمی‌کنی؛ eventual consistency را می‌پذیری و با state machine و compensation و طراحیِ idempotent، سیستم را در نهایت به یک حالتِ سازگار می‌رسانی.


۱۰) Gossip و anti-entropy: پخشِ اطلاعات مثلِ شایعه

وقتی صدها نود داری، نمی‌توانی هر تغییر را از یک نقطهٔ مرکزی به همه broadcast کنی (نقطهٔ شکستِ واحد و گلوگاه). به‌جایش از gossip (پروتکلِ اپیدمیک) استفاده می‌کنی: هر نود هر ثانیه با چند نودِ تصادفی حرف می‌زند و اطلاعاتش را رد و بدل می‌کند. مثلِ پخشِ شایعه، اطلاعات به‌صورتِ نمایی (exponential) و مقاوم به خرابی در کل cluster پخش می‌شود. Cassandra، Consul و Serf با gossip عضویتِ cluster و وضعیتِ سلامت را می‌فهمند.

gossip مثلِ پخشِ شایعه در دفتر

یک خبر را به دو نفر می‌گویی؛ هرکدام به دو نفرِ دیگر؛ و ظرفِ چند دقیقه کلِ دفتر می‌داند — بدونِ اینکه کسی به همه اعلانِ عمومی کند. اگر یک نفر هم غایب باشد، بالاخره از یکی دیگر می‌شنود. gossip همین است: مقاوم، غیرمتمرکز، و «بالاخره» همه را باخبر می‌کند (eventual). عیبش این است که «بالاخره» زمان می‌برد و در همان فاصله نودها نظرِ متفاوتی دارند.

anti-entropy فرایندِ پس‌زمینه‌ای است که اختلافِ داده بین نسخه‌ها را پیدا و ترمیم می‌کند. برای اینکه مقایسهٔ کلِ داده گران است، از درختِ Merkle استفاده می‌کنند: یک درختِ هش که برگ‌هایش هشِ بلوک‌های داده و گره‌های بالاتر هشِ فرزندانشان‌اند. دو نود فقط ریشه را مقایسه می‌کنند؛ اگر یکی بود، هیچ اختلافی نیست و تمام. اگر فرق داشت، درخت را پایین می‌روند و فقط شاخه‌هایی که هششان فرق دارد را همگام می‌کنند — به‌جای انتقالِ کلِ دیتاست، فقط تفاوت‌ها منتقل می‌شوند.


کِی eventual consistency قابل‌قبول است و کِی خطرناک؟

eventual consistency وقتی قابل‌قبول است که یک پنجرهٔ کوتاهِ ناسازگاری آسیبِ واقعی نزند و عملیات‌ها جابه‌جایی‌پذیر یا idempotent باشند. مثال‌های امن: شمارِ لایک و بازدید، feedِ شبکهٔ اجتماعی، آمار و متریک، cache، وضعیتِ «آخرین‌بار آنلاین» — چند ثانیه اختلاف بین نودها هیچ‌کس را ناراحت نمی‌کند. اما خطرناک است هرجا یک invariantِ سختِ کسب‌وکار روی داده هست: موجودیِ انبار (نباید منفی شود)، مانده‌حساب (نباید دو بار خرج شود)، رزروِ صندلی (نباید double-book شود). آنجا به سازگاریِ قوی یا حداقل به یک مکانیزمِ reservation/compensating نیاز داری. قاعدهٔ سنیور: بپرس «اگر دو کاربر هم‌زمان این را ببینند و هر دو اقدام کنند، بدترین اتفاق چیست؟» — اگر جواب «چیزِ مهمی نه» بود، eventual کافی است؛ اگر «پول یا موجودی خراب می‌شود»، نه.

۱۱) CRDT: داده‌ای که خودش conflict را حل می‌کند

فرض کن دو کاربر آفلاین یک سبدِ خرید را ویرایش می‌کنند و بعد هر دو sync می‌شوند. کدام برنده است؟ با last-write-wins یکی از تغییرات گم می‌شود. CRDT (Conflict-free Replicated Data Type) نوعِ خاصی از ساختارِ داده است که طوری طراحی شده که merge کردنِ دو نسخه همیشه یک نتیجهٔ قطعی و بدون‌conflict می‌دهد، بدونِ نیاز به هماهنگی. این پایهٔ اپلیکیشن‌های local-first و collaborative مثل Redis (نوعِ CRDT در Active-Active)، Riak و ویرایشگرهای اشتراکی است.

دو خانواده دارد. state-based (CvRDT): هر نود کلِ state خودش را می‌فرستد و یک تابعِ merge آن‌ها را ترکیب می‌کند؛ این merge باید جابه‌جایی‌پذیر (commutative)، شرکت‌پذیر (associative) و خودتوان (idempotent) باشد تا ترتیب و تکرارِ پیام‌ها مهم نباشد (یک join-semilattice). operation-based (CmRDT): فقط خودِ عملیات‌ها را پخش می‌کند که سبک‌ترند ولی تحویلِ دقیق‌تری می‌خواهند.

ساده‌ترین مثال، G-Counter (شمارندهٔ فقط-افزایشی) است: هر نود فقط خانهٔ خودش را زیاد می‌کند و merge یعنی maximum خانه‌به‌خانه:

// G-Counter: شمارندهٔ فقط-افزایشیِ توزیع‌شده، merge بدونِ conflict
public class GCounter {
    private final int nodeId;
    private final long[] counts;

    public GCounter(int nodeId, int numNodes) {
        this.nodeId = nodeId;
        this.counts = new long[numNodes];
    }

    public void increment() { counts[nodeId]++; }

    public long value() {
        long sum = 0;
        for (long c : counts) sum += c;
        return sum;
    }

    // merge: خانه‌به‌خانه max — commutative، associative، idempotent
    public void merge(GCounter other) {
        for (int i = 0; i < counts.length; i++) {
            counts[i] = Math.max(counts[i], other.counts[i]);
        }
    }
}

چون merge فقط max است، هر ترتیب و هر تکرارِ merge به یک جواب می‌رسد — دقیقاً خاصیتی که برای eventual consistency بدونِ هماهنگی می‌خواهی. برای شمارندهٔ کاهش‌پذیر PN-Counter (دو G-Counter، یکی برای + و یکی برای -)، برای مجموعه OR-Set (که add و remove را با تگِ یکتا مدیریت می‌کند) و برای مقدارِ تکی LWW-Register داری.

CRDT رایگان نیست — metadata رشد می‌کند

CRDTها جادو نیستند. قیمتشان metadata است: یک G-Counter به‌ازای هر نود یک خانه دارد، و یک OR-Set باید تگِ هر عنصرِ حذف‌شده (tombstone) را نگه دارد تا حذف‌ها درست merge شوند. در مقیاسِ بزرگ این metadata می‌تواند از خودِ داده بزرگ‌تر شود. CRDT وقتی درخشان است که merge خودکارِ بدونِ هماهنگی واقعاً برایت ارزش دارد (collaborative editing، شمارنده‌های آفلاین، multi-region active-active) — نه به‌عنوانِ راه‌حلِ پیش‌فرضِ هر مسئلهٔ سازگاری.


۱۲) Split-brain و fencing token: خطرناک‌ترین باگِ توزیع‌شده

split-brain یعنی به‌خاطرِ پارتیشنِ شبکه، دو نود هم‌زمان فکر می‌کنند leaderاند و هر دو شروع به نوشتن می‌کنند. نتیجه: داده‌ٔ خراب، دو primary که همدیگر را overwrite می‌کنند، پول دو بار خرج‌شده. این کابوسِ هر سیستمِ توزیع‌شده است.

lease و lock به‌تنهایی جلوی split-brain را نمی‌گیرند

سناریوی کلاسیک: نودِ A یک lock/lease می‌گیرد و leader می‌شود. بعد A دچارِ یک GC pause طولانی (مثلاً ۱۵ ثانیه) می‌شود. در این مدت lease منقضی می‌شود، سیستم فکر می‌کند A مُرده و به نودِ B lease می‌دهد. حالا A از GC برمی‌گردد، هنوز فکر می‌کند leader است، و یک نوشتهٔ کهنه به دیتابیس می‌فرستد. حالا دو leader داری. صرفِ داشتنِ lock کافی نیست چون A نمی‌داند lock‌اش را از دست داده.

راه‌حلِ درست fencing token است: هر بار که lock داده می‌شود، یک عددِ یکنواخت‌افزایشی (monotonic) هم صادر می‌شود. هر نوشته باید token‌اش را همراه ببرد، و منبعِ داده (دیتابیس/storage) هر نوشته‌ای با token کوچک‌تر از آخرین token دیده‌شده را رد می‌کند.

نمودار زیر نشان می‌دهد چطور fencing token جلوی نودِ زامبی را می‌گیرد. | This diagram shows how a fencing token stops a zombie node.

sequenceDiagram
    participant A as Node A (paused)
    participant L as Lock Service
    participant B as Node B
    participant S as Storage
    A->>L: acquire lock -> token 33
    Note over A: long GC pause...
    L->>B: lease expired, grant -> token 34
    B->>S: write with token 34 (accepted)
    A->>S: write with token 33 (STALE)
    S-->>A: REJECTED (33 < 34)
// storage فقط نوشته‌هایی با token بزرگ‌تر یا مساوی آخرین را می‌پذیرد
public class FencedStorage {
    private long lastToken = 0;

    public synchronized void write(long fencingToken, byte[] data) {
        if (fencingToken < lastToken) {
            throw new StaleTokenException(
                "token " + fencingToken + " < " + lastToken + " — rejected");
        }
        lastToken = fencingToken;
        persist(data);
    }
}

چون token یکنواخت‌افزایشی است، نودِ زامبیِ A با token قدیمی (۳۳) هرگز نمی‌تواند روی نوشتهٔ Bِ زندهٔ (۳۴) بنویسد. سیستم‌های واقعی: ZooKeeper با zxid، etcd با revision، و lockهای مبتنی بر version، همگی این عدد را فراهم می‌کنند.

split-brain چیست و چطور جلویش را می‌گیری؟

split-brain وقتی است که یک پارتیشنِ شبکه باعث می‌شود دو (یا چند) نود هم‌زمان خودشان را leader بدانند و هر دو بنویسند، که به داده‌ٔ خراب منجر می‌شود. دفاعِ اول: quorum — یک نود فقط وقتی می‌تواند leader بماند که با اکثریتِ cluster در تماس باشد؛ طرفِ اقلیتِ پارتیشن خودش را کنار می‌کشد (این چیزی است که Raft/etcd می‌کنند). ولی quorum به‌تنهایی جلوی نودِ «زامبی» را — نودی که lease‌اش منقضی شده ولی به‌خاطرِ GC pause یا شبکهٔ کند خودش خبر ندارد — نمی‌گیرد. برای آن به fencing token نیاز داری: یک عددِ یکنواخت‌افزایشی که با هر اعطای lock صادر می‌شود و در هر نوشته حمل می‌شود؛ منبعِ داده هر نوشته با token کهنه را رد می‌کند. ترکیبِ quorum برای انتخابِ درستِ leader و fencing token برای رد کردنِ نوشته‌های کهنه، دفاعِ کاملِ صنعتی است.

جمع‌بندیِ فصل

۱) شکست جزئی جوهرِ سختیِ سیستمِ توزیع‌شده است: حالتِ سومِ «نمی‌دانم» همه‌چیز را عوض می‌کند، و timeout هرگز به تو نمی‌گوید نود مُرده یا فقط کُند. ۲) به ساعتِ دیوار برای ترتیب اعتماد نکن؛ از happens-before، ساعتِ Lamport (فقط ترتیب) و vector clock (تشخیصِ همزمانی) استفاده کن. ۳) اجماع با FLP محدود است ولی Raft با timeoutِ تصادفی و اکثریت عملی‌اش می‌کند؛ همیشه تعدادِ فردِ نود. ۴) R+W>N مرزِ بین strong و eventual در مدلِ quorum است. ۵) CAP یعنی هنگامِ پارتیشن بین C و A انتخاب کن (CP یا AP، نه CA)؛ PACELC اضافه می‌کند که حتی در حالتِ عادی هم بین latency و consistency معامله می‌کنی. ۶) exactly-once یک افسانه است؛ چیزی که می‌سازی effectively-once است با at-least-once + idempotency + dedup. ۷) برای تراکنشِ توزیع‌شده saga را به 2PCِ blocking ترجیح بده، ولی بدان که isolation نداری. ۸) gossip/anti-entropy/CRDT ابزارهای eventual consistency بدونِ هماهنگی‌اند. ۹) در برابرِ split-brain، quorum برای انتخابِ leader و fencing token برای رد کردنِ نودِ زامبی را با هم به‌کار ببر. اگر این نُه اصل را در استخوانت حس کنی، در هر طراحی و هر مصاحبه‌ای مثلِ یک سنیورِ واقعی قضاوت می‌کنی.

Let me start with a hard truth that separates a mid-level engineer from a real senior: on a single machine, either everything works or everything fails. In a distributed system, something is always broken and you don't know which part. This is called partial failure, and nearly all of the difficulty in this field flows from that one sentence.

When you call save() on a local database, you get two outcomes: success or error. But when that request travels over the network to another service, a third, devilish outcome appears: you don't know. Maybe the request never arrived. Maybe it arrived, was processed, but the reply was lost. Maybe it's still processing and will answer in three seconds. That "I don't know" is the beating heart of distributed systems theory.

This chapter is not going to parrot academic formulas. After reading it, when someone in an interview or a design review asks "OK, so what happens here if the network partitions?", you should answer not with fear but with a precise mental map: knowing what is theoretically impossible, what trade-off you're making, and which pattern solves your pain.

Roadmap for this chapter
  1. Why distributed is hard: partial failure, an unreliable network, lying clocks, and the "eight fallacies."
  2. Failure models: crash, omission, timing, and Byzantine.
  3. Time and ordering: happens-before, Lamport clocks, and vector clocks.
  4. Consensus: the FLP result, Paxos at a glance, and Raft in detail (leader election and log replication).
  5. Replication and quorums: the R+W>N rule and read/write consistency.
  6. CAP and PACELC in depth.
  7. The consistency spectrum: linearizable through eventual and read-your-writes.
  8. Message delivery: at-least-once vs exactly-once (and why it's really effectively-once).
  9. Distributed transactions: 2PC/3PC vs saga.
  10. Gossip, anti-entropy, CRDTs, and finally split-brain and fencing tokens.

1) Why is distributed computing so hard?

The network lies to you

In 1994, Peter Deutsch and colleagues at Sun compiled a list called "The Eight Fallacies of Distributed Computing." These are false assumptions that every junior engineer unconsciously makes, and then pays for in production with blood and bone:

  1. The network is reliable. 2) Latency is zero. 3) Bandwidth is infinite. 4) The network is secure. 5) Topology doesn't change. 6) There is one administrator. 7) Transport cost is zero. 8) The network is homogeneous.
All eight are false — and each one is an incident

Every "only happens sometimes" bug in distributed systems is usually one of these eight assumptions that someone made unconsciously in code. For example, restTemplate.getForObject(url) with no timeout means "I assumed the network is reliable and latency is zero." The result? One slow service ties up all your Tomcat threads and the whole application hangs. That's a cascading failure.

The Two Generals problem

Two generals stand on opposite sides of a valley and must attack the enemy at the same time. The only way to communicate is to send a messenger through the valley, where the messenger might be captured. General one sends: "Attack at dawn." But how does he know it arrived? He waits for an acknowledgment. General two sends the ack, but how does he know his ack arrived? He must wait for an ack of the ack... This loop never ends. The "Two Generals problem" proves that over an unreliable channel, guaranteed agreement is theoretically impossible. Every time you write catch (TimeoutException) in code, you're wrestling with the ghost of those two generals: you don't know whether the other side did its job or not.

Why a timeout is not a definitive answer

A subtle point most people miss: a timeout never tells you whether the other side is dead or just slow. If you waited 5 seconds and got no reply, you have three indistinguishable possibilities: (a) the server died, (b) the server is alive but slow and will answer in 6 seconds, (c) the answer was ready but got lost on the way back. This uncertainty makes decisions hard: if you assume it's dead and resend, the work might be done twice (a duplicate). If you assume it's alive and wait, you might wait forever.


2) Failure models: know your enemy precisely

Before you build a "resilient" system, you must define precisely: resilient against what kind of failure? Failure models range from "gentle" to "evil":

Failure model What happens Real-world example Difficulty
Crash-stop (fail-stop) Node halts and never returns pod killed with OOMKilled Easiest
Crash-recovery Node crashes but later returns (maybe with stale state) service restart, lost RAM Medium
Omission Node drops some messages but is alive full TCP queue, dropped packets Medium
Timing (performance) Node responds but too late long GC pause, slow disk High
Byzantine (arbitrary) Node behaves arbitrarily/maliciously, sends contradictory messages memory-corruption bug, hacked node Hardest
Senior judgment: almost always target crash-recovery

Most industrial systems (Kafka, etcd, PostgreSQL, most databases) are designed for the crash-recovery with omission model, not Byzantine. Tolerating Byzantine faults (which blockchains need) is enormously expensive: instead of a simple majority ($n/2+1$), you need $3f+1$ nodes to tolerate $f$ traitors. In an internal enterprise system where you trust your own nodes, spending that cost is foolish. In an interview, if you say "our system is Byzantine fault tolerant," they'll immediately ask "why? your nodes are your own" — and if you don't have a good answer, you've shown weakness.

A crucial point: a timing failure is more dangerous than a crash failure, because a slow node is worse than a dead node. You remove a dead node from the loop and you're done. But a slow node shows green in the health check, still takes traffic, and slows everyone down. This phenomenon is called "gray failure" and is the cause of many real incidents.


3) Time and ordering: why you can't trust the clock

The wall clock lies

The biggest beginner mistake in a distributed system is thinking that System.currentTimeMillis() is comparable across two machines. It isn't.

Never measure elapsed time with currentTimeMillis

System.currentTimeMillis() is the wall clock and gets synchronized by NTP, meaning it can suddenly jump backward. If you take two timestamps from it and subtract, you might get a negative number! To measure a duration, always use System.nanoTime() (the monotonic clock that only moves forward). For business logic that asks "what time is it now," use Instant.now(), but never compare two Instants from two different machines to decide ordering.

// Wrong: with NTP this can go negative
long start = System.currentTimeMillis();
doWork();
long elapsed = System.currentTimeMillis() - start; // may be wrong!

// Right: the monotonic clock only moves forward
long startNanos = System.nanoTime();
doWork();
long elapsedNanos = System.nanoTime() - startNanos; // always valid

The problem is deeper: even with NTP, two servers' clocks usually differ by a few milliseconds (and in the worst case a few seconds) — this is clock skew. So the logic "the record with the newer timestamp wins" (last-write-wins) is inherently buggy: a write that really happened later might carry a smaller timestamp and be wrongly discarded.

happens-before: causal order instead of clock order

In his legendary 1978 paper, Leslie Lamport said: forget "what time" it happened. Just ask "could A have influenced B?". We call this relation "happens-before" and denote it $\rightarrow$:

  • If A and B are in the same process and A ran earlier: $A \rightarrow B$.
  • If A is sending a message and B is receiving that message: $A \rightarrow B$.
  • It's transitive: if $A \rightarrow B$ and $B \rightarrow C$ then $A \rightarrow C$.

If neither $A \rightarrow B$ nor $B \rightarrow A$, they are concurrent — meaning they have no causal relation and their order is meaningless.

The diagram below shows the happens-before relation: sending a message creates a causal edge. | این نمودار رابطهٔ happens-before را نشان می‌دهد: فرستادن پیام یک لبهٔ علّی می‌سازد.

sequenceDiagram
    participant P1 as Process 1
    participant P2 as Process 2
    P1->>P1: event a
    P1->>P2: message m (event b = send)
    Note over P2: event c = receive m
    P2->>P2: event d
    Note over P1,P2: a -> b -> c -> d (causal chain)

Lamport clock: a logical number per event

A Lamport clock is a simple counter on each process that translates this causal relation into a number. Its rule is three lines:

public class LamportClock {
    private long counter = 0;

    // before each local event: increment
    public synchronized long tick() {
        return ++counter;
    }

    // when sending a message: tick and send the number with the message
    public synchronized long onSend() {
        return ++counter;
    }

    // on receive: max of self and sender, then increment
    public synchronized long onReceive(long senderClock) {
        counter = Math.max(counter, senderClock) + 1;
        return counter;
    }
}

Lamport's guarantee: if $A \rightarrow B$ then $L(A) < L(B)$. But the reverse is not true: if $L(A) < L(B)$ it does not necessarily mean A caused B; they might just be concurrent. That's the fundamental limitation of Lamport clocks.

A Lamport clock is one-directional only

$A \rightarrow B \implies L(A) < L(B)$ is true. But $L(A) < L(B) \implies A \rightarrow B$ is false. That is, with a Lamport clock you can't tell whether two events are causal or merely concurrent. For that distinction you need a vector clock.

Vector clock: detecting concurrency

A vector clock fixes this weakness. Instead of one number, each node keeps an array: one slot per node. It increments its own slot for each event, and on receive it takes an element-wise maximum.

public class VectorClock {
    private final int nodeId;
    private final long[] vector;

    public VectorClock(int nodeId, int numNodes) {
        this.nodeId = nodeId;
        this.vector = new long[numNodes];
    }

    public synchronized long[] tick() {
        vector[nodeId]++;
        return vector.clone();
    }

    public synchronized void onReceive(long[] other) {
        for (int i = 0; i < vector.length; i++) {
            vector[i] = Math.max(vector[i], other[i]);
        }
        vector[nodeId]++;
    }

    // A is before B if every slot is <= and at least one is <
    public static boolean happensBefore(long[] a, long[] b) {
        boolean strictlyLess = false;
        for (int i = 0; i < a.length; i++) {
            if (a[i] > b[i]) return false;
            if (a[i] < b[i]) strictlyLess = true;
        }
        return strictlyLess;
    }
}

Now if neither $A \rightarrow B$ nor $B \rightarrow A$, you know with certainty they were concurrent — and that's exactly what an eventually consistent database like early Riak/DynamoDB needs to detect a conflict. When two concurrent writes hit the same key, the system keeps both as "siblings" and asks the application to resolve the conflict.

What's the difference between a Lamport clock and a vector clock, and when do you use each?

Both capture causal order (happens-before), but with different power. A Lamport clock is a scalar and only guarantees that if A caused B, its number is smaller — but it cannot detect concurrency. It's lightweight and great for building an arbitrary total order (e.g., to break ties). A vector clock is an array the length of the number of nodes and can tell precisely whether two events are causal or concurrent — essential for conflict detection in multi-master systems. Its cost is that its size grows with the number of nodes (a scaling problem). In practice: if you only need a total order, Lamport; if you must detect and resolve conflicts, a vector clock (or its compressed form, the dotted version vector).


4) Consensus: the hardest problem in the field

Consensus means several nodes agreeing on one value, even if some nodes fail. It's the foundation of everything: leader election, committing a transaction, recording log order. We want three properties: (1) Agreement: everyone agrees on one value. (2) Validity: the agreed value must be one of the proposed values. (3) Termination: a decision is eventually made.

The FLP result: deterministic consensus is impossible in an asynchronous network

In 1985, Fischer, Lynch, and Paterson proved that in a fully asynchronous system — where there is no bound on message delay — with even one faulty node, no deterministic algorithm can always reach consensus. So how do etcd and Raft work? Because in practice they cheat: they sidestep FLP using timeouts (an assumption of partial synchrony) and randomness (randomized timeouts). FLP doesn't say consensus is impossible; it says you can't be always safe AND always live at the same time. Real algorithms never sacrifice safety, but occasionally sacrifice liveness (a round might take longer).

Explain the FLP result in simple terms — does it mean consensus is impossible?

No, it's not impossible. FLP says that in a fully asynchronous system (with no bound on message delay), no deterministic algorithm can simultaneously guarantee all three of: safety, liveness, and tolerating even one failure. The intuition: when a node doesn't respond, you can't tell whether it's "dead" or "just slow"; if you wait for safety you might block forever (violating liveness), and if you proceed the node might later return with a different decision (violating safety). Real systems' escape hatch: assume partial synchrony via timeouts plus a bit of randomness. They never sacrifice safety and only temporarily lose liveness in the worst case (an election round takes longer). So FLP is a theoretical limit, not a practical dead end.

Paxos at a glance

Paxos (also Lamport's work) was the first proven consensus algorithm and sits at the heart of systems like Google Chubby and Spanner. It has three roles: Proposer, Acceptor, Learner. It works in two phases: phase 1 (Prepare/Promise), where a proposer obtains permission with a unique number, and phase 2 (Accept/Accepted), where the value gets locked in. A majority of acceptors suffices. Paxos is famous for being hard to understand; Lamport himself had to write "Paxos Made Simple," and even then people struggle to implement it. Multi-Paxos is its practical version for a sequence of decisions.

Raft: understandable consensus

Raft was created in 2014 by Ongaro and Ousterhout with an explicit goal: understandability. Today it's the consensus engine in etcd (and therefore Kubernetes), Consul, RabbitMQ (quorum queues), ClickHouse, and Kafka's KRaft mode. Raft decomposes the problem into three independent subproblems: leader election, log replication, and safety.

Each node is in one of three states and moves between them based on the term (a logical time counter that plays the role of a "presidential term").

The state diagram below shows the lifecycle of a Raft node. | نمودار حالت زیر چرخهٔ عمرِ یک نودِ Raft را نشان می‌دهد.

stateDiagram-v2
    [*] --> Follower
    Follower --> Candidate: election timeout\n(no heartbeat)
    Candidate --> Candidate: split vote\n(new term, retry)
    Candidate --> Leader: wins majority
    Candidate --> Follower: sees higher term\nor new leader
    Leader --> Follower: discovers higher term
    Leader --> Leader: send heartbeats

Leader election: each follower has a randomized election timeout (typically 150 to 300 ms). If it receives no heartbeat from the leader within that window, it assumes the leader died: it increments its own term, becomes a candidate, votes for itself, and requests votes (RequestVote) from the others. Each node votes only once per term. If a candidate gets a majority, it becomes leader and starts sending heartbeats. The randomized timeout is the key: it makes two nodes rarely become candidates simultaneously and reduces split votes — this is exactly where Raft sidesteps FLP.

Leader election like electing a meeting chair

Imagine in a meeting the chair (leader) says every few seconds "I'm here, carry on" (heartbeat). If you hear nothing for a while (election timeout), you assume the chair left. But instead of everyone shouting "I'll be chair!" at once (chaos), each person has a random timer; the first whose timer expires raises a hand and says "new round, I'm the candidate, vote for me." If a majority of the room agrees, that's the new chair. The term is the "round number of the meeting," so nobody can disrupt it with a stale command from a previous round.

Log replication: only the leader accepts writes from clients. It appends each command as an entry to its log and sends it in AppendEntries to followers. Once a majority of nodes have written the entry to disk, the leader marks it "committed," applies it to the state machine, and replies to the client. This "majority write before commit" guarantees that even if the leader immediately fails, any new leader (which must have a majority vote) definitely has that entry.

Why does Raft need an odd number of nodes (3 or 5)?

Because consensus is based on a majority (quorum = $n/2+1$), and an odd count gives the best "fault tolerance per cost" ratio. With 3 nodes, quorum is 2 and you can tolerate 1 failure. With 4 nodes, quorum is 3 — still only 1 failure tolerated, but with an extra machine that provides no benefit and only raises the chance of failure! With 5 nodes, quorum is 3 and you tolerate 2 failures. So always odd: 3 for most workloads, 5 for critical systems. An even number not only doesn't help but increases split-brain risk in a 50-50 partition.

In production, be sensitive to consensus latency

Every write in a Raft-based system must wait for an fsync to disk on a majority of nodes. So your write throughput is bounded by the slowest disk inside the quorum and the network round-trip. That's why etcd is placed in one datacenter (or nearby AZs), not spread across continents. If you spread a cluster's members across continents, every write takes hundreds of milliseconds. Seniors account for this in advance; juniors learn it after the incident.


5) Replication and quorums: the golden R+W>N rule

When you copy data across $N$ replicas, you choose two numbers yourself: for each write, you wait for acks from $W$ replicas; for each read, you ask $R$ replicas and take the newest. This Dynamo-style model (used by Cassandra, Riak, and ScyllaDB) lets you tune consistency per request.

The R + W > N rule

If $R + W > N$, the set of nodes you last wrote to and the set you're now reading from must share at least one node (the pigeonhole principle). That shared node has the newest value, so your read definitely sees the last write: this is called "strong consistency" in the quorum model. If $R + W \le N$, you might read data that hasn't yet reached those nodes — that is, eventual consistency.

For example with $N=3$: choosing $W=2, R=2$ is balanced ($2+2>3$). If you want fast reads, $W=3, R=1$ (slow writes, instant reads). If you want fast writes, $W=1, R=3$. And $W=1, R=1$ means maximum speed but no guarantee.

Config (N=3) Meaning Good for Risk
W=3, R=1 write to all, read from one read-heavy data write blocks if one node is down
W=2, R=2 balanced quorum default for most systems balanced
W=1, R=1 fastest, no guarantee metrics, logs, cache stale reads
A sloppy quorum breaks the guarantee

For higher availability, some systems use a "sloppy quorum": if the primary nodes aren't reachable, they place the write on substitute nodes (hinted handoff). This raises availability but breaks the $R+W>N$ guarantee: a read and a write might share no node and you'd read stale data. You must know this trade-off; Cassandra controls exactly this with consistency levels like QUORUM vs LOCAL_QUORUM.

Explain intuitively why R+W>N gives strong consistency.

Imagine N=3 replicas with W=2 and R=2. When you write, the data lands on at least 2 of the 3 nodes. When you read, you ask 2 of the 3 nodes. Question: could these two sets of 2 share no node at all? No — because 2+2=4 is greater than 3, by the pigeonhole principle at least one node must be in both sets. That shared node has the last write, so your read definitely sees the newest value (provided you use versioning to pick the newest among the responses). That's the whole magic of R+W>N: it guarantees overlap. If R+W≤N, that overlap isn't guaranteed and you might read from nodes that haven't seen the new write yet.

To compensate for inconsistency, these systems have two background mechanisms: read repair (when a read sees different versions, it pushes the newest back to the lagging nodes) and anti-entropy with a Merkle tree, which I'll explain in the gossip section.


6) CAP and PACELC: the fundamental decision

CAP: the forced choice at partition time

The CAP theorem (proposed by Eric Brewer, proven by Gilbert and Lynch) names three properties: Consistency (everyone sees one value), Availability (every request gets a response), Partition tolerance (the system works despite a network cut).

The industry's biggest misconception is thinking "pick 2 of 3." That's wrong. The truth is subtler:

You don't choose P — a partition is forced on you

In a real distributed system, a network partition is a physical reality, not an option; a cable is cut, a switch fails. So you must always tolerate P. Your real choice is only this: during a partition, do you sacrifice C or A? So CAP in practice means "CP or AP." If you say "our system is CA," you've said "our system offers no guarantee when the network cuts" — which is almost always wrong. A single-node database is CA, but the moment it's distributed you must choose between CP and AP.

  • CP (like etcd, ZooKeeper, HBase, MongoDB with majority): during a partition, the minority side refuses to respond so it doesn't serve contradictory data. It preserves consistency, sacrifices availability.
  • AP (like Cassandra, DynamoDB, Riak): both sides of the partition respond, even if the data is stale; they reconcile later (eventual). Preserves availability, temporarily sacrifices consistency.
Why can't a distributed system be CA?

Because in a truly distributed system, a network partition is inevitable; sooner or later a cable is cut or a node is severed from the rest. "Being CA" means "we assume a partition never happens" — which is a false assumption. When a partition does happen, you're forced to choose: either you let both sides respond (keep A, lose C → AP), or you shut down the minority side (keep C, lose A → CP). The only truly CA system is a single-node database, because it has no network to partition in the first place. So in distributed design, CA is not an option; only CP or AP.

PACELC: what CAP left out

The CAP theorem has a big flaw: it only talks about partition time, which rarely happens. But your system is in normal state 99.9% of the time — and there too you have a trade-off. Daniel Abadi completed the picture in 2012 with PACELC:

The PACELC formula

If a partition (P) occurs, choose between Availability and Consistency; else (E) — in normal operation — choose between Latency and Consistency. Because for a read to definitely return the newest data, you must coordinate with the other replicas, and that adds latency. So always, even without a partition, you're trading between "fast" and "consistent."

System PACELC class Meaning
DynamoDB, Cassandra PA/EL availability during partition; latency in normal (fast but eventual)
CockroachDB, Spanner PC/EC always prefers consistency, even at the cost of slowness
MongoDB (majority) PC/EC consistency-oriented
PostgreSQL (async replica) PC/EL consistent on primary; reads from replica fast but possibly stale
What's the difference between CAP and PACELC, and why is PACELC more complete?

CAP only says: during a partition, choose between C and A. Its problem is that partitions are rare, so CAP is silent about the system's behavior 99% of the time — the normal state. PACELC fills that gap: it says even without a partition you have a perpetual trade-off between latency and consistency, because strong consistency requires synchronous coordination between replicas, which is inherently slow. Key example: DynamoDB is PA/EL — it prefers speed over consistency both during a partition and in normal operation; by contrast, Spanner is PC/EC and always chooses consistency (using its TrueTime atomic clock to reduce the cost). In an interview, mentioning PACELC shows you understand the trade-off isn't just a crisis-time thing — it's perpetual.


7) The consistency spectrum: from strong to weak

"Consistency" is a vague word; it's really a spectrum. From strongest (most expensive) to weakest:

  • Linearizable (atomic): the strongest. The system behaves as if there is a single copy of the data and every operation takes effect at an atomic instant between its start and end. If my write completed, every subsequent read (by any clock) definitely sees it. This is what a local variable gives you. The most expensive, because it needs global coordination.
  • Sequential: everyone sees all operations in one single total order, and each client's operation order is preserved — but that order need not match real wall-clock time.
  • Causal: only causally-related operations (happens-before) are seen in order; concurrent operations may be seen in any order. This is the "sweet spot": the strongest consistency achievable without sacrificing availability.
  • Read-your-writes: guarantees that after I write something, I myself see it (though others may not yet). Many UX bugs come from lacking exactly this.
  • Monotonic reads: if I've seen a value once, my later reads won't go backward.
  • Eventual: the weakest. It only promises "if writes stop, all replicas eventually converge." It says nothing about "when."

The diagram below orders the consistency spectrum from strong to weak. | نمودار زیر طیفِ سازگاری را از قوی به ضعیف مرتب می‌کند.

flowchart TD
    L[Linearizable\nstrongest, costliest] --> S[Sequential]
    S --> C[Causal\nbest without losing availability]
    C --> RYW[Read-your-writes]
    RYW --> MR[Monotonic reads]
    MR --> E[Eventual\nweakest, cheapest]
Read-your-writes and the read-replica trap

The classic production bug: a user saves their profile (write to primary), the page refreshes (read from a read replica that hasn't synced yet), and the user doesn't see their change and thinks the system is broken. Fixes: after a write, route that user's reads to the primary for a short window (sticky/read-from-primary); or use causal reads by tracking the last version seen. This is exactly where a deep understanding of consistency translates into a better user experience.

Buy strong consistency only where you need it

Linearizable consistency is expensive and lowers throughput. Seniors choose it per-use-case: inventory and wallet balance → strong; the like-count of a post, view counters, feeds → eventual is entirely sufficient. A common mistake is spending on strong consistency for data whose few-seconds lag nobody would notice.


What is causal consistency and why is it called the "sweet spot"?

Causal consistency guarantees that causally-related operations (happens-before) are seen in the same order by everyone; for example, if you make a post and then a comment on it, no one should see the comment before the post. But concurrent operations (with no causal relation) may be seen in any order. Why the sweet spot? Because the CAP theorem proves that strong consistency (linearizable) is incompatible with availability during a partition — yet causal is the strongest model you can achieve while staying available and non-blocking during a partition. For most applications (chat, social networks, comments), causal is exactly what a user reasonably expects, without paying the expensive price of linearizable.


8) Message delivery: why "exactly-once" is a myth

When A sends a message to B and the message might be lost, you have three possible guarantees:

  • At-most-once: fire and forget; if it's lost, it's lost. You get no duplicates but might lose messages. Good for metrics and telemetry.
  • At-least-once: retry until you get an ack. A message is never lost, but you will definitely get duplicates someday (because the returning ack might be lost and you resend). This is the realistic default for most systems.
  • Exactly-once: everyone's wish; each message takes effect exactly once.
Exactly-once is really effectively-once

Physical exactly-once delivery over an unreliable network is impossible (because of the same Two Generals problem). What real systems do: they keep delivery at-least-once (so a message might arrive multiple times), but they make the processing idempotent or deduplicate the duplicates, so the net effect is as if it happened once. This is called effectively-once. The key is always idempotency, not magic.

Idempotency: the main weapon

An idempotent operation means running it again yields the same result as running it once. balance = 100 (set) is idempotent; balance += 10 (increment) is not. The standard pattern is to give each message a unique key (idempotency key) and check before processing whether you've seen it before:

@Service
public class PaymentConsumer {

    private final ProcessedMessageRepository processed;
    private final PaymentService payments;

    @KafkaListener(topics = "payments")
    @Transactional
    public void onMessage(PaymentEvent event) {
        // dedup: if we've seen this messageId before, skip
        if (processed.existsById(event.getMessageId())) {
            return; // duplicate — safely ignore
        }
        payments.charge(event.getAccountId(), event.getAmount());
        // record dedup and effect in one atomic transaction
        processed.save(new ProcessedMessage(event.getMessageId()));
    }
}

A crucial point: recording the dedup key and the business effect must be in one atomic transaction, otherwise there's a window where the work is done but dedup isn't recorded (or vice versa). This is the inbox pattern. The dedup table looks like this — and its upsert differs between Oracle and PostgreSQL:

-- PostgreSQL: ON CONFLICT for an idempotent insert
INSERT INTO processed_message (message_id, processed_at)
VALUES (:id, now())
ON CONFLICT (message_id) DO NOTHING;
-- Oracle (19c/23ai): MERGE for the same logic
MERGE INTO processed_message t
USING (SELECT :id AS message_id FROM dual) s
ON (t.message_id = s.message_id)
WHEN NOT MATCHED THEN
  INSERT (message_id, processed_at) VALUES (s.message_id, SYSTIMESTAMP);
Dialect difference: ON CONFLICT vs MERGE

PostgreSQL uses INSERT ... ON CONFLICT ... DO NOTHING/DO UPDATE, which is readable and atomic. Oracle long had only MERGE (which is more verbose and needs FROM dual); Oracle 23ai added a PostgreSQL-like INSERT ... ON CONFLICT, but for portability MERGE is the safest common choice. Remember that in Oracle an empty string '' equals NULL — so never store an empty string in an idempotency-key column.

Exactly-once in Kafka

Kafka builds "exactly-once semantics" from version 0.11 with two mechanisms. First, the idempotent producer: each producer gets a PID and each record a sequence number; using these two, the broker deduplicates repeated retries on the log. Second, transactions: it makes writing across multiple partitions and committing consumer offsets atomic.

Kafka defaults since 3.0

Since Kafka 3.0, the producer defaults to enable.idempotence=true and acks=all. For a transaction you must set transactional.id and, on the consumer side, set isolation.level=read_committed so you don't read messages from uncommitted transactions. Important: Kafka's exactly-once is truly end-to-end only within Kafka itself (Kafka-to-Kafka, like Kafka Streams). As soon as you have an external side-effect (writing to another database, calling an API), you still need idempotency on that side. Kafka's exactly-once does not free you from idempotency.

How do you guarantee exactly-once in a message-driven microservice?

The honest answer: you don't guarantee exactly-once delivery, you guarantee effectively-once processing. Three layers: (1) make the producer idempotent (on by default in Kafka) so retries don't create duplicates. (2) On the consuming side, give each message a unique key and check it against a dedup table (the inbox pattern) before processing; the dedup record and the business effect must be in one atomic database transaction. (3) Wherever possible, design the operation to be inherently idempotent (set instead of increment, upsert instead of insert). If the interviewer insists "I want real exactly-once," explain that this is physically impossible over an unreliable network and the best you can do is effectively-once — which is exactly what all industrial systems do.


9) Distributed transactions: 2PC vs saga

When one operation must be atomic across several services/databases (all or nothing), you have two schools.

Two-Phase Commit (2PC)

A coordinator runs two phases. Phase 1 (prepare): it asks everyone "are you ready to commit?" and each participant takes a lock and says "yes/no." Phase 2 (commit): if everyone said yes, the coordinator sends commit to all; otherwise abort.

The diagram below shows a successful 2PC. | نمودار زیر یک 2PCِ موفق را نشان می‌دهد.

sequenceDiagram
    participant C as Coordinator
    participant A as Service A
    participant B as Service B
    C->>A: prepare
    C->>B: prepare
    A-->>C: yes (locked)
    B-->>C: yes (locked)
    C->>A: commit
    C->>B: commit
    A-->>C: ack
    B-->>C: ack
2PC is a blocking protocol — its Achilles' heel

If the coordinator fails after the prepare phase but before sending commit, participants get stuck in the "prepared" state: they've taken locks, don't know whether to commit or abort, and must wait for the coordinator to return. During that time those locked rows are blocked for everyone else. That's why 2PC is nearly obsolete in microservice architectures: it kills throughput and lowers availability. 3PC adds a phase (pre-commit) to reduce blocking, but it becomes unsafe against network partitions and is practically unused.

Saga: the practical distributed transaction

The saga pattern breaks a large transaction into a sequence of small local transactions, each committed on one service. If a step fails, instead of a rollback (impossible in a distributed system), you run compensating transactions in reverse order to undo the effects of previous steps.

It has two styles: orchestration (a central coordinator commands the steps; readable and traceable) and choreography (each service reacts to events with events; decoupled but harder to trace the flow).

The diagram below shows an order saga with a failure and compensation. | نمودار زیر یک saga سفارش با یک شکست و جبران را نشان می‌دهد.

sequenceDiagram
    participant O as Order
    participant P as Payment
    participant I as Inventory
    O->>P: reserve payment
    P-->>O: paid
    O->>I: reserve stock
    I-->>O: OUT OF STOCK (fail)
    Note over O,P: compensate in reverse
    O->>P: refund payment
    P-->>O: refunded
A saga is not atomic — it has no isolation

The big saga trap: because each step commits separately, mid-saga the system is in a state the outside world can see (money deducted but the order not yet complete). This lack of isolation causes anomalies like dirty reads. Fixes: use a "pending/reserved" status instead of a final one, application-level semantic locks, and design compensations to be idempotent (since they might be retried multiple times). And always design compensations up front — a "refund" is sometimes harder than the "charge" itself (money gone, user gone).

When 2PC and when saga?

2PC gives strong, atomic consistency but is blocking, holds locks long, and lowers availability; it's unsuitable for high-throughput distributed systems. Practically it only remains where participants share an XA resource manager and are few in number with a reliable network (e.g., a couple of databases in one datacenter). Saga is the practical standard for microservices: each service commits its own local transaction and reacts to failure with compensation — high availability, no distributed lock, at the cost of losing atomicity and isolation. If the interviewer asks "how do you compensate for atomicity in a saga?", say you don't; you accept eventual consistency and, with a state machine, compensations, and idempotent design, you bring the system to a consistent state in the end.


10) Gossip and anti-entropy: spreading information like a rumor

When you have hundreds of nodes, you can't broadcast every change from one central point to everyone (single point of failure and bottleneck). Instead you use gossip (an epidemic protocol): every second each node talks to a few random nodes and exchanges its information. Like a spreading rumor, information propagates exponentially and fault-tolerantly across the whole cluster. Cassandra, Consul, and Serf learn cluster membership and health status via gossip.

Gossip like rumor spreading in an office

You tell a piece of news to two people; each tells two others; and within minutes the whole office knows — without anyone making a public announcement to everyone. If someone was absent, they eventually hear it from someone else. That's gossip: resilient, decentralized, and it "eventually" informs everyone (eventual). Its downside is that "eventually" takes time, and during that window nodes hold different opinions.

Anti-entropy is a background process that finds and repairs data differences between replicas. Since comparing all the data is expensive, they use a Merkle tree: a hash tree whose leaves are hashes of data blocks and whose higher nodes are hashes of their children. Two nodes compare only the root; if they match, there's no difference and it's done. If they differ, they walk down the tree and sync only the branches whose hashes differ — instead of transferring the whole dataset, only the differences move.


When is eventual consistency acceptable, and when is it dangerous?

Eventual consistency is acceptable when a short window of inconsistency causes no real harm and operations are commutative or idempotent. Safe examples: like and view counts, a social feed, stats and metrics, caches, "last seen online" — a few seconds of divergence between nodes bothers no one. But it's dangerous wherever a hard business invariant rides on the data: inventory stock (must not go negative), account balance (must not be spent twice), seat reservation (must not double-book). There you need strong consistency, or at least a reservation/compensating mechanism. The senior rule of thumb: ask "if two users see this at the same time and both act on it, what's the worst that happens?" — if the answer is "nothing important," eventual is fine; if it's "money or inventory gets corrupted," it isn't.

11) CRDT: data that resolves conflicts itself

Suppose two users edit a shopping cart offline, then both sync. Who wins? With last-write-wins, one of the changes is lost. A CRDT (Conflict-free Replicated Data Type) is a special kind of data structure designed so that merging two versions always gives a deterministic, conflict-free result without any coordination. It's the basis of local-first and collaborative apps like Redis (CRDT type in Active-Active), Riak, and shared editors.

It has two families. State-based (CvRDT): each node sends its entire state and a merge function combines them; this merge must be commutative, associative, and idempotent so that message order and repetition don't matter (a join-semilattice). Operation-based (CmRDT): it only broadcasts the operations themselves, which are lighter but require more precise delivery.

The simplest example is the G-Counter (grow-only counter): each node only increments its own slot, and merge is an element-wise maximum:

// G-Counter: distributed grow-only counter, conflict-free merge
public class GCounter {
    private final int nodeId;
    private final long[] counts;

    public GCounter(int nodeId, int numNodes) {
        this.nodeId = nodeId;
        this.counts = new long[numNodes];
    }

    public void increment() { counts[nodeId]++; }

    public long value() {
        long sum = 0;
        for (long c : counts) sum += c;
        return sum;
    }

    // merge: element-wise max — commutative, associative, idempotent
    public void merge(GCounter other) {
        for (int i = 0; i < counts.length; i++) {
            counts[i] = Math.max(counts[i], other.counts[i]);
        }
    }
}

Because merge is just max, any order and any repetition of merges reaches the same answer — exactly the property you want for eventual consistency without coordination. For a decrementable counter there's the PN-Counter (two G-Counters, one for + and one for -), for a set there's the OR-Set (which manages add and remove with a unique tag), and for a single value there's the LWW-Register.

CRDTs aren't free — metadata grows

CRDTs aren't magic. Their price is metadata: a G-Counter has one slot per node, and an OR-Set must keep a tag (tombstone) for each removed element so removals merge correctly. At large scale this metadata can grow bigger than the data itself. CRDTs shine when automatic, coordination-free merge is genuinely worth it (collaborative editing, offline counters, multi-region active-active) — not as the default solution to every consistency problem.


12) Split-brain and fencing tokens: the most dangerous distributed bug

Split-brain means that due to a network partition, two nodes simultaneously think they're the leader and both start writing. The result: corrupted data, two primaries overwriting each other, money spent twice. This is the nightmare of every distributed system.

A lease or lock alone doesn't prevent split-brain

The classic scenario: node A acquires a lock/lease and becomes leader. Then A suffers a long GC pause (e.g., 15 seconds). During that time the lease expires, the system thinks A is dead and grants the lease to node B. Now A returns from GC, still thinks it's the leader, and sends a stale write to the database. Now you have two leaders. Merely holding a lock isn't enough, because A doesn't know it lost its lock.

The correct fix is a fencing token: every time a lock is granted, a monotonically increasing number is issued too. Every write must carry its token, and the data source (database/storage) rejects any write whose token is smaller than the last token it saw.

The diagram below shows how a fencing token stops a zombie node. | نمودار زیر نشان می‌دهد چطور fencing token جلوی نودِ زامبی را می‌گیرد.

sequenceDiagram
    participant A as Node A (paused)
    participant L as Lock Service
    participant B as Node B
    participant S as Storage
    A->>L: acquire lock -> token 33
    Note over A: long GC pause...
    L->>B: lease expired, grant -> token 34
    B->>S: write with token 34 (accepted)
    A->>S: write with token 33 (STALE)
    S-->>A: REJECTED (33 < 34)
// storage accepts only writes with a token >= the last one
public class FencedStorage {
    private long lastToken = 0;

    public synchronized void write(long fencingToken, byte[] data) {
        if (fencingToken < lastToken) {
            throw new StaleTokenException(
                "token " + fencingToken + " < " + lastToken + " — rejected");
        }
        lastToken = fencingToken;
        persist(data);
    }
}

Because the token is monotonically increasing, zombie node A with the old token (33) can never overwrite the live node B's write (34). Real systems: ZooKeeper with zxid, etcd with revision, and version-based locks all provide this number.

What is split-brain and how do you prevent it?

Split-brain is when a network partition causes two (or more) nodes to consider themselves leader simultaneously and both write, leading to corrupted data. First defense: quorum — a node can only remain leader while in contact with a majority of the cluster; the minority side of the partition steps down (this is what Raft/etcd do). But quorum alone doesn't stop a "zombie" node — one whose lease has expired but which, due to a GC pause or slow network, doesn't know it. For that you need a fencing token: a monotonically increasing number issued with each lock grant and carried in every write; the data source rejects any write with a stale token. The combination of quorum for correct leader election and a fencing token for rejecting stale writes is the complete industrial defense.

Chapter recap
  1. Partial failure is the essence of distributed difficulty: the third state, "I don't know," changes everything, and a timeout never tells you whether a node is dead or just slow. 2) Don't trust the wall clock for ordering; use happens-before, Lamport clocks (order only), and vector clocks (concurrency detection). 3) Consensus is bounded by FLP but Raft makes it practical with randomized timeouts and majorities; always an odd number of nodes. 4) R+W>N is the boundary between strong and eventual in the quorum model. 5) CAP means during a partition choose between C and A (CP or AP, not CA); PACELC adds that even in normal operation you trade between latency and consistency. 6) Exactly-once is a myth; what you build is effectively-once with at-least-once + idempotency + dedup. 7) For distributed transactions, prefer saga over blocking 2PC, but know you have no isolation. 8) Gossip/anti-entropy/CRDT are tools for coordination-free eventual consistency. 9) Against split-brain, combine quorum for leader election and a fencing token to reject the zombie node. If you feel these nine principles in your bones, you'll reason like a real senior in any design and any interview.