Architecture & Design · معماری و طراحی سنیورSenior ~47 دقیقه مطالعه~40 min read

میکروسرویس، سیستم‌های توزیع‌شده و طراحی سیستمMicroservices, Distributed Systems & System Design

از صفر یاد می‌گیری چرا توزیع‌شدگی یک هزینه است نه هدف، و با تشبیه‌های ساده CAP/PACELC، سازگاری، idempotency، Saga/outbox، CQRS و الگوهای تاب‌آوری را می‌فهمی و یک چارچوب طراحی سیستم را روی ردیابی ناوگان پیاده می‌کنی.Learn from scratch why distribution is a cost rather than a goal, then build up CAP/PACELC, consistency, idempotency, Saga/outbox, CQRS and resilience patterns through plain analogies, and apply a reusable system-design framework to a real fleet-tracking problem.


خب، بیا با هم صادق باشیم: «میکروسرویس» یکی از آن واژه‌هایی است که همه در رزومه می‌نویسند و کمتر کسی می‌تواند در مصاحبه دقیق تعریفش کند. این فصل قرار است تو را از یک نفر که اسم این الگوها را شنیده، به کسی تبدیل کند که می‌داند هر کدام چه دردی را درمان می‌کنند، چه هزینه‌ای دارند و کِی نباید از آن‌ها استفاده کند. من هیچ اصطلاحی را بدون توضیح رها نمی‌کنم — هر جا واژهٔ تازه‌ای مثل idempotent یا partition یا linearizable ببینی، همان‌جا با یک تشبیه بازش می‌کنم.

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

اول یاد می‌گیریم چرا توزیع‌شدگی یک هزینه است نه یک هدف (و آن هشت مغالطهٔ معروف). بعد سراغ تجزیهٔ درست و دام «مونولیت توزیع‌شده» می‌رویم. سپس دو قضیهٔ بنیادین — CAP و PACELC — و مدل‌های سازگاری را دقیق بیان می‌کنیم. بعد ابزارهای عملی: idempotency و افسانهٔ exactly-once، تراکنش‌های توزیع‌شده با Saga و outbox، الگوی CQRS/event sourcing، و پنج الگوی تاب‌آوری (timeout، retry، circuit breaker، bulkhead، backpressure). در ادامه API gateway/BFF/service mesh و سه ستون observability. آخر سر یک چارچوب طراحی سیستم را روی یک مسئلهٔ واقعی (ردیابی یک میلیون خودرو) پیاده می‌کنیم و با ۱۵ سؤال مصاحبه جمع‌بندی می‌کنیم.

بخش ۰ — واژه‌هایی که باید از قبل بدانی

قبل از شروع، چند واژهٔ پایه را با تشبیه در ذهنت میخ‌کوب کنیم تا بعداً بدون مکث جلو برویم.

  • تراکنش اتمیک (atomic transaction): «همه یا هیچ». مثل انتقال پول بین دو حساب؛ یا هر دو طرف (کم شدن از یکی، اضافه شدن به دیگری) با هم انجام می‌شوند یا هیچ‌کدام. نیمه‌کاره ممنوع.
  • ACID: چهار تضمین یک دیتابیس سنتی — Atomicity (اتمیک بودن)، Consistency (سالم ماندن قواعد داده)، Isolation (تراکنش‌ها روی هم اثر نگذارند) و Durability (بعد از تأیید، داده گم نشود حتی با قطع برق).
  • hop (پرش): هر بار که درخواست از یک سرویس به سرویس دیگر روی شبکه می‌رود، یک hop است. مثل یک تماس تلفنی بین دو دفتر؛ هر تماس هم زمان می‌برد هم ممکن است قطع شود.
  • latency (تأخیر) در برابر throughput (توان عبوری): latency یعنی «یک درخواست چقدر طول می‌کشد» (زمان تحویل یک پیتزا)؛ throughput یعنی «در واحد زمان چند تا می‌توانی تحویل بدهی» (پیتزا در ساعت). این دو فرق دارند.
  • node (نود): یک ماشین/سرور در سیستم توزیع‌شده — یکی از شعبه‌های شرکت.
  • replica (نسخهٔ تکراری): یک کپی از داده روی نودی دیگر، برای اینکه اگر یکی خراب شد یا شلوغ شد، بقیه جواب بدهند.
  • partition (پارتیشن شبکه): وقتی ارتباط شبکه بین دو گروه از نودها قطع می‌شود ولی خودشان زنده‌اند — مثل قطع شدن خط تلفن بین دو شعبه در حالی که هر دو شعبه باز و مشغول کارند.
  • idempotent (خودتکرارپذیر): عملیاتی که اگر یک بار یا صد بار انجامش بدهی، نتیجه یکی است. دکمهٔ طبقهٔ ۳ آسانسور idempotent است: صد بار فشارش بدهی، باز هم می‌روی طبقهٔ ۳.
نگه‌دار این جمله را

هر چیزی که در مونولیت «رایگان» بود — تماس مطمئن، تراکنش اتمیک — به محض عبور از مرز شبکه پولی می‌شود. کل این فصل دربارهٔ این است که آن هزینه را چطور بپردازی یا از پرداختش طفره بروی.

توزیع‌شدگی یک هزینه است، نه یک هدف

یک آشپزخانه در برابر چند رستوران

تصور کن یک آشپزخانهٔ واحد داری. سرآشپز می‌تواند بگوید «این سفارش را همه با هم بزنید» و مطمئن باشد یا کل سفارش آماده می‌شود یا هیچ (تراکنش اتمیک)؛ صدا زدن آشپز کناری فوری و مطمئن است (تماس درون‌حافظه‌ای)؛ و فقط یک آشپزخانه هست که باید حواست بهش باشد. حالا تصمیم می‌گیری کارها را به پنج رستوران در پنج نقطهٔ شهر بسپاری. حالا هر رستوران می‌تواند جدا رشد کند و جدا تعطیل/باز شود — اما دیگر آن «همه با هم یا هیچ» رایگان نیست، و هر «آشپز کناری» حالا سر آن طرف شهر است و باید بهش زنگ بزنی و امیدوار باشی جواب بدهد.

یک مونولیت که در یک پروسه جا می‌شود، سه امتیاز را رایگان می‌دهد: تراکنش‌های اتمیک، فراخوانی‌های همگام درون‌حافظه‌ای (تابع را صدا می‌زنی، همان‌جا جواب می‌گیری)، و یک آرتیفکت قابل‌استقرار (deployable) واحد که راحت می‌شود درباره‌اش استدلال کرد.

هر مرز شبکه‌ای که اضافه می‌کنی، این سه را می‌فروشی تا در عوض مقیاس‌پذیری مستقل و استقرار مستقل بخری. پس پرسش حاکم بر کل این فصل این است:

آیا سود سازمانی و مقیاسیِ این جداسازی، ارزشِ از دست دادن تراکنش رایگانِ ACID و فراخوانی تابعِ مطمئن را دارد؟

اگر نتوانی سودِ یک جداسازی مشخص را با جمله بیان کنی، جدا نکن.

هشت مغالطهٔ محاسبات توزیع‌شده

این هشت باور غلط (معروف به fallacies of distributed computing، از Deutsch و Gosling) چک‌لیستی است که هر طرح باید رعایتش کند. هر کدام یک فرضِ راحت‌طلبانه است که در دنیای واقعی نقض می‌شود:

  1. شبکه قابل‌اتکا نیست (بسته گم می‌شود).
  2. تأخیر صفر نیست (رفت‌وبرگشت زمان می‌برد).
  3. پهنای باند بی‌نهایت نیست.
  4. شبکه امن نیست.
  5. توپولوژی تغییر می‌کند (نودها می‌آیند و می‌روند).
  6. یک ادمین واحد وجود ندارد.
  7. هزینهٔ انتقال صفر نیست.
  8. شبکه همگن نیست (ماشین‌ها و نسخه‌ها فرق دارند).
تقریباً هر incident به یکی از این هشت برمی‌گردد

وقتی یک سیستم توزیع‌شده در تولید (production) می‌ترکد، تقریباً همیشه یک نفر جایی یکی از این هشت را «فرض گرفته». «فرض کردیم شبکه سریع است»، «فرض کردیم پیام حتماً می‌رسد». این لیست را مثل یک چک‌لیست ایمنی هوانوردی نگه‌دار.

تجزیه (decomposition) و دام مونولیت توزیع‌شده

دپارتمان‌ها را بر اساس «کار»، نه «عنوان شغلی» بچین

یک شرکت را تصور کن. تقسیم درست این است: دپارتمان «فروش»، دپارتمان «انبار»، دپارتمان «حسابداری» — هر کدام یک قابلیت کامل کسب‌وکار را از اول تا آخر انجام می‌دهد و مالک کار خودش است. تقسیم غلط این است: دپارتمان «کسانی که تلفن جواب می‌دهند»، دپارتمان «کسانی که فرم پر می‌کنند»، دپارتمان «کسانی که بایگانی می‌کنند» — این‌ها لایه‌های فنی‌اند، نه قابلیت. برای انجام یک کار ساده باید بین هر سه دست‌به‌دست شود، بدون آنکه هیچ‌کدام واقعاً مستقل باشند.

مرزهای خوبِ سرویس از قابلیت‌های کسب‌وکار (business capabilities) و کانتکست‌های محدود (bounded contexts) در طراحی دامنه‌محور (DDD) پیروی می‌کنند، نه از لایه‌های فنی. «سرویس سفارش، سرویس پرداخت، سرویس موجودی» یک جداسازیِ قابلیتی است. «سرویس کنترلر، سرویس ریپازیتوری، سرویس دیتابیس» یک جداسازیِ لایه‌ای است که فقط hop شبکه اضافه می‌کند، بدون هیچ خودمختاری.

bounded context یعنی چه؟ مرزی که در آن یک واژه یک معنای دقیق دارد. مثلاً «مشتری» در دپارتمان فروش با «مشتری» در دپارتمان پشتیبانی دو مدل متفاوت است. هر سرویس داخل مرز خودش زبان و مدل خودش را دارد و بیرون را قاطی نمی‌کند.

حالتِ شکستی که باید از آن بترسی مونولیت توزیع‌شده (distributed monolith) است: بدترین هر دو دنیا — پیچیدگیِ توزیع‌شدگی، بدون استقلالش. نشانه‌ها:

  • دیتابیس مشترک. دو سرویس که روی جدول‌های یکسان می‌نویسند، یعنی دیگر نمی‌توانی schema را تغییر بدهی، مستقل مستقر شوی یا مستقل مقیاس بدهی. هر تغییر باید با هماهنگیِ همه انجام شود. این بزرگ‌ترین ضدالگو است.
  • زنجیره‌های همگام پرحرف. اگر پاسخ به یک درخواست یعنی A → B → C → D به‌صورت همگام، آنگاه دسترس‌پذیریِ کلِ تو حاصل‌ضرب دسترس‌پذیری هر hop است، و تأخیرت مجموع آن‌ها. مثال عددی: اگر هر سرویس ۹۹٪ در دسترس باشد، ۰٫۹۹⁴ ≈ ۰٫۹۶ — یعنی زنجیرهٔ چهارتایی از هر کدام بی‌اعتمادتر است.
  • انتشار قفل‌شده (lockstep). اگر نتوانی سرویس X را بدون استقرارِ همزمانِ Y مستقر کنی، این‌ها در واقع یک سرویس‌اند با دو نقاب.
دیتابیس مشترک، ریشهٔ بیشتر بدبختی‌ها

قانون طلایی: هر سرویس مالک دادهٔ خودش است. سرویس‌های دیگر فقط از طریق API یا رویدادهای منتشرشده‌اش به آن داده می‌رسند، نه با خواندن مستقیم جدولش. لحظه‌ای که دو تیم به یک جدول write می‌کنند، دیگر دو سرویس نداری؛ یک مونولیت داری که تظاهر می‌کند دو تاست.

قواعد سرانگشتی برای تجزیه: درشت‌دانه شروع کن (یک «مونولیت ماژولار» — یک deployable ولی با ماژول‌های تمیز و مرزبندی‌شدهٔ داخلی) و سرویس را فقط زمانی استخراج کن که یک ماژول یکی از این سه را داشته باشد: پروفایل مقیاسیِ متمایز، تیم متمایز، یا نرخ تغییرِ متمایز.

قانون کانوی اختیاری نیست

قانون کانوی (Conway's Law): معماریِ سیستمِ تو آینهٔ چارت سازمانی‌ات خواهد شد. تیم‌هایی که ساختار سیستم را طراحی می‌کنند، ناخودآگاه ساختار ارتباطیِ خودشان را در آن کپی می‌کنند. نتیجهٔ عملی: اگر مرزهای سرویس با مرزهای تیم هم‌راستا نباشد، دائم با اصطکاک می‌جنگی. سازمان و مرزها را با هم طراحی کن.

CAP و PACELC — قضیه را دقیق بیان کن

خط تلفن بین دو شعبه قطع شده

دو شعبهٔ یک بانک داری که باید موجودیِ یک حساب را همگام نگه دارند و مدام به هم زنگ می‌زنند. یک روز خط تلفن بینشان قطع می‌شود (partition) — اما مشتری‌ها هنوز جلوی هر دو باجه صف کشیده‌اند. حالا دو انتخاب داری: (الف) هر دو شعبه به کار ادامه بدهند و برداشت انجام بدهند — ولی ریسک اینکه یک حساب دو بار خالی شود؛ یا (ب) شعبه‌ای که مطمئن نیست بگوید «ببخشید الان نمی‌توانم سرویس بدهم» و در را ببندد تا خط وصل شود. نمی‌توانی هم‌زمان هم «همیشه سرویس بده» و هم «هرگز داده خراب نشو» باشی وقتی خط قطع است. این دقیقاً CAP است.

قضیهٔ CAP: در طول یک پارتیشن شبکه (P)، یک سیستم توزیع‌شده باید بین این دو یکی را انتخاب کند:

  • سازگاری (Consistency, C): هر خواندن آخرین نوشتنِ تأییدشده را می‌بیند. این خاصیت را linearizability (خطی‌پذیری) می‌گویند — یعنی انگار تنها یک نسخه از داده وجود دارد و همه به همان ترتیبِ زمانِ واقعی می‌بینندش.
  • دسترس‌پذیری (Availability, A): هر درخواست یک پاسخِ بدونِ‌خطا می‌گیرد (نه لزوماً جدیدترین داده، ولی خطا نمی‌دهد).

حین پارتیشن نمی‌توانی هر دو را داشته باشی. اشتباهِ رایج این است که CAP را «دو تا از سه تا را بردار» تصور کنیم. اما پارتیشن‌ها اختیاری نیستند؛ در دنیای واقعی رخ می‌دهند. پس انتخابِ واقعی این است: CP در برابر AP، هنگام وقوعِ پارتیشن.

  • CP (مثلاً یک store با سازگاریِ قوی مثل ZooKeeper/etcd، یا یک RDBMS تک‌رهبر): هنگام پارتیشن، سمتِ اقلیت از نوشتن سر باز می‌زند تا داده واگرا نشود. دسترس‌پذیری را از دست می‌دهی، درستی را نگه می‌داری.
  • AP (مثل store‌های سبک Dynamo، یا Cassandra با سازگاریِ پایین): هنگام پارتیشن، هر دو سمت به سرویس‌دهی ادامه می‌دهند و بعداً آشتی می‌دهی (reconcile). دسترس‌پذیری را نگه می‌داری، ریسکِ خواندنِ کهنه یا متضاد را می‌پذیری.

PACELC — عدسیِ کامل‌تر

CAP فقط دربارهٔ لحظهٔ پارتیشن حرف می‌زند، اما پارتیشن نادر است. PACELC (از Daniel Abadi) تصویر را کامل می‌کند:

اگر Partition بود، بین A یا C انتخاب کن؛ در غیر این صورت (Else، کارکردِ عادی) بین Latency یا Consistency.

این عدسیِ روزمرهٔ مفیدتری است، چون پارتیشن‌ها نادرند اما تعاملِ تأخیر در برابر سازگاری همیشه حاضر است. برای اینکه یک خواندن آخرین نوشتن را قطعی ببیند، باید منتظرِ هماهنگیِ چند replica بمانی — و آن انتظار یعنی تأخیر.

  • Cassandra به‌صورت پیش‌فرض PA/EL است: هنگام پارتیشن به‌نفعِ دسترس‌پذیری، و در حالت عادی به‌نفعِ تأخیرِ کم (سرعت).
  • یک RDBMS سنتی معمولاً PC/EC است: به‌نفعِ سازگاری در هر دو حالت.
  • اگر Cassandra را با خواندن+نوشتنِ QUORUM تنظیم کنی (یعنی اکثریتِ replicaها باید تأیید کنند)، به سمتِ EC حرکت می‌کند، به بهای تأخیرِ بیشتر.
این جمله را در مصاحبه بگو

«چرا این خواندن کهنه است؟» اغلب جواب یک باگ نیست، بلکه یک انتخابِ عمدیِ EL است — سیستم عمداً سرعت را به تازگیِ داده ترجیح داده. دانستنِ PACELC یعنی می‌توانی رفتارِ سیستم را توضیح بدهی به‌جای اینکه دنبال باگِ خیالی بگردی.

مدل‌های سازگاری (consistency models)

سازگاری یک دکمهٔ روشن/خاموش نیست؛ یک طیف است، از قوی‌ترین (گران‌ترین) تا ضعیف‌ترین (ارزان‌ترین و مقیاس‌پذیرترین):

مدل تضمین کاربرد معمول
خطی‌پذیر / قوی (linearizable) خواندن آخرین نوشتنِ کامیت‌شده را با ترتیبِ زمانِ واقعی می‌بیند انتخاب رهبر، قفل، موجودی حساب
ترتیبی (sequential) همهٔ نودها عملیات را با ترتیبِ یکسان می‌بینند (نه لزوماً زمانِ واقعی) برخی لاگ‌های replicated
علّی (causal) عملیاتِ علّی‌مرتبط به ترتیب دیده می‌شوند؛ همزمان‌ها ممکن است جابه‌جا شوند اپ‌های مشارکتی، کامنت‌ها
خواندن-نوشته‌های-خود (read-your-writes) نوشته‌های خودت را می‌بینی تجربهٔ سطح session (پست بگذار، رفرش کن، ببینش)
نهایی (eventual) در نبودِ نوشتنِ جدید، replicaها هم‌گرا می‌شوند کش، DNS، کاتالوگ

بیشتر سیستم‌های بزرگ سازگارِ نهایی (eventually consistent) هستند، با سازگاریِ قویِ هدف‌مند فقط آن‌جا که پول یا ایمنی طلب می‌کند. «قوی در همه‌جا» مقیاس نمی‌گیرد (خیلی گران است)؛ «نهایی در همه‌جا» ثابت‌ها (invariants — قواعدی که همیشه باید درست بمانند، مثل «موجودی منفی نشود») را خراب می‌کند. مهارتِ سنیور، انتخابِ سازگاری به‌ازای هر عملیات است، نه یک تصمیمِ سراسری.

Idempotency و افسانهٔ «دقیقاً یک‌بار (exactly-once)»

بسته‌ای که رسیدش گم شده

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

تحویلِ دقیقاً-یک‌بار روی شبکهٔ نامطمئن وجود ندارد. فرستنده‌ای که ack نمی‌گیرد نمی‌داند پیام گم شده یا ack گم شده، پس مجبور است retry کند — یعنی گیرنده گاهی تکراری می‌بیند. آنچه می‌توانی بسازی پردازشِ effectively-once (عملاً یک‌بار) است:

پردازشِ effectively-once = تحویلِ at-least-once (حداقل یک‌بار) + هندلرِ idempotent.

این تمایز یکی از محبوب‌ترین تله‌های مصاحبه است. «تحویل» را نمی‌توانی یک‌بار کنی؛ «اثرِ پردازش» را می‌توانی.

idempotency یعنی اعمالِ یک عملیات N≥۱ بار همان اثری را دارد که یک‌بار اعمالش. سه تکنیک:

  1. کلیدهای idempotency. کلاینت به‌ازای هر عملیاتِ منطقی یک کلیدِ یکتا می‌سازد؛ سرور کلید + نتیجه را ذخیره می‌کند. در retry با همان کلید، به‌جای اجرای دوباره، نتیجهٔ ذخیره‌شده را برمی‌گرداند. API استرایپ (Stripe) دقیقاً همین‌طور کار می‌کند.
  2. جدول dedup / الگوی inbox. مصرف‌کننده شناسهٔ پیام‌های پردازش‌شده را ذخیره می‌کند (با یک TTL یا watermark) و تکراری‌ها را رد می‌کند.
  3. idempotency طبیعی. SET balance = 500 ذاتاً idempotent است (هر چند بار اجرایش کنی، ۵۰۰ می‌شود)؛ اما balance += 100 نیست (هر بار اضافه می‌کند). پس دستورها را طوری طراحی کن که حالتِ مطلق یا یک نسخه (version) را حمل کنند، نه یک دلتای نسبی.
// اندپوینت شارژ idempotent با کلید idempotency + محدودیت یکتا.
@Transactional
public ChargeResult charge(String idempotencyKey, Money amount, String account) {
    // مسیر سریع: آیا همین درخواست دقیق قبلاً پردازش شده؟
    var existing = idempotencyRepo.findByKey(idempotencyKey);
    if (existing != null) {
        return existing.result();          // نتیجهٔ قبلی برگردان، دوباره شارژ نکن
    }
    var result = ledger.debit(account, amount);
    try {
        // UNIQUE(idempotency_key) در دیتابیس گاردِ واقعی در برابر race است
        idempotencyRepo.save(new IdempotencyRecord(idempotencyKey, result));
    } catch (DataIntegrityViolationException race) {
        // تکراریِ همزمان در insert برنده شد؛ نتیجه‌اش را بخوان و برگردان
        return idempotencyRepo.findByKey(idempotencyKey).result();
    }
    return result;
}
چه چیزی واقعاً این کد را درست می‌کند؟

نکتهٔ ظریفی که مصاحبه‌گر دنبالش است: چکِ if (existing != null) گاردِ واقعی نیست. تصور کن دو درخواستِ تکراری دقیقاً هم‌زمان بیایند؛ هر دو چک را رد می‌کنند (چون هنوز هیچ‌کدام ذخیره نشده) و هر دو شارژ می‌کنند! چیزی که واقعاً زیرِ همزمانی نجاتت می‌دهد، محدودیتِ UNIQUE روی ستونِ کلید در دیتابیس است: دیتابیس اجازه نمی‌دهد دو ردیف با یک کلید ثبت شوند، پس یکی از دو درخواست در save شکست می‌خورد، آن exception را می‌گیری و نتیجهٔ برنده را می‌خوانی. چکِ ابتدایی فقط یک بهینه‌سازی است تا در حالت عادی به دیتابیس فشار نیاوری.

تراکنش‌های توزیع‌شده: Saga و الگوی outbox

در مونولیت، یک @Transactional دور همه‌چیز می‌گذاشتی و تمام. اما بینِ دیتابیس‌های چند میکروسرویس نمی‌توانی از two-phase commit (2PC) استفاده کنی: مسدودکننده (blocking) است، دسترس‌پذیریِ کلِ تراکنش را به کندترین شرکت‌کننده گره می‌زند، و منابع را روی شبکه قفل نگه می‌دارد. به‌جایش از Saga استفاده می‌کنی.

رزرو یک سفر با امکان لغو

یک سفر رزرو می‌کنی: اول بلیط هواپیما، بعد هتل، بعد ماشین کرایه‌ای. هر کدام یک تراکنشِ محلیِ جداست. اگر وسط راه ماشین کرایه‌ای پیدا نشد، نمی‌توانی جادویی همه را با هم برگردانی — به‌جایش برای هر مرحله یک اقدامِ جبرانی (compensating action) داری: هتل را کنسل کن، بلیط را کنسل کن. Saga یعنی همین: توالی‌ای از تراکنش‌های محلی که هرکدام یک «دکمهٔ undo» دارند.

دو سبکِ هماهنگی برای Saga:

  • کوریوگرافی (choreography): هر سرویس به رویدادها واکنش نشان می‌دهد و رویدادِ بعدی را منتشر می‌کند. غیرمتمرکز و بدونِ مغزِ مرکزی — مثل رقصنده‌هایی که هر کدام حرکتِ بعدی را از حرکتِ کناری می‌فهمند. اما جریانِ کلی ضمنی است و ردیابی‌اش سخت.
  • ارکستراسیون (orchestration): یک هماهنگ‌کنندهٔ مرکزی (یک ماشینِ حالت) دستورها را صادر می‌کند و جبران‌ها را مدیریت می‌کند — مثل رهبرِ ارکستر. صریح و قابل‌مشاهده، اما خودِ ارکستراتور یک وابستگیِ اضافه است.

مشکلِ dual-write و راه‌حلِ outbox

بخشِ واقعاً سختِ Saga این است: چطور یک رویداد را منتشر کنی و هم‌زمان نوشتنِ دیتابیس را به‌صورت اتمیک کامیت کنی؟

هرگز اول به دیتابیس ننویس و بعد Kafka را صدا نزن

اگر بنویسی db.save() و بعد kafka.send()، یک کرش در فاصلهٔ این دو خط، یا رویداد را گم می‌کند (دیتابیس ذخیره شد ولی رویداد نرفت) یا دوبار می‌فرستد. این «مشکلِ dual-write» است: دو سیستمِ جدا را نمی‌توانی اتمیک کامیت کنی.

راه‌حل: Transactional Outbox. تغییرِ دامنه و یک ردیفِ outbox را در همان یک تراکنشِ محلیِ دیتابیس بنویس. حالا یا هر دو ثبت می‌شوند یا هیچ‌کدام — چون همان دیتابیس، همان تراکنش. سپس یک relay جدا (یا polling روی جدولِ outbox، یا CDC — Change Data Capture — مثلِ Debezium که تغییراتِ دیتابیس را می‌خواند) ردیف‌های outbox را برمی‌دارد و به Kafka منتشر می‌کند.

-- یک تراکنش محلی هم تغییر حالت و هم رویدادِ منتشرشدنی را می‌نویسد.
BEGIN;
  UPDATE orders SET status = 'CONFIRMED' WHERE id = 42;
  INSERT INTO outbox (id, aggregate, type, payload, created_at)
  VALUES (gen_random_uuid(), 'order', 'OrderConfirmed',
          '{"orderId":42}', now());
COMMIT;
-- یک relay جدا (Debezium CDC یا poller) ردیف‌های outbox را به Kafka می‌فرستد،
-- حذف/علامت‌گذاری می‌کند و در شکست retry می‌کند => at-least-once.

چون relay در صورتِ شکست دوباره تلاش می‌کند، تحویل at-least-once است — و همین ما را برمی‌گرداند به درسِ قبلی: پس مصرف‌کننده‌ها باید idempotent باشند. این الگوها زنجیره‌وار به هم وصل‌اند.

CQRS و event sourcing

CQRS مخففِ Command Query Responsibility Segregation است: «تفکیکِ مسئولیتِ دستور و پرس‌وجو». ساده‌اش کن: مدلِ نوشتن را از مدلِ خواندن جدا کن.

آشپزخانه در برابر منوی روی میز

در رستوران، آشپزخانه (سمتِ نوشتن) با دستورِ پخت و مواد خام کار می‌کند — نرمال‌شده، دقیق، پر از قواعد. اما مشتری منوی خوش‌عکسِ روی میز (سمتِ خواندن) را می‌بیند که برای خواندنِ سریع بهینه شده، نه برای پختن. CQRS همین است: با یک مدل می‌نویسی و با یک یا چند مدلِ متفاوتِ بهینه‌شده می‌خوانی.

نوشتن‌ها از یک مدلِ نرمال‌شده و اعمال‌کنندهٔ invariant عبور می‌کنند (قواعد را تضمین می‌کند). خواندن‌ها از یک یا چند projection غیرنرمال (denormalized) سرو می‌شوند که هر کدام برای یک پرس‌وجوی خاص بهینه شده — یک search index این‌جا، یک جدولِ پهن آن‌جا.

  • مزایا: خواندن را مستقل مقیاس می‌دهی؛ storeِ خواندن را به‌ازای هر use-case شکل می‌دهی؛ هر طرف ساده‌تر می‌شود.
  • هزینه: سمتِ خواندن با سمتِ نوشتن سازگارِ نهایی است. یعنی بعد از یک نوشتن، UI ممکن است چند لحظه دادهٔ کهنه نشان بدهد و باید برایش طراحی کنی (مثلاً نتیجهٔ خودِ دستور را خوش‌بینانه به کاربر برگردانی).

Event sourcing یک قدم فراتر می‌رود: حالت را به‌صورتِ یک لاگِ فقط-افزودنی (append-only) از رویدادها ذخیره می‌کند — «OrderPlaced»، «ItemAdded»، «OrderShipped» — نه ردیف‌های فعلی.

دفترِ حسابداری، نه فقط ماندهٔ حساب

یک حسابدارِ خوب فقط عددِ نهاییِ حساب را نگه نمی‌دارد؛ او هر تراکنش را به ترتیب ثبت می‌کند و ماندهٔ فعلی، حاصلِ جمع‌زدنِ (fold) کلِ آن دفتر است. event sourcing همین است: «حالتِ فعلی» را ذخیره نمی‌کنی، بلکه هر تغییر را ثبت می‌کنی و حالت را با پخش‌کردنِ (replay) رویدادها بازمی‌سازی.

مزایا: یک لاگِ audit کامل و بی‌نقص؛ سفر در زمان و دیباگ؛ ساختنِ آسانِ projectionهای جدید فقط با replayِ رویدادها؛ و تناسبِ طبیعی با outbox. اما هزینه‌ها سنگین‌اند و مصاحبه‌گرها دقیقاً همین‌ها را کاوش می‌کنند:

  • تکاملِ schema رویدادها: رویدادها تغییرناپذیر و ابدی‌اند؛ وقتی فرمتِ یک رویداد عوض شود، به versioning/upcasting نیاز داری (رویدادِ قدیمی را هنگام خواندن به فرمتِ جدید «ارتقا» می‌دهی).
  • بازسازیِ projectionها از یک لاگِ طولانی گران است؛ به snapshot (عکسِ فوریِ حالت در یک نقطه، تا از آن‌جا به بعد replay کنی) نیاز داری.
  • پرس‌وجوی موقتیِ (ad-hoc) آسان روی event store وجود ندارد؛ باید projection را پرس‌وجو کنی.
  • «حقِ فراموش‌شدن» GDPR با یک لاگِ تغییرناپذیر تضاد دارد. راهِ‌فرارِ معمول crypto-shredding است: داده را رمز کن و برای «حذف»، فقط کلیدِ رمز را دور بینداز تا داده برای همیشه ناخوانا شود.
اول ساده، بعد پیچیده

CQRS به event sourcing نیاز ندارد (می‌توانی خواندن/نوشتن را جدا کنی بدون لاگِ رویداد). و event sourcing یک تعهدِ جدی است، نه انتخابِ پیش‌فرض. هر دو را فقط جایی به کار ببر که نیازِ audit یا زمانی، پیچیدگی‌شان را توجیه کند — نه چون «باحال» است.

پنج الگوی تاب‌آوری (resilience)

قاعدهٔ ذهنی: فراخواننده باید از خودش محافظت کند. هرگز فرض نکن وابستگی‌هایت بالا و سریع‌اند. پنج ابزار داری.

۱. Timeout (مهلت)

هر فراخوانیِ شبکه به یک timeout نیاز دارد. نبودِ timeout یک قطعیِ نهفته است: یک وابستگیِ کند، thread poolِ تو را یکی‌یکی اشغال می‌کند تا تمام شود، و بعد شکست آبشاری (cascading) می‌شود — سرویسِ تو هم می‌میرد چون همهٔ threadهایش منتظرِ آن وابستگیِ کند مانده‌اند. نکتهٔ مهم: timeout را از بودجهٔ تأخیرِ کلاینت تعیین کن (مشتری چقدر حاضر است صبر کند)، نه از مسیرِ خوش‌بینانهٔ سرور.

۲. Retry — با backoff و jitter

همه با هم به یک درِ گیرکرده هجوم نبرید

تصور کن یک در گیر کرده و صد نفر پشتش‌اند. اگر همه هم‌زمان هُل بدهند، در بدتر گیر می‌کند و همه له می‌شوند. عاقلانه این است که هر کس کمی صبر کند — و مهم‌تر، هر کس یک مقدارِ تصادفیِ متفاوت صبر کند تا همه با هم دوباره هجوم نبرند. این دقیقاً backoff (صبرِ فزاینده) به‌علاوهٔ jitter (پراکندنِ تصادفی) است.

فقط عملیاتِ idempotent و فقط خطاهای قابل‌retry را retry کن (timeout، ۵۰۳ — نه ۴۰۰ یا ۴۰۹ که با تکرار درست نمی‌شوند). retryِ ساده و با فاصلهٔ ثابت باعثِ دو فاجعه می‌شود:

  • retry storm (طوفانِ تکرار): همه بارِ اضافه را روی سرویسی که همین الان در حال زمین‌خوردن است چند برابر می‌کنند.
  • thundering herd (رمهٔ غرّان): بدونِ jitter همهٔ کلاینت‌ها دقیقاً در یک لحظه دوباره تلاش می‌کنند و مثلِ یک حملهٔ DDoS به سرویسِ در حالِ بهبودِ خودت عمل می‌کنند.

راه‌حل: exponential backoff با jitter (فاصله‌ها را نمایی زیاد کن و تصادفی بپاش)، به‌علاوهٔ سقف روی تعدادِ تلاش و زمانِ کل.

۳. Circuit breaker (قطع‌کنندهٔ مدار)

فیوزِ برقِ خانه

وقتی اتصالِ کوتاه می‌شود، فیوز می‌پرد و جریان را قطع می‌کند تا سیم‌کشی آتش نگیرد. بعد از رفعِ مشکل، فیوز را وصل می‌کنی. circuit breaker همین است برای فراخوانیِ سرویس‌ها.

نرخِ شکست به یک وابستگی را رصد کن؛ وقتی از آستانه گذشت، مدار را باز (open) کن و فوراً fail کن (دیگر هیچ فراخوانی‌ای به سرویسِ بیمار نمی‌رسد). بعد از یک دورهٔ خنک‌شدن، چند فراخوانیِ آزمایشیِ نیمه‌باز (half-open) اجازه بده؛ اگر موفق بودند، مدار بسته (closed) و عادی می‌شود. این هم منابعِ تو را از هدررفتن روی سرویسِ مرده نجات می‌دهد، هم به آن سرویس مجالِ بهبود می‌دهد.

۴. Bulkhead (دیوارهٔ آب‌بند)

محفظه‌های ضدغرقِ کشتی

بدنهٔ کشتی به محفظه‌های جدا تقسیم شده؛ اگر یکی سوراخ شد و پر از آب شد، بقیه آب‌بند می‌مانند و کشتی غرق نمی‌شود. bulkhead دقیقاً از همین ایده اسم گرفته.

منابع را جدا کن — thread pool یا connection pool یا سمافورِ جدا به‌ازای هر وابستگی — تا یک downstreamِ اشباع‌شده نتواند همهٔ threadهای تو را ببلعد و ترافیکِ نامرتبط را هم غرق کند. اگر سرویسِ «توصیه‌ها» کند شد، نباید سرویسِ «پرداخت» را با خودش پایین بکشد.

۵. Backpressure (فشارِ معکوس)

وقتی نمی‌توانی هم‌پای ورودی شوی، به‌جای بافر کردنِ نامحدود تا OOM (تمام‌شدنِ حافظه)، به بالادست سیگنال بده که کند شود: صف‌های محدود، رد کردن با کدِ ۴۲۹ (Too Many Requests)، یا سیگنال‌های demand در Flow جاوا و Reactor. پسرعموی خشنِ backpressure، load shedding است: رها کردنِ عمدیِ کارهای کم‌اولویت به‌صورت زودهنگام تا سیستم زیرِ بار نمیرد.

همه با هم در Resilience4j

// Resilience4j: exponential backoff با jitter + circuit breaker، ترکیب‌شده.
IntervalFunction backoff = IntervalFunction
    .ofExponentialRandomBackoff(Duration.ofMillis(100), 2.0, 0.5); // پایه، ضریب، jitter
RetryConfig retryCfg = RetryConfig.custom()
    .maxAttempts(3)
    .retryExceptions(TimeoutException.class, IOException.class) // فقط قابل‌retry
    .intervalFunction(backoff)
    .build();
CircuitBreakerConfig cbCfg = CircuitBreakerConfig.custom()
    .failureRateThreshold(50)                       // باز شدن در ۵۰٪ شکست
    .slowCallRateThreshold(80)
    .slowCallDurationThreshold(Duration.ofSeconds(1))
    .waitDurationInOpenState(Duration.ofSeconds(10)) // خنک‌شدن پیش از half-open
    .slidingWindowSize(20)
    .build();

var cb = CircuitBreaker.of("inventory", cbCfg);
var retry = Retry.of("inventory", retryCfg);

Supplier<Stock> call = () -> inventoryClient.getStock(sku); // timeout مخصوص خود را دارد
Supplier<Stock> guarded = Retry.decorateSupplier(retry,
        CircuitBreaker.decorateSupplier(cb, call));         // CB درون retry
Stock stock = Try.ofSupplier(guarded)
        .recover(ex -> Stock.unknown())                     // fallback نرم
        .get();
ترتیبِ قرار گرفتن مهم است

circuit breaker را درونِ retry بگذار (همان‌طور که در کد decorateSupplierها تودرتو هستند)، تا وقتی breaker باز است، retry بیهوده شلیک نکند و سرویسِ مرده را نکوبد. اما این را صریح در ذهن داشته باش: برخی تیم‌ها عکسش را می‌خواهند (retry درونِ CB) تا چند تلاشِ یک retry، برای آمارِ breaker به‌عنوانِ یک فراخوانیِ منطقی شمرده شود. جوابِ درست بستگی به این دارد که می‌خواهی breaker چه چیزی را بشمارد.

API gateway، BFF و service mesh

سه ابزار برای مدیریتِ ترافیک، در دو مرزِ متفاوت: لبهٔ سیستم (بین کلاینت و سرویس‌ها) و داخلِ سیستم (بین سرویس‌ها).

میزِ پذیرشِ هتل

مهمان مستقیم به آشپزخانه یا اتاقِ تأسیسات نمی‌رود؛ همه‌چیز از میزِ پذیرش می‌گذرد: کارتِ شناسایی چک می‌شود، درخواست‌ها مسیریابی می‌شوند، و مهمان اصلاً نمی‌داند پشتِ صحنه چند بخش هست. API gateway همان میزِ پذیرش است.

  • API gateway. یک نقطهٔ ورودِ واحد که دغدغه‌های عرضیِ لبه را مدیریت می‌کند: پایان‌دهیِ TLS (رمزگشاییِ HTTPS)، احراز هویت/مجوز (authn/authz)، rate limiting (محدودکردنِ نرخ)، مسیریابی، و تجمیعِ درخواست. کلاینت‌ها را از توپولوژیِ داخلی بی‌خبر نگه می‌دارد. ریسک: خودش به یک مونولیتِ هوشمند و گلوگاه تبدیل شود — منطقِ کسب‌وکار را از آن دور نگه دار.
  • BFF (بک‌اند برای فرانت‌اند). گونه‌ای از gateway با یک بک‌اند به‌ازای هر نوعِ کلاینت (وب، iOS، اندروید). هر BFF شکلِ payload و تجمیع را به کلاینتِ خودش می‌دوزد و از یک APIِ «کمترین مخرجِ مشترک» و فراخوانی‌های پرحرفِ موبایل جلوگیری می‌کند. موبایل روی شبکهٔ ضعیف نباید ۱۰ تا درخواست بزند؛ BFF یکی می‌فرستد و بقیه را پشتِ صحنه جمع می‌کند.
  • Service mesh (مثلِ Istio یا Linkerd). دغدغه‌های سرویس-به-سرویس — mTLS (رمزنگاریِ دوطرفه بینِ سرویس‌ها)، retry، timeout، تقسیمِ ترافیک، تله‌متری — را از کدِ تو بیرون می‌کشد و به پراکسی‌های sidecar (مثلِ Envoy) کنارِ هر pod می‌سپارد. یعنی تاب‌آوری و مشاهده‌پذیریِ یکنواخت بدونِ انبوهِ کتابخانه در هر زبان می‌گیری — به بهای پیچیدگیِ عملیاتی و کمی تأخیرِ اضافه در هر hop. وقتی سرویس‌های زیاد در زبان‌های زیاد داری ترجیحش بده؛ یک ناوگانِ کوچک و همگن می‌تواند فقط از کتابخانه استفاده کند.

مشاهده‌پذیری (observability): لاگ، متریک، trace

جعبهٔ سیاهِ هواپیما

وقتی هواپیما سقوط می‌کند، جعبهٔ سیاه به تو اجازه می‌دهد بعداً بفهمی دقیقاً چه شد — حتی سؤالی که از قبل پیش‌بینی نکرده بودی. observability یعنی همین: توانِ پاسخ به مجهول‌های ناشناخته (unknown-unknowns)، نه فقط داشبوردهایی که از پیش ساخته‌ای برای مشکلاتی که انتظارشان را داشتی.

سه ستون داری:

  • لاگ‌ها (logs) — رویدادهای گسسته. ساختاریافته (JSON) بسازشان و همیشه یک شناسهٔ همبستگی / trace ID حمل کنند، تا بتوانی یک درخواست را در عبور از چند سرویس به هم بدوزی.
  • متریک‌ها (metrics) — اعدادِ قابل‌تجمیع (counter شمارنده، gauge سنجه، histogram توزیع). ارزان و عالی برای هشدار و روند. صدک‌ها (percentiles) مثل p95/p99 را ببین، هرگز میانگین را — میانگین آن دنبالهٔ (tail) کندی را که کاربر واقعاً حس می‌کند پنهان می‌کند.
  • Traceها (رد) — مسیرِ علّیِ یک درخواستِ واحد در عبور از سرویس‌ها، با زمان‌بندیِ هر span (هر قطعهٔ کار). OpenTelemetry استانداردِ بی‌طرفِ vendor برای هر سه ستون است؛ کانتکست را با هدرهای W3C traceparent منتشر کن. tracingِ توزیع‌شده روشی است که می‌فهمی کدام hop در زنجیره کند است.

دو مجموعهٔ متریکِ کلاسیک را حفظ کن: روشِ RED برای سرویس‌ها (Rate نرخ، Errors خطاها، Duration مدت) و روشِ USE برای منابع (Utilization بهره‌وری، Saturation اشباع، Errors خطاها). و مهم: روی SLO و error budget هشدار بده (هدفِ سطحِ سرویس و بودجهٔ خطای مجاز)، نه روی CPUِ خام — کاربر به عددِ CPU اهمیت نمی‌دهد، به این اهمیت می‌دهد که سرویس کار کند.


یک چارچوبِ قابل‌استفادهٔ طراحیِ سیستم در مصاحبه

در مصاحبهٔ طراحیِ سیستم، مصاحبه‌گر فرایندِ تو را به اندازهٔ جوابِ نهایی نمره می‌دهد. هر بار همین اسکلتِ شش‌مرحله‌ای را دنبال کن:

شش گامِ همیشگی

۱) نیازمندی‌ها و دامنه ← ۲) تخمینِ سرانگشتی ← ۳) طراحیِ API ← ۴) معماریِ سطح‌بالا و مدلِ داده ← ۵) مقیاس‌دهی و کاوشِ عمیق ← ۶) گلوگاه‌ها، حالت‌های شکست و تعادل‌ها. هرگز پیش از گامِ ۱ جعبه نکش.

۱. نیازمندی‌ها و دامنه. اول نیازمندی‌های کارکردی (سیستم چه کاری می‌کند)، بعد غیرکارکردی (NFR) را روشن کن: مقیاس، تأخیر، دسترس‌پذیری، سازگاری، دوام. و صریحاً بگو چه چیزی خارج از دامنه است. ۲. تخمین (back-of-envelope). QPS (اوج ≈ ۲–۳× میانگین؛ نسبتِ خواندن به نوشتن)، ذخیره‌سازیِ روزانه و ۵ساله، پهنای باند، حافظهٔ کش. اعدادِ گرد به کار ببر؛ هدف، توجیهِ انتخابِ کامپوننت‌هاست نه دقتِ نجومی. ۳. طراحیِ API. اندپوینت‌ها/قراردادهای کلیدی (REST/gRPC/streaming)، شکلِ درخواست/پاسخ، و idempotency. ۴. معماریِ سطح‌بالا و مدلِ داده. کامپوننت‌ها را بکش و storage را به‌ازای الگوی دسترسی انتخاب کن. SQL در برابر NoSQL را با شکلِ پرس‌وجو برگزین، و partition key را عمداً انتخاب کن (کلیدی که داده بر اساسش بین نودها پخش می‌شود). ۵. مقیاس‌دهی و کاوش‌های عمیق. لایه‌های کش، sharding (تکه‌تکه‌کردنِ داده)، replication، load balancing، پردازشِ async، توزیعِ جغرافیایی. ۶. گلوگاه‌ها، حالت‌های شکست، تعادل‌ها. کجا در ۱۰× بار می‌شکند؟ در شکستِ یک نود یا یک ناحیه (region) چه می‌شود؟ انتخابِ CAP/PACELCات را با صدای بلند نام ببر.

مثالِ کامل: ردیابیِ بلادرنگِ ناوگان (fleet tracking)

بیایید این چارچوب را روی یک مسئلهٔ واقعی اجرا کنیم.

۱. نیازمندی‌ها. ۱٬۰۰۰٬۰۰۰ خودرو هر ۵ ثانیه GPS می‌فرستند. قابلیت‌ها: نمایشِ موقعیتِ زندهٔ خودرو، بازپخشِ تاریخچهٔ سفر، و هشدارِ geofence (وقتی خودرو از یک محدودهٔ جغرافیایی خارج شد). NFRها: موقعیتِ زنده می‌تواند سازگارِ نهایی باشد (چند ثانیه کهنگی مشکلی نیست ← روی مسیرِ زنده PACELC از نوعِ EL)؛ تاریخچهٔ سفر باید بادوام باشد؛ مسیرِ نوشتن هرگز نباید نقطه‌ای را بیندازد؛ و p99 پرس‌وجوی نقشه باید < ۳۰۰ms باشد.

۲. تخمین.

  • QPS نوشتن = ۱٬۰۰۰٬۰۰۰ / ۵ = ۲۰۰٬۰۰۰ نوشتن/ثانیه. اوج ≈ ۲× = ۴۰۰k/ثانیه. این سیستم سنگین از نظرِ نوشتن است — و همین نیروی حاکمِ کلِ طراحی است.
  • هر نقطه ≈ ۵۰ بایت (id، lat، lon، ts، speed). در روز: ۲۰۰k × ۸۶٬۴۰۰ × ۵۰ B ≈ ۸۶۴ گیگابایت/روز خام ← فشرده‌سازیِ time-series این را بسیار کوچک‌تر می‌کند.
  • حالتِ زنده: ۱M خودرو × ~۱۰۰ B آخرین موقعیت ≈ ۱۰۰ مگابایت ← به‌راحتی در حافظه جا می‌شود (Redis).

۳. API.

POST /v1/positions            # ingest دسته‌ای از دستگاه‌ها (idempotent با (vehicleId, ts))
GET  /v1/vehicles/{id}/live   # آخرین موقعیت، سرو از کش
GET  /v1/vehicles/{id}/trip?from=..&to=..   # مسیر تاریخی (store سری‌زمانی)
WS   /v1/stream?bbox=...       # پوش آپدیت برای خودروهای داخل viewport نقشه

۴. معماری.

دستگاه‌ها --> Load Balancer --> سرویس Ingest --> Kafka (پارتیشن با vehicleId)
                                                   |            |
                                        (consumer) |            | (consumer)
                                                   v            v
                                       Redis (آخرین موقعیت)   دیتابیس سری‌زمانی
                                       geo-index/حالت داغ     (Cassandra/Timescale)
                                                   |
   کلاینت نقشه <-- فن‌اوت WebSocket <-- سرویس Stream <-- Kafka
  • Ingest بی‌حالت (stateless) و افقی‌مقیاس پشتِ LB است؛ فقط اعتبارسنجی می‌کند و به Kafka تولید می‌کند — اوج‌ها را جذب و یک بافرِ بادوامِ at-least-once فراهم می‌کند (backpressure این‌جا زندگی می‌کند — محدود، و خودِ Kafka بافر است).
  • پارتیشنِ Kafka را با vehicleId بزن تا همهٔ نقاطِ یک خودرو مرتب بمانند و یک خودروی داغ نتواند ترتیبِ خودروی دیگری را به‌هم بریزد.
  • مسیرِ داغ: یک consumer با هر نقطه، Redis را با آخرین موقعیتِ خودرو آپدیت می‌کند (ایندکسِ GEO در Redis پرس‌وجوی «خودروهای داخلِ این bounding box» را پشتیبانی می‌کند). GET /live از Redis می‌خواند ← سریع و سازگارِ نهایی.
  • مسیرِ سرد: consumer دیگری به یک storeِ سری‌زمانی (TimescaleDB یا Cassandra) می‌نویسد، با پارتیشنِ (vehicleId, day) — کلیدی ایده‌آل برای «سفرِ خودروی X در روزِ D» که بار را هم پخش می‌کند. پارتیشن‌های قدیمی به object storage منتقل می‌شوند (retentionِ لایه‌ای).
  • نقشهٔ زنده: کلاینت‌ها روی WebSocket با viewport (محدودهٔ دیدِ نقشه) سابسکرایب می‌کنند؛ سرویسِ stream آپدیت‌های Kafka را به bounding boxِ همان کلاینت فیلتر می‌کند.

۵. مقیاس و گلوگاه.

  • مسیرِ نوشتن با افزودنِ پارتیشنِ Kafka و نمونه‌های consumer مقیاس می‌گیرد؛ نقطهٔ خفگی، توانِ نوشتنِ دیتابیسِ سری‌زمانی است ← نوشتنِ دسته‌ای (batch) و انتخابِ storeِ پرتوان.
  • Idempotency: retry دستگاه می‌تواند نقاط را تکراری کند؛ dedup روی (vehicleId, ts).
  • پارتیشنِ داغ: یک viewport روی شهری متراکم می‌تواند یک stream server را بار کند ← فن‌اوت را با geohash (کدِ جغرافیایی) shard کن، نه با vehicle.
  • حالت‌های شکست: اگر دیتابیسِ سری‌زمانی قطع شود، Kafka داده را نگه می‌دارد (retention = بافرِ دوامِ تو) و consumerها بعداً جبران می‌کنند — و مسیرِ زنده (Redis) هنوز کار می‌کند. اگر Redis بمیرد، آخرین حالت را با مصرفِ انتهای Kafka بازبساز. این پاداشِ قرار دادنِ یک لاگِ بادوام در مرکزِ طراحی است.
  • موضعِ CAP: مسیرِ زنده AP/EL است (موقعیتِ احتمالاً-کهنه سرو کن، هرگز بلاک نکن)؛ ingest+Kafka بادوام و مرتب per key است. این را صریح بگو — همین پاسخِ سؤالِ «چه سازگاری‌ای می‌دهد؟» است.

دام‌ها و نکاتِ ظریف (این‌ها را به‌خاطر بسپار)

  • تقسیم به میکروسرویس به دلایلِ رزومه‌ای، و رسیدن به مونولیتِ توزیع‌شده با دیتابیسِ مشترک.
  • بی‌timeout بودنِ یک فراخوانی ← تمام‌شدنِ thread pool ← شکستِ آبشاری.
  • retryِ عملیاتِ غیر-idempotent ← شارژِ دوگانه / سفارشِ تکراری.
  • retry بدونِ jitter ← retry stormِ همزمان که سرویسِ در حالِ بهبودِ خودت را DDoS می‌کند.
  • باور به وجودِ «تحویلِ دقیقاً-یک‌بار»؛ وجود ندارد — مصرف‌کنندهٔ idempotent مهندسی کن.
  • dual-write (اول دیتابیس بعد Kafka) بدونِ outbox ← رویدادهای گم‌شده یا شبح.
  • هشدار روی میانگین به‌جای p99؛ دنباله همان چیزی است که کاربر تجربه می‌کند.
  • قفلِ توزیع‌شده برای درستی بدونِ fencing token (یک ژتونِ شماره‌دار که هر بار زیاد می‌شود). چرا؟ اگر یک GCِ مکث‌کرده (توقفِ لحظه‌ای برای جمع‌آوریِ زباله) طولانی شود، سرویس فکر می‌کند هنوز قفل را «در دست» دارد در حالی که قفل منقضی شده و به دیگری داده شده — و بدونِ fencing token، نوشتنِ کهنه‌اش داده را خراب می‌کند.

بهترین شیوه‌ها

  • پیش‌فرض یک مونولیتِ ماژولار؛ سرویس را فقط وقتی مقیاس/تیم/نرخ‌تغییر طلب کند استخراج کن.
  • یک دیتابیس به‌ازای هر سرویس؛ یکپارچگی از طریقِ API و رویداد، هرگز schemaِ مشترک.
  • هر مصرف‌کننده را idempotent و هر تولیدکننده را با الگوی outbox بساز.
  • timeout روی همه‌چیز؛ retry فقط برای idempotent + قابل‌retry، همیشه با backoff + jitter؛ وابستگی‌های لغزنده را در circuit breaker و bulkhead بپیچ.
  • یک trace ID را سرتاسری منتشر کن؛ با OpenTelemetry ابزارگذاری کن؛ روی SLO هشدار بده.
  • سازگاری را به‌ازای هر عملیات انتخاب کن و موضعِ CAP/PACELCِ هر مسیر را بنویس.

سؤالات مصاحبه

هر سؤال را در ذهنت مثلِ یک تمرین جواب بده، بعد جوابِ مرجع را ببین.

۱) چرا «تحویلِ دقیقاً-یک‌بار» یک افسانه است و به‌جایش چه می‌سازیم؟ (تله)

روی شبکهٔ نامطمئن، فرستنده نمی‌تواند پیامِ گم‌شده را از ackِ گم‌شده تشخیص دهد، پس باید retry کند و تکراری تولید می‌شود. به‌جایش پردازشِ effectively-once می‌سازی با تحویلِ at-least-once + هندلرِ idempotent (کلیدِ dedup، جدولِ inbox، یا دستورِ طبیعتاً idempotent).

۲) CAP را دقیق بیان کن و بگو چرا «دو از سه» غلط است.

CAP می‌گوید حینِ partition باید C یا A را انتخاب کنی؛ نمی‌توانی سازگاریِ خطی‌پذیر و دسترس‌پذیریِ کامل را هم‌زمان حینِ پارتیشن حفظ کنی. پارتیشن اختیاری نیست، پس محورِ واقعی CP در برابر AP هنگامِ پارتیشن است. PACELC حالتِ Else را می‌افزاید: حتی بدونِ پارتیشن، latency در برابر consistency را معامله می‌کنی.

۳) یک retry برای رفعِ لغزندگی اضافه می‌کنی و قطعی بدتر می‌شود. چرا؟ (تله)

retryها بار را روی وابستگیِ از قبل درمانده چند برابر می‌کنند (retry storm)، و retryهای همزمان یک thundering herd می‌سازند. بدونِ jitter همهٔ کلاینت‌ها با هم retry می‌کنند. راه‌حل: exponential backoff با jitter، سقفِ تلاش، circuit breaker برای failِ سریع، و retry فقط برای خطای idempotent/قابل‌retry.

۴) circuit breaker را نسبت به retry کجا می‌گذاری و چرا؟

معمولاً CB را درونِ دکوریتورِ retry، تا وقتی breaker باز است retryها سریع fail شوند نه اینکه سرویسِ مرده را بکوبند. باید بتوانی حالتِ جایگزین (retry درونِ CB) را هم توجیه کنی، اگر می‌خواهی N retry برای آمارِ breaker یک فراخوانیِ منطقی شمرده شود.

۵) یک اندپوینتِ پرداختِ idempotent طراحی کن. واقعاً چه چیزی درستی را تحتِ همزمانی تضمین می‌کند؟ (سخت)

کلاینت کلیدِ idempotency می‌فرستد؛ سرور کلید+نتیجه را ذخیره می‌کند. تضمینِ درستی یک محدودیتِ UNIQUE روی کلید در دیتابیس است — چکِ «قبلاً پردازش شده؟» فقط بهینه‌سازی است؛ دو تکراریِ همزمان روی insert مسابقه می‌دهند و بازنده نتیجهٔ برنده را می‌خواند.

۶) transactional outbox چیست و چه مشکلی را حل می‌کند؟

مشکلِ dual-write: نمی‌توانی اتمیک هم به دیتابیس کامیت کنی هم به Kafka منتشر. outbox تغییرِ حالت و یک ردیفِ رویداد را در یک تراکنشِ محلی می‌نویسد؛ یک relay (poller یا CDC/Debezium) ردیف‌های outbox را منتشر و retry می‌کند. تحویل at-least-once می‌شود، پس مصرف‌کننده باید idempotent باشد.

۷) CQRS با و بدونِ event sourcing را مقایسه کن. کِی event sourcing ایدهٔ بدی است؟ (سخت)

CQRS فقط مدلِ خواندن/نوشتن را جدا می‌کند؛ سمتِ خواندن سازگارِ نهایی است. event sourcing علاوه بر آن حالت را به‌صورتِ لاگِ رویدادِ تغییرناپذیر ذخیره می‌کند. وقتی نیازِ audit/زمانی نداری بد است، چون هزینهٔ versioning/upcastingِ رویداد، snapshot، بازسازیِ projection، نبودِ پرس‌وجوی موقتی، و تضاد با پاک‌سازیِ GDPR (crypto-shredding) را می‌پردازی.

۸) مونولیتِ توزیع‌شده چیست و چطور بویش را حس می‌کنی؟

سرویس‌هایی که باید lockstep مستقر شوند، دیتابیسِ مشترک دارند، یا زنجیرهٔ همگامِ طولانی می‌سازند. بو: تغییرِ schema در تیم‌ها موج می‌اندازد، دسترس‌پذیری = حاصل‌ضربِ چند hop، هیچ سرویسی تنها منتشر نمی‌شود. مرزها را حولِ bounded context اصلاح کن و به هرکدام دادهٔ خودش را بده.

۹) bulkhead در برابر circuit breaker در برابر backpressure را توضیح بده.

bulkhead منابع را جدا می‌کند (pool مجزا) تا یک وابستگیِ بد بقیه را گرسنه نکند. circuit breaker روی وابستگیِ قطع‌شده سریع fail می‌کند و به آن مجالِ بهبود می‌دهد. backpressure به بالادست سیگنالِ کند شدن می‌دهد (صفِ محدود، ۴۲۹، demandِ ری‌اکتیو) به‌جای بافر تا OOM.

۱۰) چرا روی p99 هشدار دهیم نه میانگینِ تأخیر؟ (تله)

میانگین دنباله را پنهان می‌کند. اگر ۱٪ درخواست‌ها ۵ ثانیه بگیرند، میانگین به‌سختی تکان می‌خورد اما ۱ از ۱۰۰ کاربر تجربهٔ بدی دارد — و در زنجیرهٔ ۱۰ سرویس، احتمالِ برخورد به یک hop از نوعِ p99 حدودِ ۱−۰٫۹۹¹⁰ ≈ ۱۰٪ است. tail latency بر عملکردِ درک‌شدهٔ کاربر غلبه دارد.

۱۱) باگ را پیدا کن.
public void transfer(Account from, Account to, Money amt) {
    inventoryService.reserve(from, amt); // فراخوانی راه‌دور، بدون timeout
    to.credit(amt);                       // محلی
}

(تله) بدونِ timeout روی فراخوانیِ راه‌دور — یک reserve کند، threadِ فراخواننده را نامحدود بلاک و pool را زیرِ بار تمام می‌کند (شکستِ آبشاری). همچنین اگر credit بعد از موفقیتِ reserve شکست بخورد، هیچ جبرانی نیست — به یک saga با اقدامِ جبرانیِ «release» نیاز است، و reserve باید برای retryِ امن idempotent باشد.

۱۲) سیستمِ ردیابیِ خودرو، ماشینی را روی نقشه «به عقب می‌پرد» نشان می‌دهد. محتمل‌ترین علتِ سیستم‌های توزیع‌شده چیست؟ (سخت)

تحویلِ خارج از ترتیب: دو نقطهٔ GPS مسیر/پارتیشنِ متفاوت رفتند و جابه‌جا رسیدند، یا skewِ ساعت میانِ دستگاه‌ها. رفع با پارتیشن‌بندیِ stream با vehicleId (ترتیبِ per-key)، مُهر زدنِ نقاط با زمانِ دستگاه، و انداختنِ نقاطِ قدیمی‌تر از آخرین timestampِ دیده‌شده روی مسیرِ داغ.

۱۳) مدلِ خواندنِ CQRS تو درست بعد از نوشتن کهنه است و کاربر شکایت می‌کند. گزینه‌ها؟

بپذیر و انتظارِ UX را تنظیم کن؛ نتیجهٔ خودِ دستور را خوش‌بینانه برگردان («read-your-writes» از سمتِ نوشتن برای همان actor)؛ یا خواندنِ بعدیِ آن کاربر را کوتاه‌مدت به مدلِ نوشتن مسیریابی کن؛ یا lagِ پروجکشن را کم کن. کلِ سیستم را برای «رفعِ» یک صفحه همگام نکن.

۱۴) قرار دادنِ Kafka در مرکزِ طراحی واقعاً چه سازگاری‌ای می‌خرد؟

یک لاگِ بادوام و مرتب (per-partition) که تولیدکننده را از مصرف‌کننده جدا می‌کند، اوجِ بار را جذب (بافرِ backpressure) و replay برای بازسازیِ حالت یا افزودنِ consumerِ جدید را ممکن می‌کند، و پایهٔ outbox است. ترتیبِ میان‌پارتیشنی یا جادوی exactly-once نمی‌دهد — ترتیب per key و تحویل at-least-once است.

۱۵) کِی از میکروسرویس استفاده *نمی‌کنی*؟ (سنیور)

تیمِ کوچک، مقیاسِ یک deployable کافی، مرزهای دامنه هنوز سیال، یا نبودِ بلوغِ پلتفرم (CI/CD، observability، on-call) برای اجرای سرویس‌های زیاد. میکروسرویسِ زودهنگام مالیاتِ سیستمِ توزیع‌شده را تحمیل می‌کند بدونِ پاداش — با یک مونولیتِ ماژولار شروع کن.

جمع‌بندی
  • توزیع یک هزینه است. مونولیت تراکنشِ اتمیک و تماسِ مطمئن را رایگان می‌داد؛ هر مرزِ شبکه آن‌ها را می‌فروشد. اگر سودِ جداسازی را نمی‌توانی بیان کنی، جدا نکن. هشت مغالطه را رعایت کن.
  • مرزها را حولِ قابلیتِ کسب‌وکار و bounded context بکش، نه لایهٔ فنی. از مونولیتِ توزیع‌شده (دیتابیسِ مشترک، زنجیرهٔ همگامِ پرحرف، انتشارِ lockstep) بترس. با مونولیتِ ماژولار شروع کن.
  • CAP = حینِ پارتیشن بین C و A انتخاب کن؛ PACELC = در حالتِ عادی بین Latency و Consistency. سازگاری یک طیف است — آن را به‌ازای هر عملیات انتخاب کن.
  • exactly-once تحویل وجود ندارد؛ effectively-once = at-least-once + idempotency. گاردِ واقعیِ idempotency یک محدودیتِ UNIQUE در دیتابیس است.
  • Saga جای 2PC را می‌گیرد؛ outbox مشکلِ dual-write را حل می‌کند و تحویل را at-least-once می‌کند.
  • CQRS خواندن/نوشتن را جدا می‌کند (سمتِ خواندن سازگارِ نهایی)؛ event sourcing حالت را به‌صورتِ لاگِ رویداد ذخیره می‌کند — قدرتمند ولی گران.
  • پنج ابزارِ تاب‌آوری: timeout، retry (با backoff+jitter، فقط idempotent/قابل‌retry)، circuit breaker (درونِ retry)، bulkhead، backpressure.
  • observability یعنی پاسخ به مجهول‌های ناشناخته: لاگِ ساختاریافته + متریک (p99 نه میانگین) + traceِ توزیع‌شده با OpenTelemetry. روی SLO هشدار بده.
  • چارچوبِ طراحیِ سیستم: نیازمندی → تخمین → API → معماری/داده → مقیاس → گلوگاه. موضعِ CAP/PACELC هر مسیر را صریح بگو — همان‌طور که در ردیابیِ ناوگان مسیرِ زنده را AP/EL و لاگِ Kafka را بادوام و مرتب per key نامیدیم.

Let us be honest: "microservices" is one of those words everyone puts on a résumé and few can define precisely in an interview. This chapter is going to move you from someone who has heard these patterns to someone who knows which pain each one cures, what it costs, and when NOT to use it. I will never drop a term cold — the moment you see a new word like idempotent, partition, or linearizable, I unpack it in the same breath with an analogy.

Roadmap for this chapter

First we learn why distribution is a cost, not a goal (and those eight famous fallacies). Then proper decomposition and the "distributed monolith" trap. Next the two foundational theorems — CAP and PACELC — and the consistency models, stated precisely. Then practical tools: idempotency and the exactly-once myth, distributed transactions via Saga and outbox, the CQRS/event-sourcing pattern, and five resilience patterns (timeout, retry, circuit breaker, bulkhead, backpressure). After that API gateway/BFF/service mesh and the three pillars of observability. Finally we apply a system-design framework to a real problem (tracking a million vehicles) and close with 15 interview questions.

Part 0 — words you must know first

Before we start, let us nail a few base words with analogies so we never stumble later.

  • Atomic transaction: "all or nothing." Like transferring money between two accounts; either both halves (debit one, credit the other) happen together or neither does. No half-done state allowed.
  • ACID: the four guarantees of a classic database — Atomicity, Consistency (data rules stay valid), Isolation (transactions don't interfere), Durability (once confirmed, data survives even a power cut).
  • hop: every time a request crosses the network from one service to another, that's a hop. Like a phone call between two offices; each call takes time and can drop.
  • latency vs throughput: latency is "how long one request takes" (delivery time of one pizza); throughput is "how many you can serve per unit time" (pizzas per hour). They are different axes.
  • node: one machine/server in the distributed system — one branch office of the company.
  • replica: a copy of data on another node, so if one breaks or gets busy, others can answer.
  • partition (network partition): when network communication between two groups of nodes is cut but the nodes themselves are alive — like the phone line between two branches going down while both branches stay open and busy.
  • idempotent: an operation that yields the same result whether you apply it once or a hundred times. An elevator's "floor 3" button is idempotent: press it a hundred times, you still go to floor 3.
Keep this sentence

Everything that was "free" in the monolith — the reliable call, the atomic transaction — becomes paid the instant you cross a network boundary. This whole chapter is about how to pay that cost or dodge paying it.

Distribution is a cost, not a goal

One kitchen vs several restaurants

Picture a single kitchen. The head chef can say "fire this whole order together" and be sure it's all-or-nothing (atomic transaction); calling to the cook beside you is instant and reliable (in-memory call); and there is only one kitchen to reason about. Now you decide to spread the work across five restaurants across town. Each can now grow and open/close independently — but that "all-together-or-none" is no longer free, and every "cook beside you" is now across town: you must phone them and hope they answer.

A monolith that fits in one process gives you three luxuries for free: atomic transactions, synchronous in-memory calls (call a function, get the answer right there), and one deployable that's easy to reason about.

Every network boundary you add sells those three away in exchange for independent scaling and independent deployment. So the governing question of this whole chapter is:

Is the organizational and scaling benefit of this split worth giving up the free ACID transaction and the reliable function call?

If you can't articulate the benefit of a given split in a sentence, don't split.

The eight fallacies of distributed computing

These eight false beliefs (the "fallacies of distributed computing," from Deutsch and Gosling) are a checklist every design must respect. Each is a comfortable assumption that reality violates:

  1. The network is not reliable (packets get lost).
  2. Latency is not zero (round trips take time).
  3. Bandwidth is not infinite.
  4. The network is not secure.
  5. Topology does change (nodes come and go).
  6. There is not one administrator.
  7. Transport cost is not zero.
  8. The network is not homogeneous (machines and versions differ).
Nearly every incident traces to one of these eight

When a distributed system blows up in production, almost always someone somewhere assumed one of these away. "We assumed the network was fast," "we assumed the message always arrives." Keep this list like an aviation safety checklist.

Decomposition and the distributed-monolith trap

Split departments by "job done," not "job title"

Picture a company. The right split: a "Sales" department, a "Warehouse" department, an "Accounting" department — each does one complete business capability end to end and owns its own work. The wrong split: a "people who answer phones" department, a "people who fill forms" department, a "people who file paperwork" department — these are technical layers, not capabilities. Doing one simple task hand-passes across all three, and none is truly independent.

Good service boundaries follow business capabilities and Domain-Driven Design bounded contexts, not technical layers. "Order service, Payment service, Inventory service" is a capability split. "Controller service, Repository service, Database service" is a layered split that just adds network hops with none of the autonomy.

What's a bounded context? A boundary inside which a word has one precise meaning. "Customer" in Sales is a different model from "Customer" in Support. Each service speaks its own language and model inside its border and doesn't leak them.

The failure mode to fear is the distributed monolith: the worst of both worlds — the complexity of distribution without its independence. Symptoms:

  • Shared database. Two services writing the same tables means you can no longer change a schema, deploy, or scale independently. Every change needs everyone's coordination. This is the single biggest anti-pattern.
  • Chatty synchronous chains. If serving one request means A → B → C → D synchronously, your availability is the product of each hop's availability, and your latency is the sum. A number: if each service is 99% up, 0.99⁴ ≈ 0.96 — a four-hop chain is less trustworthy than any of its links.
  • Lockstep releases. If you can't deploy service X without simultaneously deploying Y, they are really one service wearing two hats.
The shared database is the root of most misery

Golden rule: each service owns its own data. Other services reach that data only through its API or its published events, never by reading its tables directly. The moment two teams write to one table, you no longer have two services; you have a monolith pretending to be two.

Heuristics for decomposition: start coarse-grained (a "modular monolith" — one deployable but with clean, well-bounded internal modules) and extract a service only when a module has one of these three: a distinct scaling profile, a distinct team, or a distinct rate of change.

Conway's Law is not optional

Conway's Law: your system's architecture will mirror your org chart. Teams that design a system unconsciously copy their own communication structure into it. The practical upshot: if service boundaries don't align with team boundaries, you fight constant friction. Design the org and the boundaries together.

CAP and PACELC — state the theorem precisely

The phone line between two branches is cut

You have two bank branches that must keep an account's balance in sync and constantly call each other. One day the phone line between them drops (partition) — but customers are still queued at both counters. Now you have two choices: (a) let both branches keep serving withdrawals — but risk one account being drained twice; or (b) have the branch that isn't sure say "sorry, I can't serve right now" and shut the counter until the line is back. You can't be both "always serve" and "never corrupt data" while the line is cut. That is exactly CAP.

CAP theorem: during a network partition (P), a distributed system must choose one of:

  • Consistency (C): every read sees the latest committed write. This property is called linearizability — it's as if only one copy of the data exists and everyone sees it in the same real-time order.
  • Availability (A): every request gets a non-error response (not necessarily the freshest data, but no error).

While partitioned you can't have both. The common mistake is treating CAP as "pick two of three." But partitions are not optional; they happen in the real world. So the real choice is: CP vs AP, when a partition occurs.

  • CP (e.g., a strongly-consistent store like ZooKeeper/etcd, or a single-leader RDBMS): on partition, the minority side refuses writes so data doesn't diverge. You lose availability, keep correctness.
  • AP (e.g., Dynamo-style stores, or Cassandra with low consistency): on partition, both sides keep serving, and you reconcile later. You keep availability, accept the risk of stale or conflicting reads.

PACELC — the fuller lens

CAP only talks about the moment of partition, but partitions are rare. PACELC (from Daniel Abadi) completes the picture:

If Partition, choose A or C; Else (normal operation), choose Latency or Consistency.

This is the more useful day-to-day lens, because partitions are rare but the latency-vs-consistency trade is always present. To guarantee a read sees the latest write, you must wait for several replicas to coordinate — and that wait is latency.

  • Cassandra by default is PA/EL: on partition favors availability, in normal operation favors low latency (speed).
  • A traditional RDBMS is usually PC/EC: favors consistency in both states.
  • Tune Cassandra with QUORUM reads+writes (a majority of replicas must acknowledge) and it moves toward EC, at the cost of higher latency.
Say this in the interview

"Why is this read stale?" is often not a bug — it's a deliberate EL choice; the system intentionally preferred speed over freshness. Knowing PACELC means you can explain the system's behavior instead of hunting an imaginary bug.

Consistency models

Consistency isn't an on/off switch; it's a spectrum, from strongest (most expensive) to weakest (cheapest and most scalable):

Model Guarantee Typical use
Linearizable / strong Reads see the latest committed write, in real-time order Leader election, locks, account balance
Sequential All nodes see ops in the same order (not necessarily real-time) Some replicated logs
Causal Causally-related ops are seen in order; concurrent ops may reorder Collaborative apps, comments
Read-your-writes You see your own writes Session-level UX (post then refresh)
Eventual Absent new writes, replicas converge Caches, DNS, catalogs

Most large systems are eventually consistent with targeted strong consistency only where money or safety demands it. "Strong everywhere" doesn't scale (too expensive); "eventual everywhere" corrupts invariants (rules that must always hold, like "balance never goes negative"). The senior skill is choosing consistency per operation, not one global decision.

Idempotency and the "exactly-once" myth

A parcel whose delivery receipt got lost

You mail a registered parcel and wait for the delivery receipt (ack). It never arrives. Now you don't know: was the parcel lost, or did it arrive but the receipt was lost? You have no choice but to resend — and if the first one did arrive, the recipient now has two parcels. On any unreliable network you cannot remove this ambiguity. So instead of chasing "exactly one delivery," you make it so that receiving the parcel twice has the same effect as once — e.g., the recipient checks the parcel's number and throws away the duplicate.

There is no exactly-once delivery over an unreliable network. A sender that gets no ack can't tell whether the message was lost or the ack was lost, so it must retry — meaning the receiver sometimes sees duplicates. What you can build is effectively-once processing:

Effectively-once processing = at-least-once delivery + an idempotent handler.

This distinction is one of the most popular interview traps. You can't make delivery once; you can make the effect of processing once.

Idempotency means applying an operation N≥1 times has the same effect as applying it once. Three techniques:

  1. Idempotency keys. The client generates a unique key per logical operation; the server records the key + result. On retry with the same key, it returns the stored result instead of re-executing. Stripe's API works exactly this way.
  2. Dedup table / inbox pattern. The consumer stores processed message IDs (with a TTL or watermark) and skips duplicates.
  3. Natural idempotency. SET balance = 500 is inherently idempotent (run it any number of times, it's 500); balance += 100 is not (it adds each time). So design commands to carry absolute state or a version, not a relative delta.
// Idempotent charge endpoint using an idempotency key + a unique constraint.
@Transactional
public ChargeResult charge(String idempotencyKey, Money amount, String account) {
    // Fast path: was this exact request already processed?
    var existing = idempotencyRepo.findByKey(idempotencyKey);
    if (existing != null) {
        return existing.result();          // return prior result, do NOT charge again
    }
    var result = ledger.debit(account, amount);
    try {
        // UNIQUE(idempotency_key) in DB is the real guard against a race
        idempotencyRepo.save(new IdempotencyRecord(idempotencyKey, result));
    } catch (DataIntegrityViolationException race) {
        // Concurrent duplicate won the insert; reload and return its result
        return idempotencyRepo.findByKey(idempotencyKey).result();
    }
    return result;
}
What actually makes this code correct?

The subtle point the interviewer is hunting for: the if (existing != null) check is not the real guard. Imagine two duplicate requests arrive at exactly the same instant; both pass the check (neither is stored yet) and both charge! What actually saves you under concurrency is the UNIQUE constraint on the key column in the database: the DB won't allow two rows with the same key, so one of the two requests fails at save, you catch that exception, and read the winner's result. The upfront check is only an optimization to avoid hitting the DB in the normal case.

Distributed transactions: Saga and the outbox

In a monolith you wrapped everything in one @Transactional and were done. But across the databases of several microservices you can't use two-phase commit (2PC): it's blocking, it ties the whole transaction's availability to the slowest participant, and it locks resources across the network. Instead you use a Saga.

Booking a trip with cancellations

You book a trip: first the flight, then the hotel, then the rental car. Each is a separate local transaction. If halfway through no car is available, you can't magically roll everything back together — instead each step has a compensating action: cancel the hotel, cancel the flight. A Saga is exactly this: a sequence of local transactions, each with its own "undo button."

Two coordination styles for a Saga:

  • Choreography: each service reacts to events and emits the next event. Decentralized, no central brain — like dancers who each take their next move from the one beside them. But the overall flow is implicit and hard to trace.
  • Orchestration: a central coordinator (a state machine) issues commands and handles compensations — like a conductor. Explicit and observable, but the orchestrator is an extra dependency.

The dual-write problem and the outbox solution

The genuinely hard part of a Saga: how do you publish an event and atomically commit your DB write together?

Never write to the DB then call Kafka

If you do db.save() then kafka.send(), a crash between those two lines either loses the event (DB saved but event never sent) or double-sends it. This is the "dual-write problem": you can't atomically commit two separate systems.

The solution: the Transactional Outbox. Write the domain change and an outbox row in the same local database transaction. Now either both commit or neither does — it's one database, one transaction. Then a separate relay (either polling the outbox table, or CDC — Change Data Capture — like Debezium, which reads the database's change stream) picks up the outbox rows and publishes them to Kafka.

-- One local transaction writes both the state change and the event to publish.
BEGIN;
  UPDATE orders SET status = 'CONFIRMED' WHERE id = 42;
  INSERT INTO outbox (id, aggregate, type, payload, created_at)
  VALUES (gen_random_uuid(), 'order', 'OrderConfirmed',
          '{"orderId":42}', now());
COMMIT;
-- A separate relay (Debezium CDC or a poller) ships outbox rows to Kafka,
-- deletes/marks them, and retries on failure => at-least-once.

Because the relay retries on failure, delivery is at-least-once — which loops us right back to the previous lesson: so consumers must be idempotent. These patterns chain together.

CQRS and event sourcing

CQRS stands for Command Query Responsibility Segregation. Simplify it: separate the write model from the read model.

The kitchen vs the menu on the table

In a restaurant, the kitchen (the write side) works with recipes and raw ingredients — normalized, precise, full of rules. But the customer sees the glossy menu on the table (the read side), optimized for fast reading, not for cooking. CQRS is exactly this: you write through one model and read through one or more different, purpose-optimized models.

Writes go through a normalized, invariant-enforcing model (it guarantees the rules). Reads are served from one or more denormalized projections, each optimized for a specific query — a search index here, a wide table there.

  • Benefits: you scale reads independently; you shape the read store per use-case; each side gets simpler.
  • Cost: the read side is eventually consistent with the write side. After a write, the UI may show slightly stale data for a moment, and you must design for it (e.g., optimistically return the command's own result to the user).

Event sourcing goes one step further: it stores state as an append-only log of events — "OrderPlaced," "ItemAdded," "OrderShipped" — rather than current rows.

The accountant's ledger, not just the balance

A good accountant doesn't just keep the final account balance; they record every transaction in order, and the current balance is the result of folding over the whole ledger. Event sourcing is this: you don't store "current state," you record every change and rebuild state by replaying the events.

Benefits: a perfect, complete audit log; time-travel and debugging; trivially building new projections just by replaying events; and a natural fit with the outbox. But the costs are heavy, and interviewers probe exactly these:

  • Event schema evolution: events are immutable and forever; when an event's format changes you need versioning/upcasting (you "upgrade" an old event to the new format when reading it).
  • Rebuilding projections from a long log is expensive; you need snapshots (a point-in-time capture of state, so you replay only from there).
  • No easy ad-hoc query over the event store; you must query the projections.
  • GDPR's "right to be forgotten" conflicts with an immutable log. The usual escape is crypto-shredding: encrypt the data, and to "delete" it, just throw away the encryption key so the data is permanently unreadable.
Simple first, complex later

CQRS does not require event sourcing (you can separate read/write without an event log). And event sourcing is a serious commitment, not a default. Use either only where audit or temporal requirements justify the complexity — not because it's "cool."

The five resilience patterns

Mental rule: the caller must protect itself. Never assume your dependencies are up and fast. You have five tools.

1. Timeout

Every network call needs a timeout. A missing timeout is a latent outage: one slow dependency occupies your thread pool one by one until it's exhausted, and then the failure goes cascading — your service dies too because all its threads are stuck waiting on that slow dependency. Key point: set the timeout from the client's latency budget (how long the customer will wait), not from the server's happy path.

2. Retry — with backoff and jitter

Don't all rush a stuck door at once

Picture a stuck door with a hundred people behind it. If everyone shoves at once, the door jams harder and everyone gets crushed. The sensible thing is for each person to wait a bit — and, crucially, each to wait a different random amount so they don't all rush again together. That's exactly backoff (increasing wait) plus jitter (random spread).

Retry only idempotent operations and only retryable errors (timeouts, 503 — not 400 or 409, which won't fix on repeat). Naive fixed-interval retries cause two disasters:

  • Retry storm: everyone multiplies the extra load onto a service that is right now collapsing.
  • Thundering herd: without jitter, all clients retry at exactly the same instant and act like a DDoS against your own recovering service.

The fix: exponential backoff with jitter (grow the intervals exponentially and scatter them randomly), plus caps on the number of attempts and total time.

3. Circuit breaker

The fuse in your house

When there's a short circuit, the fuse trips and cuts the current so the wiring doesn't catch fire. After fixing the problem, you reset the fuse. A circuit breaker is this, for service calls.

Track the failure rate to a dependency; when it crosses a threshold, open the circuit and fail fast immediately (no calls reach the sick service). After a cooldown, allow a few half-open trial calls; if they succeed, the circuit closes and returns to normal. This both saves your resources from being wasted on a dead service and gives that service room to recover.

4. Bulkhead

A ship's watertight compartments

A ship's hull is divided into separate compartments; if one is breached and floods, the others stay watertight and the ship doesn't sink. The bulkhead pattern takes its name from exactly this.

Isolate resources — a separate thread pool, connection pool, or semaphore per dependency — so one saturated downstream can't consume all your threads and drown unrelated traffic too. If the "recommendations" service slows down, it must not drag the "payments" service down with it.

5. Backpressure

When you can't keep up with the input, instead of buffering unboundedly until OOM (out of memory), signal upstream to slow down: bounded queues, rejecting with 429 (Too Many Requests), or demand signals in Java's Flow and Reactor. Backpressure's blunt cousin is load shedding: deliberately dropping low-priority work early so the system doesn't die under load.

All together in Resilience4j

// Resilience4j: exponential backoff with jitter + circuit breaker, composed.
IntervalFunction backoff = IntervalFunction
    .ofExponentialRandomBackoff(Duration.ofMillis(100), 2.0, 0.5); // base, multiplier, jitter
RetryConfig retryCfg = RetryConfig.custom()
    .maxAttempts(3)
    .retryExceptions(TimeoutException.class, IOException.class) // only retryable
    .intervalFunction(backoff)
    .build();
CircuitBreakerConfig cbCfg = CircuitBreakerConfig.custom()
    .failureRateThreshold(50)                       // open at 50% failures
    .slowCallRateThreshold(80)
    .slowCallDurationThreshold(Duration.ofSeconds(1))
    .waitDurationInOpenState(Duration.ofSeconds(10)) // cooldown before half-open
    .slidingWindowSize(20)
    .build();

var cb = CircuitBreaker.of("inventory", cbCfg);
var retry = Retry.of("inventory", retryCfg);

Supplier<Stock> call = () -> inventoryClient.getStock(sku); // has its own timeout
Supplier<Stock> guarded = Retry.decorateSupplier(retry,
        CircuitBreaker.decorateSupplier(cb, call));         // CB inside retry
Stock stock = Try.ofSupplier(guarded)
        .recover(ex -> Stock.unknown())                     // graceful fallback
        .get();
Order of composition matters

Put the circuit breaker inside the retry (as the nested decorateSupplier calls show), so that when the breaker is open, a retry doesn't fire uselessly and pound a dead service. But hold this explicitly: some teams want the reverse (retry inside CB) so that a retry's several attempts count as one logical call for the breaker's stats. The right answer depends on what you want the breaker to count.

API gateway, BFF, and service mesh

Three tools for managing traffic at two different boundaries: the edge (between clients and services) and the interior (between services).

The hotel front desk

A guest doesn't walk straight into the kitchen or the boiler room; everything passes through the front desk: ID is checked, requests are routed, and the guest never knows how many departments are behind the scenes. The API gateway is that front desk.

  • API gateway. A single entry point handling cross-cutting edge concerns: TLS termination (decrypting HTTPS), authentication/authorization (authn/authz), rate limiting, routing, and request aggregation. It keeps clients ignorant of internal topology. Risk: it becomes a smart, bottlenecked monolith — keep business logic out of it.
  • BFF (Backend for Frontend). A gateway variant with one backend per client type (web, iOS, Android). Each BFF tailors payload shape and aggregation to its own client, avoiding a lowest-common-denominator API and chatty mobile calls. A phone on a weak network shouldn't fire 10 requests; the BFF sends one and aggregates the rest behind the scenes.
  • Service mesh (e.g., Istio or Linkerd). It pulls service-to-service concerns — mTLS (mutual encryption between services), retries, timeouts, traffic splitting, telemetry — out of your code and into sidecar proxies (like Envoy) beside each pod. So you get uniform resilience and observability without library sprawl across languages — at the cost of operational complexity and a little added latency per hop. Prefer it when you have many services in many languages; a small homogeneous fleet can just use libraries.

Observability: logs, metrics, traces

The airplane's black box

When a plane crashes, the black box lets you figure out afterward exactly what happened — even a question you hadn't anticipated. Observability is this: the power to answer unknown-unknowns, not just the pre-built dashboards for problems you expected.

You have three pillars:

  • Logs — discrete events. Make them structured (JSON) and always carry a correlation / trace ID, so you can stitch one request together as it crosses several services.
  • Metrics — aggregatable numbers (counters, gauges, histograms). Cheap and great for alerting and trends. Watch percentiles like p95/p99, never the average — the average hides the slow tail that the user actually feels.
  • Traces — the causal path of one single request across services, with the timing of each span (each unit of work). OpenTelemetry is the vendor-neutral standard for all three pillars; propagate context via W3C traceparent headers. Distributed tracing is how you find which hop in the chain is slow.

Keep two classic metric sets: the RED method for services (Rate, Errors, Duration) and the USE method for resources (Utilization, Saturation, Errors). And crucially: alert on SLOs and error budgets (service-level objective and allowed error budget), not on raw CPU — the user doesn't care about a CPU number, they care that the service works.


A reusable system-design interview framework

In a system-design interview, the interviewer grades your process as much as the final answer. Every time, follow this same six-step skeleton:

The six ever-present steps
  1. Requirements & scope → 2) Back-of-envelope estimation → 3) API design → 4) High-level architecture & data model → 5) Scaling & deep dives → 6) Bottlenecks, failure modes & trade-offs. Never draw a box before step 1.
  1. Requirements & scope. Clarify functional requirements first (what the system does), then the non-functional (NFR) ones: scale, latency, availability, consistency, durability. And explicitly state what's out of scope.
  2. Estimation (back-of-envelope). QPS (peak ≈ 2–3× average; read:write ratio), storage per day and per 5yr, bandwidth, cache memory. Use round numbers; the point is to justify component choices, not astronomical precision.
  3. API design. The key endpoints/contracts (REST/gRPC/streaming), request/response shapes, and idempotency.
  4. High-level architecture & data model. Draw the components and pick storage per access pattern. Choose SQL vs NoSQL by the query shape, and pick a partition key deliberately (the key that spreads data across nodes).
  5. Scaling & deep dives. Caching layers, sharding (splitting data into pieces), replication, load balancing, async processing, geo-distribution.
  6. Bottlenecks, failure modes, trade-offs. Where does it break at 10× load? What happens on a node or region failure? Name your CAP/PACELC choice out loud.

Worked example: real-time fleet tracking

Let's run this framework on a real problem.

1. Requirements. 1,000,000 vehicles emit GPS every 5 seconds. Features: show a vehicle's live position, replay a trip's history, and geofencing alerts (when a vehicle leaves a geographic zone). NFRs: live position may be eventually consistent (a couple seconds stale is fine → PACELC EL on the live path); trip history must be durable; the write path must never drop a point; and p99 map query must be < 300 ms.

2. Estimation.

  • Write QPS = 1,000,000 / 5 = 200,000 writes/sec. Peak ≈ 2× = 400k/sec. This system is write-heavy — and that's the governing force of the entire design.
  • Each point ≈ 50 bytes (id, lat, lon, ts, speed). Per day: 200k × 86,400 × 50 B ≈ 864 GB/day raw → time-series compression makes this far smaller.
  • Live state: 1M vehicles × ~100 B latest position ≈ 100 MB → fits easily in memory (Redis).

3. API.

POST /v1/positions            # batched ingest from devices (idempotent by (vehicleId, ts))
GET  /v1/vehicles/{id}/live   # latest position, served from cache
GET  /v1/vehicles/{id}/trip?from=..&to=..   # historical path (time-series store)
WS   /v1/stream?bbox=...       # push updates for vehicles in a map viewport

4. Architecture.

Devices --> Load Balancer --> Ingest service --> Kafka (partitioned by vehicleId)
                                                   |            |
                                        (consumer) |            | (consumer)
                                                   v            v
                                       Redis (latest pos)   Time-series DB
                                       geo-index/hot state  (Cassandra/Timescale)
                                                   |
   Map clients <-- WebSocket fanout <-- Stream service <-- Kafka
  • Ingest is stateless and horizontally scaled behind the LB; it only validates and produces to Kafka, absorbing spikes and giving a durable at-least-once buffer (backpressure lives here — bounded, and Kafka itself is the buffer).
  • Partition Kafka by vehicleId so all points for one vehicle stay ordered and one hot vehicle can't reorder another.
  • Hot path: a consumer updates Redis with each vehicle's latest position on every point (Redis's GEO index supports "vehicles in this bounding box"). GET /live reads Redis → fast, eventually consistent.
  • Cold path: another consumer writes to a time-series store (TimescaleDB or Cassandra) partitioned by (vehicleId, day) — an ideal key for "vehicle X's trip on day D" that also spreads load. Old partitions roll off to object storage (tiered retention).
  • Live map: clients subscribe over WebSocket by viewport (the map's visible area); the stream service filters Kafka updates down to that client's bounding box.

5. Scaling & bottlenecks.

  • The write path scales by adding Kafka partitions and consumer instances; the choke point is the time-series DB's write throughput → batch the writes and pick a high-ingest store.
  • Idempotency: device retries can duplicate points; dedup on (vehicleId, ts).
  • Hot partition: a viewport over a dense city can overload one stream server → shard the fanout by geohash (geographic code), not by vehicle.
  • Failure modes: if the time-series DB is down, Kafka retains the data (retention = your durability buffer) and consumers catch up later — and the live path (Redis) still works. If Redis dies, rebuild latest state by consuming the tail of Kafka. This is the payoff of putting a durable log at the center of the design.
  • CAP stance: the live path is AP/EL (serve possibly-stale positions, never block); ingest+Kafka is durable and ordered per key. State this explicitly — it's the answer to "what consistency does this give?"

Pitfalls & gotchas (memorize these)

  • Splitting into microservices for résumé reasons, ending with a distributed monolith and a shared DB.
  • No timeout on a call → thread-pool exhaustion → cascading failure.
  • Retrying non-idempotent operations → double charges / duplicate orders.
  • Retries without jitter → synchronized retry storms that DDoS your own recovering service.
  • Believing "exactly-once delivery" exists; it doesn't — engineer idempotent consumers.
  • Dual-write (DB then Kafka) without an outbox → lost or phantom events.
  • Alerting on averages instead of p99; the tail is what users experience.
  • Distributed locks used for correctness without fencing tokens (a numbered token that increments each time). Why? If a paused GC (a brief stop-the-world for garbage collection) runs long, a service thinks it still "holds" the lock while it has actually expired and been handed to another — and without a fencing token its stale write corrupts the data.

Best practices

  • Default to a modular monolith; extract a service only when scaling/team/change-rate demands it.
  • One database per service; integrate via APIs and events, never a shared schema.
  • Make every consumer idempotent and every producer use the outbox pattern.
  • Timeouts on everything; retries only for idempotent + retryable, always with backoff + jitter; wrap flaky deps in circuit breakers and bulkheads.
  • Propagate a trace ID end-to-end; instrument with OpenTelemetry; alert on SLOs.
  • Choose consistency per operation, and write down the CAP/PACELC posture of each path.

Interview Questions

Answer each in your head like an exercise, then check the reference answer.

1) Why is "exactly-once delivery" a myth, and what do we build instead? (gotcha)

Over an unreliable network the sender can't distinguish a lost message from a lost ack, so it must retry, producing duplicates. Instead you build effectively-once processing via at-least-once delivery + an idempotent handler (dedup keys, an inbox table, or naturally idempotent commands).

2) State CAP precisely, and explain why "pick two of three" is wrong.

CAP says during a partition you must choose C or A; you can't preserve linearizable consistency and full availability while partitioned. Partitions aren't optional, so the real axis is CP vs AP on partition. PACELC adds the Else: even without a partition, you trade latency vs consistency.

3) You add a retry to fix flakiness and the outage gets worse. Why? (gotcha)

Retries multiply load on an already-struggling dependency (retry storm), and synchronized retries create a thundering herd. Without jitter every client retries at once. Fixes: exponential backoff with jitter, capped attempts, a circuit breaker to fail fast, and retrying only idempotent/retryable errors.

4) Where do you place the circuit breaker relative to the retry, and why?

Typically the CB inside the retry decorator, so that when the breaker is open, retries fail fast instead of hammering a dead service. Be able to justify the alternative (retry inside CB) if you want N retries to count as one logical call for the breaker's stats.

5) Design an idempotent payment endpoint. What actually guarantees correctness under concurrency? (hard)

The client sends an idempotency key; the server stores key+result. The correctness guarantee is a UNIQUE constraint on the key in the database — the "already processed?" check is only an optimization; two concurrent duplicates race on the insert and the loser reads the winner's result.

6) What's the transactional outbox and what problem does it solve?

The dual-write problem: you can't atomically commit to your DB and publish to Kafka. The outbox writes the state change and an event row in one local transaction; a relay (poller or CDC/Debezium) publishes the outbox rows and retries. Delivery becomes at-least-once, so consumers must be idempotent.

7) Compare CQRS with and without event sourcing. When is event sourcing a bad idea? (hard)

CQRS just separates read/write models; the read side is eventually consistent. Event sourcing additionally stores state as an immutable event log. It's a bad fit when you don't need audit/temporal features but will pay for event versioning/upcasting, snapshotting, projection rebuilds, no ad-hoc queries, and GDPR-erasure conflicts (crypto-shredding).

8) What is a distributed monolith and how do you smell one?

Services that must deploy in lockstep, share a database, or form long synchronous call chains. Smells: schema changes ripple across teams, availability = product of many hops, no service can release alone. Fix the boundaries around bounded contexts and give each its own data.

9) Explain bulkhead vs circuit breaker vs backpressure.

Bulkhead isolates resources (separate pools) so one bad dependency can't starve the others. Circuit breaker fails fast on a known-down dependency and gives it room to recover. Backpressure signals upstream to slow down (bounded queues, 429, reactive demand) instead of buffering to OOM.

10) Why alert on p99, not average latency? (gotcha)

Averages hide the tail. If 1% of requests take 5 s, the average barely moves but 1-in-100 users has a terrible experience — and in a chain of 10 services, the probability of hitting some p99 hop is ~1−0.99¹⁰ ≈ 10%. Tail latency dominates user-perceived performance.

11) Find the bug.
public void transfer(Account from, Account to, Money amt) {
    inventoryService.reserve(from, amt); // remote call, no timeout
    to.credit(amt);                       // local
}

(gotcha) No timeout on the remote call — a slow reserve blocks the calling thread indefinitely, exhausting the pool under load (cascading failure). Also, there's no compensation if credit fails after reserve succeeds — this needs a saga with a compensating "release" action, and reserve should be idempotent for safe retries.

12) A vehicle-tracking system shows a car "jumping backward" on the map. What distributed-systems cause is most likely? (hard)

Out-of-order delivery: two GPS points took different paths/partitions and arrived reordered, or clock skew across devices. Fix by partitioning the stream by vehicleId (per-key ordering), stamping points with device time, and dropping points older than the last-seen timestamp on the hot path.

13) Your CQRS read model is stale right after a write and users complain. Options?

Accept it and set UX expectations; return the command's own result optimistically ("read-your-writes" from the write side for that actor); or route that user's next read to the write model briefly; or reduce projection lag. Don't make the whole system synchronous to "fix" one screen.

14) What consistency does putting Kafka at the center of a design actually buy you?

A durable, ordered (per-partition) log that decouples producers from consumers, absorbs load spikes (backpressure buffer), enables replay to rebuild state or add new consumers, and underpins the outbox. It does not give cross-partition ordering or exactly-once magic — order is per key, delivery is at-least-once.

15) When would you NOT use microservices? (senior)

Small team, a single deployable's scaling is adequate, domain boundaries still fluid, or you lack the platform maturity (CI/CD, observability, on-call) to operate many services. Premature microservices impose the distributed-systems tax with none of the payoff — start with a modular monolith.

In a nutshell
  • Distribution is a cost. The monolith gave you atomic transactions and reliable calls for free; every network boundary sells them away. If you can't articulate a split's benefit, don't split. Respect the eight fallacies.
  • Draw boundaries around business capabilities and bounded contexts, not technical layers. Fear the distributed monolith (shared DB, chatty synchronous chains, lockstep releases). Start with a modular monolith.
  • CAP = during a partition choose C or A; PACELC = in normal operation choose Latency or Consistency. Consistency is a spectrum — choose it per operation.
  • Exactly-once delivery doesn't exist; effectively-once = at-least-once + idempotency. Idempotency's real guard is a UNIQUE constraint in the DB.
  • Saga replaces 2PC; the outbox solves the dual-write problem and makes delivery at-least-once.
  • CQRS separates read/write (read side eventually consistent); event sourcing stores state as an event log — powerful but expensive.
  • Five resilience tools: timeout, retry (with backoff+jitter, idempotent/retryable only), circuit breaker (inside the retry), bulkhead, backpressure.
  • Observability is answering unknown-unknowns: structured logs + metrics (p99 not average) + distributed traces with OpenTelemetry. Alert on SLOs.
  • System-design framework: requirements → estimation → API → architecture/data → scaling → bottlenecks. State each path's CAP/PACELC posture explicitly — just as, in fleet tracking, we named the live path AP/EL and the Kafka log durable and ordered per key.