Databases & SQL · پایگاه‌داده و SQL متوسطIntermediate ~40 دقیقه مطالعه~33 min read

Redis، ClickHouse، ScyllaDB و ElasticsearchRedis, ClickHouse, ScyllaDB & Elasticsearch

از صفر یاد می‌گیری که چهار پایگاه‌دادهٔ تخصصی — Redis، ClickHouse، ScyllaDB و Elasticsearch — هرکدام برای چه کاری ساخته شده‌اند، کِی درست و کِی فاجعه‌اند، و درونیّاتی که مصاحبه‌گر سنیور واقعاً می‌کاود.Learn from scratch what each of four specialized datastores — Redis, ClickHouse, ScyllaDB, and Elasticsearch — is actually built for, when each is right, when each is a disaster, and the internals senior interviewers really probe.


خب، بیا یک تصویر ذهنی درست کنیم. یک آشپزخانهٔ حرفه‌ای را تصور کن. آنجا یک اجاقِ همه‌کاره داری که تقریباً هر غذایی را می‌پزد — این همان پایگاه‌دادهٔ رابطه‌ای (Postgres، MySQL، Oracle) توست: منبعِ حقیقتِ همه‌چیز. اما کنارش ابزارهای تخصصی هم داری: یک سرخ‌کنِ فوری، یک فِر پیتزای داغِ داغ، یک یخچالِ ویترینی. هیچ‌کدام «اجاقِ بهتر» نیستند؛ هرکدام یک کارِ خاص را ده برابر بهتر انجام می‌دهند. Redis، ClickHouse، ScyllaDB و Elasticsearch دقیقاً همان ابزارهای تخصصیِ کنارِ اجاق‌اند.

نقشهٔ راه این درس

در این فصل چهار پایگاه‌دادهٔ تخصصی را از پایه یاد می‌گیری:

  1. Redis — انبارِ ساختار-دادهٔ درون‌حافظه: کش، صف، قفل، leaderboard.
  2. ClickHouse — پایگاه‌دادهٔ ستونیِ تحلیلی (OLAP) روی موتور MergeTree.
  3. ScyllaDB / Cassandra — پایگاه‌دادهٔ ستون‌پهن با نوشتنِ عظیم و سازگاریِ قابل‌تنظیم.
  4. Elasticsearch — موتور جستجوی متن روی ایندکس معکوس.

برای هرکدام سه چیز را یاد می‌گیری: مدل داده (شکلش)، کِی درست است، و کِی فاجعه است. مهم‌ترین سیگنالِ سنیوری همین سومی است.

بخش ۰ — واژه‌هایی که باید بشناسی

قبل از هر چیز چند واژه را از صفر باز کنیم تا بعداً سرد رهایشان نکنم.

OLTP در برابر OLAP: صندوق‌دار در برابر حسابدار

یک فروشگاه را تصور کن. صندوق‌دار هزاران بار در روز یک تراکنش کوچک انجام می‌دهد: «این یک قلم را بفروش، موجودی این یک محصول را کم کن». سریع، تک‌ردیفی، دقیق. این OLTP است (Online Transaction Processing) — کارِ روزمرهٔ پایگاه‌دادهٔ رابطه‌ای.

آخر ماه حسابدار می‌آید و می‌پرسد: «مجموعِ فروشِ هر شعبه در هر ماهِ سال گذشته چقدر بوده؟» او یک ردیف را عوض نمی‌کند؛ میلیون‌ها ردیف را می‌خواند و جمع می‌زند. این OLAP است (Online Analytical Processing) — کارِ ClickHouse.

همان دیتابیسی که برای صندوق‌دار عالی است، برای حسابدار کند است و برعکس. کلِ این فصل حول همین تفاوت می‌چرخد.

  • ACID — تضمینِ اینکه یک تراکنش یا کاملاً انجام می‌شود یا اصلاً (اتمیک)، داده هیچ‌وقت خراب نمی‌شود (سازگار)، تراکنش‌های هم‌زمان توی هم نمی‌روند (ایزوله) و چیزی که commit شد گم نمی‌شود (بادوام). این تضمینِ طلاییِ پایگاه‌دادهٔ رابطه‌ای است؛ هر چهار ابزارِ این فصل بخشی از آن را کنار می‌گذارند تا سرعت بگیرند.
  • system of record (منبع حقیقت) — آن یک انباری که «نسخهٔ درست و رسمیِ» داده در آن است. پولِ کاربر، سفارش، هویت. اگر جایی و اینجا با هم فرق کنند، حرفِ اینجا درست است.
  • در‌حافظه (in-memory) — داده در RAM زندگی می‌کند نه روی دیسک. برای همین Redis سریع است و برای همین اگر برق برود ممکن است داده برود (مگر تدبیر کنی).
  • کاردینالیتی (cardinality) — تعدادِ مقدارهای متمایز در یک ستون. ستونِ «جنسیت» کاردینالیتیِ پایین (چند مقدار) و ستونِ «ایمیل» کاردینالیتیِ بالا (تقریباً یکتا) دارد. جلوتر مهم می‌شود.
  • idempotent (خنثی‌به‌تکرار) — عملیاتی که اگر دو بار اجرا شود همان نتیجهٔ یک بار را می‌دهد. «مقدار را روی ۱۰ بگذار» idempotent است؛ «مقدار را یکی زیاد کن» نیست.
هیچ‌کدام رقیب هم نیستند

اشتباهِ رایجِ جونیورها این است که فکر می‌کنند باید «یکی» از این‌ها را انتخاب کنند. نه. این چهار، راه‌حل‌های نقطه‌ای برای مسائلی‌اند که Postgres/MySQL در مقیاس بالا بد حل می‌کنند. در معماریِ واقعی هر چهار می‌توانند کنارِ همان RDBMS باشند، هرکدام یک دردِ خاص را دوا می‌کنند.

تصویرِ بزرگ: چهار ابزار، چهار وظیفه

این جدول را در ذهنت قاب کن؛ کلِ فصل بسطِ همین است:

ابزار شکل داده نقطهٔ قوت برای این‌ها استفاده نکن
Redis کلید→مقدار درون‌حافظه (مقدارهای غنی) کش، نشست، محدودیت نرخ، صف، leaderboard منبع حقیقت برای دادهٔ حیاتی و بادوام
ClickHouse ستونی، عمدتاً append OLAP: تحلیل و تجمیع روی میلیاردها سطر OLTP: آپدیت تک‌سطری، تراکنش، خواندن تک‌سطری با هم‌روندی بالا
ScyllaDB / Cassandra ستون‌پهن، پارتیشن‌بندی‌شده نوشتن سنگین، مقیاس افقی، الگوی دسترسی مشخص کوئری دلخواه، join، تراکنش قوی چند‌کلیدی
Elasticsearch ایندکس معکوس (سند) جستجوی متن کامل، مرتبط‌سازی، تحلیل لاگ منبع حقیقت؛ هر چیزی که ACID می‌خواهد

یک سیگنالِ همیشگیِ سنیوری: دانستنِ اینکه هر کدام کجا اشتباه است. در دنیای علی — APIهای پرترافیک Spring Boot در نوآوشگران و کشِ Redis در نشان — همچنان RDBMS منبع حقیقت است و این چهار، شتاب‌دهنده‌هایی هستند که پیرامونش سوار می‌شوند.


Redis

مدل داده — چرا «فقط کلید-مقدار» نیست

یک قفسهٔ کمدِ باشگاه، اما هر کمد یک شکل دارد

یک HashMap ساده مثل ردیفی از کمدهای یک‌شکل است: هر کلید یک جعبه که فقط یک رشتهٔ ساده تویش می‌گذاری. Redis فرق دارد: هر کمد می‌تواند شکلِ متفاوتی داشته باشد — یکی صف است، یکی مجموعهٔ رتبه‌بندی‌شده، یکی دفترچهٔ فیلد→مقدار. تو نه فقط داده، بلکه ساختارِ درستِ داده را انتخاب می‌کنی. همین انتخاب، خودِ بهینه‌سازی است.

Redis یک سرورِ ساختار-دادهٔ درون‌حافظه است که اجرای دستورها در آن تک‌رشته‌ای (single-threaded) است. تک‌رشته‌ای یعنی در هر لحظه فقط یک دستور اجرا می‌شود — که چون هیچ قفلی لازم نیست، هر دستور به‌شکلِ طبیعی اتمیک می‌شود. مقدارِ پشتِ یک کلید فقط رشته نیست؛ یک ساختارِ نوع‌دار است:

  • String — بایت تا ۵۱۲ مگابایت؛ به‌عنوان شمارنده (INCR) و bitmap هم به کار می‌رود.
  • Hash — نگاشتِ فیلد→مقدار؛ یک آبجکت را بدونِ ساختنِ N کلید جدا ذخیره می‌کنی.
  • List — لیستِ پیوندی؛ LPUSH از یک سر و BRPOP از سرِ دیگر یعنی یک صف.
  • Set / Sorted Set (ZSet) — یکتایی. ZSet یک skiplist + hash است و درجِ رتبه‌دار را در O(log N) انجام می‌دهد → همان چیزی که leaderboard و ایندکسِ زمان‌مرتب می‌خواهد.
  • Stream — لاگِ فقط-افزودنی با گروه مصرف‌کننده (یک کافکای کوچک درونِ Redis).
  • به‌علاوهٔ HyperLogLog (شمارشِ کاردینالیتی در حدود ۱۲ کیلوبایت)، Geo، Bitfield و Pub/Sub.
ساختار را درست انتخاب کن، نصفِ کار تمام است

یک leaderboard روی ZSet فقط یک دستورِ ZREVRANGE است. اگر بخواهی همان را با Stringها شبیه‌سازی کنی — که هر بار همه را بخوانی، مرتب کنی و بنویسی — یک فاجعهٔ کارایی می‌سازی. در Redis، «کدام ساختار؟» مهم‌ترین تصمیم است.

الگوهای کش — نامی که در مصاحبه می‌بری مهم است

انباردار و انبارِ عقب

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

  • Cache-aside (بارگذاریِ تنبل) — اپ خودش کش را چک می‌کند؛ در miss از DB می‌خواند و کش را پر می‌کند. ساده و مقاوم است (خرابیِ کش ≠ خرابیِ اپ)، اما اولین درخواست کند است و داده می‌تواند کهنه (stale) باشد. پیش‌فرض است و همان چیزی که علی در نشان استفاده کرد.
  • Read-through — کتابخانهٔ کش به‌جای تو، در miss به‌صورتِ شفاف از DB بارگذاری می‌کند.
  • Write-through — نوشتن هم‌زمان به کش و DB می‌رود. کش همیشه تازه، اما نوشتن کندتر.
  • Write-behind (write-back) — نوشتن به کش، و فلاشِ ناهم‌گام (async) به DB. نوشتن سریع، اما ریسکِ ازدست‌رفتنِ داده در کرش.

حالا همان cache-aside را در Java ببین. به کامنت‌ها دقت کن — دو تدبیرِ ظریف تویش هست که جلوتر توضیح می‌دهم:

// Cache-aside با Spring Data Redis + TTL و کش‌کردن null برای جلوگیری از stampede/penetration
public Product getProduct(long id) {
    String key = "product:" + id;
    String cached = redis.opsForValue().get(key);
    if (cached != null) {
        return cached.equals("__NULL__") ? null           // کش منفی: cache penetration را می‌کشد
                                          : deserialize(cached);
    }
    Product p = repository.findById(id).orElse(null);
    // TTL با نویز (jitter) از بهمنِ هم‌زمان (cache stampede) جلوگیری می‌کند
    long ttl = 300 + ThreadLocalRandom.current().nextLong(60);
    redis.opsForValue().set(key,
            p == null ? "__NULL__" : serialize(p),
            Duration.ofSeconds(p == null ? 30 : ttl));
    return p;
}

سه بلای کلاسیکِ کش

اسمِ این سه را باید بلد باشی؛ مصاحبه‌گر عاشقشان است.

درِ پشتی، تعطیلیِ هم‌زمان، و هجوم به تنها باجه
  • Cache penetration (نفوذ): یک نفر مدام سراغِ محصولی می‌آید که اصلاً وجود ندارد. قفسهٔ جلو همیشه خالی است، پس هر بار می‌روی انبارِ عقب و دست‌خالی برمی‌گردی. مهاجم می‌تواند با همین، دیتابیس را از پا درآورد. رفع: خودِ «نبودن» را هم کش کن (__NULL__) یا از Bloom filter استفاده کن.
  • Cache avalanche (بهمن): همهٔ قفسه‌ها را یک‌جا شبِ عید چیدی، پس همه هم‌زمان یک سال بعد خالی می‌شوند و در یک لحظه همه‌چیز به انبار هجوم می‌برد. رفع: به TTLها نویز (jitter) بده تا انقضاها پخش شوند.
  • Cache breakdown (شکستِ کلید داغ): یک کالای فوق‌محبوب از قفسه تمام می‌شود و در همان لحظه هزار مشتری آن را می‌خواهند؛ هزار نفر هم‌زمان می‌روند انبار تا بازسازی‌اش کنند. رفع: mutex/singleflight بگذار تا فقط یک رشته بازسازی کند و بقیه منتظر بمانند یا دادهٔ کهنه بگیرند.

حالا کدِ بالا را دوباره نگاه کن: __NULL__ همان کشِ منفی است (رفعِ penetration)، و ttl = 300 + random(60) همان jitter است (رفعِ avalanche). این‌ها ترفندهای تولیدی‌اند، نه تزیین.

حذف (Eviction) — وقتی حافظه پر می‌شود

RAM بی‌نهایت نیست. وقتی مصرف به maxmemory می‌رسد، تنظیمِ maxmemory-policy تصمیم می‌گیرد چه چیزی بیرون انداخته شود:

  • noeviction (پیش‌فرض) — دیگر چیزی حذف نمی‌شود؛ نوشتن‌های جدید خطا می‌دهند.
  • allkeys-lru / allkeys-lfu — کم‌استفاده‌ترین (بر اساسِ زمانِ آخرین استفاده / فراوانیِ استفاده) را از میانِ همهٔ کلیدها حذف کن.
  • volatile-lru / volatile-lfu / volatile-ttl — همان، ولی فقط از میانِ کلیدهایی که TTL دارند.
  • allkeys-random / volatile-random — تصادفی.
LRU در Redis تقریبی است، نه دقیق

LRU یعنی «کم‌استفاده‌ترین از نظرِ زمان» و LFU یعنی «کم‌استفاده‌ترین از نظرِ تعداد دفعات». اما Redis برای پیداکردنِ دقیقِ آن‌ها همهٔ کلیدها را اسکن نمی‌کند — فقط چند کلید را نمونه‌برداری می‌کند (maxmemory-samples) و بهترین نامزد را بین همان‌ها می‌اندازد بیرون. یعنی کمی دقت را با سرعت معامله می‌کند. برای یک کشِ خالص، allkeys-lru درست است؛ و از Redis 4، allkeys-lfu معمولاً برای دسترسیِ چوله (وقتی چند کلید خیلی داغ‌ترند) بهتر است.

پایداری: RDB در برابر AOF

عکسِ خانوادگی در برابر دفترچهٔ خاطرات

Redis درون‌حافظه است، اما دو راه دارد که اگر ری‌استارت شد داده برنگردد به صفر:

  • RDB مثل گرفتنِ یک عکسِ خانوادگیِ هرچند‌ساعت‌یک‌بار است: تصویری کامل از یک لحظه. اگر بین دو عکس اتفاقی بیفتد، آن لحظه‌ها را نداری.
  • AOF مثل نوشتنِ دفترچهٔ خاطرات است: هر کاری که کردی را همان لحظه یادداشت می‌کنی. با بازخوانیِ دفترچه می‌توانی همه‌چیز را بازبسازی.
  • RDB — عکسِ فوریِ (snapshot) دودوییِ دوره‌ای در یک نقطهٔ زمانی. Redis برای این کار خودش را fork() می‌کند (یک کپیِ فرزند از فرایند می‌سازد) و از copy-on-write استفاده می‌کند. فشرده است و ری‌استارتش سریع، اما در کرش هر چیزی بعد از آخرین snapshot را از دست می‌دهی.
  • AOF — لاگِ فقط-افزودنی از دستورهای نوشتن. appendfsync everysec (پیش‌فرض) حداکثر حدود ۱ ثانیه داده از دست می‌دهد؛ always بادوام ولی کند است. AOF به‌صورت دوره‌ای بازنویسی/فشرده می‌شود.
  • بهترین روش: هر دو — RDB برای بازیابیِ سریع، AOF برای دوام. Redis 7 حالتِ چند-بخشیِ AOF را اضافه کرد (یک RDB پایه + یک AOF افزایشی).
«Redis فقط کش است پس بادوام نیست» — این حرف غلط است

Redis می‌تواند بادوام باشد. اما تلهٔ واقعی جای دیگری است: همان fork() برای RDB یا بازنویسیِ AOF، روی دیتاستِ بزرگ می‌تواند حافظه را تقریباً دو برابر کند (به‌خاطرِ copy-on-write) و اسپایکِ تأخیر بسازد. پس بادوام‌بودن رایگان نیست.

Pub/Sub و Streams

پخشِ زندهٔ رادیو در برابر صندوقِ پستی

PUBLISH/SUBSCRIBE مثلِ پخشِ زندهٔ رادیو است: بفرست و فراموش کن. اگر رادیوی تو آن لحظه روشن نباشد، آن قطعهٔ موسیقی برای همیشه رفت — نه ضبط می‌شود نه دوباره پخش. Streams مثلِ صندوقِ پستی است: نامه‌ها می‌مانند تا بیایی برداری، رسید می‌گیری (XACK) و اگر لازم شد دوباره می‌خوانی.

پس PUBLISH/SUBSCRIBE بدونِ پایداری و بدونِ تضمینِ تحویل است؛ مشترکِ آفلاین پیام‌ها را از دست می‌دهد. برای پیام‌رسانیِ بادوام از Streams (XADD/XREADGROUP) با گروهِ مصرف‌کننده، ack (XACK) و بازپخش استفاده کن.

قفل توزیع‌شده — و نکاتی که سنیور را از جونیور جدا می‌کند

SET key value NX PX 30000 یک قفل می‌دهد: مقدار را فقط اگر وجود نداشت بگذار (NX)، با انقضای ۳۰ ثانیه (PX). آزادسازی‌اش را هم باید با احتیاط انجام دهی — با یک اسکریپتِ Lua که «فقط اگر توکن مالِ خودت بود، حذف کن»:

// آزادسازی امن: مالکیت را چک کن و اتمیک حذف کن — هرگز DELِ ساده نزن
String lua = "if redis.call('get', KEYS[1]) == ARGV[1] " +
             "then return redis.call('del', KEYS[1]) else return 0 end";
redis.execute(new DefaultRedisScript<>(lua, Long.class), List.of(key), token);

اما دو نکتهٔ ظریف هست که مصاحبهٔ سنیوری دقیقاً روی همین دو انگشت می‌گذارد:

  1. یک Redisِ تکی یک SPOF است و در failover خطی‌پذیر (linearizable) نیست. SPOF یعنی «نقطهٔ شکستِ واحد» — اگر بیفتد، همه‌چیز می‌افتد. و چون همتاسازی به replica ناهم‌گام است، قفلی که روی master داده شد، اگر master قبل از رساندنِ آن به replica failover کند، ناپدید می‌شود.
  2. Redlock (گرفتنِ قفل از اکثریتِ N مسترِ مستقل) وجود دارد، اما نقدِ Martin Kleppmann پابرجاست: Redlock به این فرض تکیه دارد که انحرافِ ساعت و مکث‌های GC/توقفِ برنامه کران‌دار (bounded) است. حالا فرض کن یک فرایند به‌خاطرِ یک وقفهٔ طولانیِ GC، بیشتر از TTL متوقف بماند: کلیدش منقضی می‌شود و یکی دیگر قفل را می‌گیرد، در حالی که اولی هنوز فکر می‌کند قفل را دارد. حالا دو نفر توی ناحیهٔ بحرانی‌اند.
fencing token: تنها راهِ درستی

برای درستیِ واقعی به یک fencing token نیاز داری: یک عددِ صعودیِ یکنوا که هر بار قفل داده می‌شود یکی زیاد می‌شود، و خودِ منبعِ محافظت‌شده آن را چک می‌کند و هر درخواستی با توکنِ قدیمی‌تر را رد می‌کند. با این، حتی اگر دو نفر فکر کنند قفل را دارند، فقط آن‌که تازه‌ترین توکن را دارد اجازهٔ نوشتن می‌گیرد. جمع‌بندیِ ذهنی: قفلِ Redis طردِ متقابلِ best-effort است، نه تضمینِ درستی. از آن برای کاهشِ کارِ تکراری استفاده کن، نه برای تضمینِ عدمِ وقوعش.


ClickHouse

مدلِ داده و چرا این‌قدر سریع است

کتابخانهٔ سطری در برابر کتابخانهٔ ستونی

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

ClickHouse یک پایگاه‌دادهٔ OLAPِ ستونی (columnar) است. برای کوئریِ تحلیلی‌ای که فقط ۳ ستون از ۲۰۰ ستون را لمس می‌کند، فقط فایلِ همان ۳ ستون خوانده می‌شود — بدونِ I/Oِ هدررفته. مزیتِ دوم: وقتی همهٔ مقدارهای یک ستون هم‌نوع‌اند، به‌شدت خوب فشرده می‌شوند (LZ4/ZSTD، و برای timestampها الگوهای delta و double-delta). مزیتِ سوم: موتورِ برداری (vectorized) داده را در بلوک‌هایی پردازش می‌کند که در کشِ CPU جا می‌شوند. همین سه با هم، جوابِ تجمیع روی میلیاردها سطر را در میلی‌ثانیه برمی‌گردانند.

خانوادهٔ MergeTree

MergeTree موتورِ اصلی است؛ تقریباً همه‌چیز از آن مشتق می‌شود. اول کد را ببین، بعد خط‌به‌خط بازش می‌کنیم:

CREATE TABLE events (
    event_date  Date,
    user_id     UInt64,
    event_type  LowCardinality(String),
    ts          DateTime,
    payload     String
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_date)          -- پارتیشن درشت (برای TTL/حذف)، ایندکس نیست
ORDER BY (event_type, user_id, ts)         -- کلید مرتب‌سازی = ترتیب فیزیکی + ایندکس sparse
SETTINGS index_granularity = 8192;         -- سطر در هر granule (mark)
فهرستِ کتاب که هر ۸۱۹۲ صفحه یک ورودی دارد

یک B-tree معمولی مثلِ فهرستی است که برای هر صفحه یک ورودی دارد — دقیق، اما اگر کتاب میلیاردها صفحه باشد، خودِ فهرست هم غول‌آسا می‌شود. ClickHouse فهرستی می‌سازد که فقط برای هر ۸۱۹۲ صفحه یک ورودی دارد؛ به هر بلوکِ ۸۱۹۲ سطری می‌گویند یک granule. این فهرستِ «تُنُک (sparse)» آن‌قدر کوچک است که حتی برای میلیاردها سطر در RAM جا می‌شود. وقتی دنبالِ چیزی می‌گردی، فهرست می‌گوید «در حدودِ این granule است»، و ClickHouse فقط همان بلوکِ ۸۱۹۲ سطری را می‌خواند نه کلِ جدول را.

رفتارِ جاریِ تأییدشده (مستنداتِ ClickHouse، ۲۰۲۵):

  • ORDER BY همان کلیدِ مرتب‌سازی (sorting key) است: ترتیبِ فیزیکیِ واقعیِ سطرها روی دیسک درونِ هر part را تعریف می‌کند.
  • کلیدِ اصلی (primary key) یک ایندکسِ sparse است، جداگانه ذخیره می‌شود، با یک ورودی به‌ازای هر granule — و granule پیش‌فرض ۸۱۹۲ سطر است. پس ایندکس حدودِ یک ورودی به‌ازای هر ۸۱۹۲ سطر دارد. این کاملاً برخلافِ B-tree است که هر سطر را ایندکس می‌کند.
  • PRIMARY KEY می‌تواند با ORDER BY فرق کند، فقط اگر پیشوند (prefix) تاپلِ ORDER BY باشد. این می‌گذارد ایندکسِ درونِ RAM را کوچک نگه داری در حالی که با ستون‌های بیشتری مرتب می‌کنی.
  • granularityِ تطبیقی (adaptive) به‌طورِ پیش‌فرض از طریقِ index_granularity_bytes (حدودِ ۱۰ مِبی‌بایت) فعال است: یک granule یا در ۸۱۹۲ سطر یا در حدودِ ۱۰ MiB — هرکدام زودتر رسید — پایان می‌یابد. این وقتی سطرها بزرگ‌اند از حافظه محافظت می‌کند.
  • کوئری‌هایی که WHEREشان با پیشوندِ کلیدِ مرتب‌سازی می‌خواند، از ایندکسِ sparse برای رد کردنِ کلِ granuleها استفاده می‌کنند؛ در غیرِ این صورت ClickHouse اسکنِ کامل می‌کند (هنوز سریع، اما بدونِ هرس).
Partها مثلِ لایه‌های زمین‌شناسی؛ merge مثلِ فشرده‌شدنشان

هر INSERT یک partِ تغییرناپذیرِ تازه می‌سازد — یک جدولِ کوچکِ مرتب، مثلِ یک لایهٔ جدید رسوب. یک فرایندِ پس‌زمینه مدام این لایه‌ها را در لایه‌های بزرگ‌تر ادغام (merge) می‌کند — اسمِ موتور از همین است: MergeTree. حالا نکته: اگر تو هزاران درجِ ریز بزنی (سطر‌به‌سطر)، هزاران لایهٔ نازک می‌سازی و فرایندِ merge از پا می‌افتد و خطای «too many parts» می‌گیری. برای همین باید دسته‌ای و بزرگ درج کنی.

اعضای کلیدیِ خانواده:

  • ReplacingMergeTree — سطرهای با کلیدِ مرتب‌سازیِ یکسان را هنگامِ merge حذفِ تکراری می‌کند (در نهایت؛ تا OPTIMIZE ... FINAL یا FINAL در کوئری تضمین‌شده نیست).
  • SummingMergeTree / AggregatingMergeTree — هنگامِ merge پیش‌تجمیع می‌کنند.
  • CollapsingMergeTree / VersionedCollapsingMergeTree — آپدیت/حذف را با سطرهای علامتِ ‎+۱/−۱‎ مدل می‌کنند.
  • ReplicatedMergeTree — همتاسازی از طریقِ ZooKeeper/ClickHouse Keeper.

و skip indexها (minmax, set, bloom_filter) با GRANULARITY مخصوصِ خودشان، اجازه می‌دهند روی ستون‌هایی که کلیدِ مرتب‌سازی نیستند هم granuleها را رد کنی.

چه زمانی ClickHouse اشتباه است

ClickHouse را برای OLTP به کار نبر — اشتباهِ کلاسیک همین است
  • معنای واقعیِ UPDATE/DELETE ندارد. آپدیت‌ها (ALTER TABLE ... UPDATE) که اسمشان mutation است، کلِ partها را ناهم‌گام بازنویسی می‌کنند — یک عملیاتِ ادمینیِ سنگین، نه ویرایشِ تراکنشیِ یک سطر.
  • تراکنشِ ACIDِ چند‌عبارتی ندارد، سازگاریِ تک‌سطری‌اش ضعیف است، و حذفِ تکراری‌اش نهایی (eventual) است.
  • جست‌وجوی نقطه‌ای برای یک سطر بر اساسِ کلیدِ اصلی، در مقایسه با یک انبارِ OLTP ناکارآمد است — برای اسکنِ بازه ساخته شده، نه واکشیِ یک سطر.
  • بارِ کاری‌اش «چند کوئریِ تحلیلیِ سنگین» است، نه «هزاران کوئریِ کوچکِ هم‌زمان».

شکلِ درستِ معماری: OLTP در MySQL/Oracle (منبعِ حقیقت)، و بعد stream/ETL به ClickHouse برای داشبورد و تحلیل. هرگز مسیرِ نوشتنت را برای ویرایشِ رکوردِ تکی به آن نشانه نگیر.


ScyllaDB / Cassandra

ScyllaDB یک بازنویسیِ C++ از Cassandra است: همان CQL، همان مدلِ داده، اما با معماریِ shard-per-core و «shard-aware» و بدونِ توقف‌های GCِ JVM. هر چه دربارهٔ مدلِ داده می‌گویم برای هر دو صادق است.

مدلِ دادهٔ ستون‌پهن — بر اساسِ کوئری مدل کن، نه موجودیت

آشپزخانهٔ فست‌فود که منو را از قبل می‌داند

در یک رستورانِ معمولی (دیتابیسِ رابطه‌ای) مواد را نرمال و جدا نگه می‌داری و هر سفارش را در لحظه با join سرِهم می‌کنی. یک فست‌فودِ پرترافیک این‌طور کار نمی‌کند: منو از قبل معلوم است، پس هر آیتم را از پیش آماده و بسته‌بندی‌شده نگه می‌دارد تا سفارش در یک ثانیه بیرون برود. ScyllaDB همین است: تو کوئری‌هایت را از قبل می‌دانی و به‌ازای هر کوئری یک جدول می‌سازی، حتی اگر داده تکراری شود.

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

CREATE TABLE messages_by_room (
    room_id    uuid,
    bucket     int,
    ts         timeuuid,
    user_id    uuid,
    body       text,
    PRIMARY KEY ((room_id, bucket), ts, user_id)  -- (کلید پارتیشن)(کلیدهای clustering)
) WITH CLUSTERING ORDER BY (ts DESC);
  • کلیدِ پارتیشن (partition key) ((room_id, bucket)) — هش می‌شود تا تصمیم بگیرد داده روی کدام نود برود. همهٔ سطرهای یک پارتیشن روی همان replicaها کنارِ هم زندگی می‌کنند. این واحدِ توزیع و واحدِ اتمیک‌بودنِ تک‌پارتیشنی است.
  • کلیدهای clustering (ts, user_id)ترتیبِ مرتب‌سازیِ درونِ پارتیشن را روی دیسک تعریف می‌کنند و اسکنِ بازه درونِ پارتیشن را ارزان می‌کنند.

قواعدِ طلایی:

  1. هر کوئری باید با کلیدِ پارتیشنِ کاملش به یک پارتیشن بخورد. کوئریِ بدونِ آن، دیتابیس را مجبور به scatter-gather روی همهٔ نودها می‌کند — که علامتش ALLOW FILTERING است، یک پرچمِ قرمز در تولید.
  2. از پارتیشن‌های داغ (چولگیِ ترافیک روی یک پارتیشن) و پارتیشن‌های بی‌کران (پارتیشنی که برای همیشه رشد می‌کند) پرهیز کن — به همین دلیل آن bucket را گذاشتیم: تا یک اتاقِ پرترافیک را روی چند پارتیشن بشکنیم.
  3. غیرنرمال‌سازی کن: به‌ازای هر الگوی دسترسی یک جدول بساز و دادهٔ تکراری را بپذیر.

سازگاریِ قابل‌تنظیم (tunable consistency)

رأی‌گیری میانِ نسخه‌های برابر

در Cassandra/Scylla هیچ رهبرِ واحدی نیست؛ هر replica برابرِ بقیه است و داده RF بار کپی می‌شود (replication factor، مثلاً ۳ نسخه). برای هر کوئری، تو تصمیم می‌گیری چند نسخه باید موافقت کنند — این همان سطحِ سازگاری (CL) است. مثلِ رأی‌گیری: هرچه رأیِ بیشتری بخواهی، مطمئن‌تری اما کندتر.

  • نوشتن به همهٔ replicaها فرستاده می‌شود؛ وقتی به‌تعدادِ CL تأیید برگشت، نوشتن موفق است (ONE, QUORUM, ALL, LOCAL_QUORUM…).
  • خواندن با CL replica تماس می‌گیرد و جواب‌ها را آشتی می‌دهد.
  • سازگاریِ قوی وقتی برقرار است که R + W > RF باشد. مثلاً با RF=3، خواندنِ QUORUM (R=2) + نوشتنِ QUORUM (W=2) می‌دهد 2+2 > 3، پس هر خواندن حتماً آخرین نوشتن را می‌بیند. اگر این نامساوی برقرار نباشد، سازگاریِ نهایی (eventual) می‌گیری و ممکن است دادهٔ کهنه بخوانی.
  • LOCAL_QUORUM رأی‌گیری را برای کاهشِ تأخیر درونِ یک دیتاسنتر نگه می‌دارد (در استقرارهای چند-دیتاسنتری).

سازگاری به‌صورتِ تنبل ترمیم می‌شود: از طریقِ read repair (هنگامِ خواندن)، hinted handoff (نگه‌داشتنِ نوشتن برای نودی که موقتاً پایین است) و ترمیمِ anti-entropy با nodetool repair.

تراکنش‌های سبک‌وزن (LWT)

نوشتنِ پیش‌فرض «last-write-wins» است و هیچ شرطی نمی‌پذیرد. اما گاهی واقعاً به compare-and-set نیاز داری — مثلِ «این نامِ کاربری را فقط اگر کسی نگرفته، بگیر». اینجا از LWT با IF استفاده می‌کنی:

INSERT INTO users (username, id) VALUES ('ali', ...) IF NOT EXISTS;
UPDATE accounts SET balance = 90 WHERE id = ? IF balance = 100;
LWT قدرت دارد اما گران است — تأییدشده (مستندات ScyllaDB، ۲۰۲۵)
  • LWT برای اجماع میانِ replicaهای یک پارتیشنِ واحد از Paxos استفاده می‌کند — همهٔ شرط‌ها باید به همان پارتیشن اشاره کنند؛ تراکنشِ بین‌پارتیشنی وجود ندارد.
  • سازگاریِ سریال (serial) می‌دهد، اما تقریباً سه تا چهار رفت‌وبرگشت هزینه دارد (prepare/read/propose/commit) در برابرِ یک رفت‌وبرگشت برای نوشتنِ عادی → اغلب یک مرتبهٔ بزرگی کندتر. کم استفاده کن.
  • برای خواندنِ آخرین مقدار در وسطِ تراکنش باید با سازگاریِ SERIAL/LOCAL_SERIAL بخوانی؛ خواندنِ QUORUMِ عادی ممکن است مقدارِ در حالِ پروازِ Paxos را نبیند.

چه زمانی اشتباه است

بدونِ join، بدونِ WHEREِ دلخواه روی ستون‌های اختیاری، بدونِ تراکنشِ قویِ چند‌پارتیشنی، و برای الگوهای کوئریِ در حالِ تحول یا نامعلوم دردسر است. قانونِ سرانگشتی: اگر نمی‌توانی کوئری‌هایت را از پیش فهرست کنی، این پایگاه‌داده اشتباه است.


Elasticsearch

ایندکسِ معکوس — قلبِ جستجو

فهرستِ اعلامِ آخرِ کتاب

پشتِ جلدِ یک کتابِ درسی، فهرستِ اعلام (index) هست: برای هر واژهٔ مهم، فهرستِ صفحه‌هایی که آن واژه در آن‌هاست. تو نمی‌آیی کلِ کتاب را ورق بزنی تا دنبالِ «مهاجرت» بگردی؛ می‌روی سراغِ فهرستِ اعلام و مستقیم صفحه‌ها را می‌بینی. ایندکسِ معکوس دقیقاً همین است: نگاشتِ ترم → فهرستِ سندهایی که آن ترم را دارند.

یک ایندکسِ رابطه‌ای یک سطر را به مقدارهای ستونش نگاشت می‌کند؛ ایندکسِ معکوس عکسش را می‌کند — هر ترم → فهرستِ سندها (postings) حاویِ آن، همراه با موقعیت و فراوانی. همین است که «هر سندِ حاویِ migration را پیدا کن» را تقریباً O(1) می‌کند نه اسکنِ کامل. خودِ Elasticsearch یک لایهٔ توزیع‌شده روی Apache Lucene است؛ Lucene صاحبِ ساختارهای ایندکس (segment، postings، doc values) است.

تحلیل‌گرها (analyzers)

کارخانهٔ خردکردنِ متن

قبل از اینکه متن وارد فهرستِ اعلام شود، از یک خطِ تولید می‌گذرد: character filter → tokenizer → token filter. مثلاً "The Migrations!" وارد می‌شود، کوچک می‌شود، روی نویسه‌های غیرحرف بریده می‌شود، stopwordها (مثلِ «the») حذف می‌شوند، و ریشه‌یابی (stem) می‌شود تا در نهایت بشود [migrat]. نکتهٔ حیاتی: همان خطِ تولید باید هنگامِ کوئری هم اجرا شود، وگرنه ترم‌های تو با ترم‌های ذخیره‌شده جور درنمی‌آیند.

به همین دلیل دو نوع فیلد داری: فیلدِ تحلیل‌نشدهٔ keyword (دقیق و توکنایز‌نشده) برای فیلتر/تجمیع/مرتب‌سازی، و فیلدِ text برای جستجوی متنِ کامل. یک باگِ رایج: انتظار داری روی یک فیلدِ text تطبیقِ دقیق یا تجمیع بگیری — در حالی که به زیرفیلدِ .keyword نیاز داری.

مرتبط‌سازی و BM25

امتیازدهیِ منصفانه به نتایج

وقتی «مهاجرتِ دیتابیس» را جستجو می‌کنی، Elasticsearch باید تصمیم بگیرد کدام سند بالاتر بیاید. سه شهودِ ساده را ترکیب می‌کند: (۱) سندی که واژه را بیشتر دارد بهتر است — اما نه بی‌حساب؛ (۲) واژهٔ نادر مهم‌تر از واژهٔ رایج است؛ (۳) تطبیق در یک عنوانِ کوتاه ارزشمندتر از همان تطبیق در یک متنِ بلند و پرحرف است. مدلی که این‌ها را جمع می‌کند BM25 نام دارد.

از Elasticsearch 5 / Lucene 6، مدلِ پیش‌فرضِ مرتبط‌سازی BM25 است (جایگزینِ TF-IDFِ کلاسیک). BM25 امتیازِ یک سند برای یک کوئری را با ترکیبِ این‌ها می‌دهد:

  • فراوانیِ ترم (TF) — اما با اشباع (saturation): دهمین تکرارِ یک واژه بسیار کمتر از دومین می‌افزاید (پارامترِ k1).
  • فراوانیِ معکوسِ سند (IDF) — ترم‌های نادر وزنِ بیشتری از رایج‌ها دارند.
  • نرمال‌سازیِ طولِ فیلد — تطبیق در یک عنوانِ کوتاه بر همان تطبیقِ مدفون در بدنهٔ بلند می‌چربد (پارامترِ b).

همان اشباع + نرمالِ طول دقیقاً دلیلِ برتریِ BM25 بر TF-IDFِ ساده است (چون keyword-stuffing دیگر بی‌حساب امتیاز نمی‌گیرد)، و یک سؤالِ کلاسیکِ سنیور «چرا عوض کردند» است.

شاردها و رِپلیکاها

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

یک ایندکسِ بزرگ را نمی‌شود در یک ماشین جا داد، پس به چند تکه (primary shard) تقسیمش می‌کنی و هر تکه را روی یک ماشین می‌گذاری — این مقیاسِ افقی و موازی‌سازی است. بعد از هر تکه چند کپی (replica) می‌سازی برای دوامِ بالا و توانِ خواندنِ بیشتر.

  • یک ایندکس به primary shardها تقسیم می‌شود (هرکدام یک ایندکسِ Luceneِ مستقل) — واحدِ مقیاسِ افقی و موازی‌سازی. تعدادِ primary shard در زمانِ ساخت ثابت است (برای تغییرش باید reindex/split کنی).
  • هر primary، replica shard دارد — کپی برای HA و توانِ خواندن؛ تعدادِ replica به‌صورتِ زنده تغییرپذیر است.
  • یک سند با hash(routing) % number_of_primary_shards به یک شارد مسیریابی می‌شود — به همین دلیل تعدادِ primary تغییرناپذیر است.
  • segmentهای Lucene تغییرناپذیراند؛ نوشتن اول به یک بافرِ درون‌حافظه + یک translog (برای دوام) می‌رود، بعد با refresh به segmentهای تازهٔ قابل‌جستجو تبدیل می‌شود (پیش‌فرض هر ۱ ثانیه → «near real-time») و در پس‌زمینه merge می‌شود. حذف‌ها tombstone‌اند که هنگامِ merge بازپس‌گرفته می‌شوند.
oversharding: شاردهای کوچکِ زیاد

تعدادِ زیادِ شاردِ کوچک، heap و state خوشه را هدر می‌دهد. هدف را روی شاردهایی در محدودهٔ ده‌ها گیگابایت بگذار، نه صدها شاردِ چند‌مگابایتی.

چه زمانی اشتباه است

Elasticsearch منبعِ حقیقت نیست: بدونِ تراکنشِ ACID، بدونِ join (فقط parent/child و nestedِ محدود)، و near-real-time است نه بلافاصله سازگار. آن را به‌عنوانِ ایندکسِ جستجو/تحلیل که از DBِ اصلیِ تو (یا Kafka/CDC) تغذیه می‌شود به کار ببر، هرگز به‌عنوانِ تنها انبارِ حقیقت برای دادهٔ حیاتی.


انتخاب بین این‌ها (سؤالِ واقعیِ مصاحبه)

اگر فقط یک جمله را از این فصل به خاطر بسپاری، این جدولِ تصمیم است:

  • نیاز به خواندنِ زیرمیلی‌ثانیهٔ دادهٔ داغ / نشست / شمارنده / قفل → Redis.
  • نیاز به GROUP BY روی میلیاردها سطر برای داشبورد → ClickHouse.
  • نیاز به توانِ نوشتنِ عظیم با کوئریِ قابل‌پیش‌بینیِ پارتیشن‌محور و بدونِ نقطهٔ شکستِ واحد → ScyllaDB/Cassandra.
  • نیاز به «جستجو با واژه‌ها»، تحملِ غلطِ تایپی، رتبه‌بندیِ مرتبط، کاوشِ لاگ → Elasticsearch.
  • منبعِ حقیقتِ تو برای دادهٔ تراکنشیِ رابطه‌ایِ حیاتیِ مالی تقریباً همیشه Postgres/MySQL/Oracle می‌ماند — و این چهار، لایه‌های تخصصیِ خواندن/مقیاس پیرامونِ آن‌اند که با ETL، CDC یا نوشتنِ دوگانه هم‌گام نگه داشته می‌شوند.

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

۱. Redis تک‌رشته‌ای است — چطور بیش از ۱۰۰ هزار عملیات در ثانیه سرو می‌کند و چه چیزی آن را می‌شکند؟

اجرای دستور تک‌رشته‌ای است، پس عملیات‌ها اتمیک‌اند و سربارِ قفل نیست؛ کار CPU-ارزان و درون‌حافظه است. با I/O multiplexing (epoll) و از Redis 6 با I/Oِ چندرشته‌ای برای خواندن/نوشتنِ سوکت (پارس همچنان تک‌رشته‌ای) مقیاس می‌گیرد. چه چیزی می‌شکندش: یک دستورِ O(N) (KEYS *، SMEMBERSِ بزرگ، SORT) همه‌چیز را بلاک می‌کند. از SCAN استفاده کن و از دستورهای بزرگِ بلاک‌کننده بپرهیز.

۲. (سخت) قفلِ توزیع‌شدهٔ Redis گاهی دو worker را به ناحیهٔ بحرانی راه می‌دهد. چرا؟

چند دلیل: (الف) failoverِ مستر قبل از همتاسازیِ قفل به replica؛ (ب) توقفِ فرایندِ نگه‌دارنده (GC/گرسنگیِ CPU) فراتر از TTL، پس کلید منقضی شد و کلاینتِ دیگری آن را گرفت در حالی که اولی فکر می‌کند قفل را دارد؛ (ج) DELِ ساده که قفلِ دیگری را بعد از انقضای TTLِ تو حذف می‌کند. رفع: توکنِ یکتا + Luaِ compare-and-delete، و برای درستی یک fencing token که منبع اعتبارسنجی کند. طبقِ Kleppmann، Redlock تحتِ توقف‌های بی‌کران برای درستی امن نیست — best-effort است.

۳. تفاوت `ORDER BY` و `PRIMARY KEY` در MergeTree کلیک‌هاوس؟

ORDER BY کلیدِ مرتب‌سازی است: ترتیبِ فیزیکیِ سطرها روی دیسک + پایهٔ ایندکسِ sparse. PRIMARY KEY، اگر جداگانه مشخص شود، باید پیشوندِ ORDER BY باشد؛ می‌گذارد ایندکسِ درونِ RAM را کوچک کنی در حالی که با ستون‌های اضافی مرتب می‌کنی. ایندکس یک mark به‌ازای هر granule (پیش‌فرض ۸۱۹۲ سطر) ذخیره می‌کند، نه به‌ازای هر سطر.

۴. (سخت) ۵۰۰۰ سطر در ثانیه را سطر‌به‌سطر به کلیک‌هاوس درج می‌کنی و از پا می‌افتد. چرا و راه‌حل؟

هر درج یک partِ تغییرناپذیرِ تازه می‌سازد؛ mergerِ پس‌زمینه نمی‌رسد → «too many parts» و توقف. کلیک‌هاوس برای درجِ دسته‌ای ساخته شده. رفع: بافر کن و در بلوک‌های بزرگ (ده‌ها هزار سطر) درج کن، async insert، یا جدولِ Buffer / موتورِ Kafka. رایج‌ترین اشتباهِ تولیدیِ کلیک‌هاوس همین است.

۵. چرا کلیک‌هاوس برای OLAP سریع اما برای OLTP اشتباه است؟

ذخیره‌سازیِ ستونی فقط ستون‌های ارجاع‌شده را می‌خواند و به‌شدت فشرده می‌کند؛ موتورِ برداری بازه‌ها را با ایندکسِ sparse اسکن می‌کند. اما آپدیتِ تراکنشیِ سطر نیست (mutation کلِ part را ناهم‌گام بازنویسی می‌کند)، ACID نیست، و جست‌وجوی نقطه‌ایِ تک‌سطری ناکارآمد است. برای اسکن و تجمیع ساخته شده، نه ویرایش و واکشیِ یک سطر.

۶. در کاساندرا/اسکایلا تفاوتِ partition key و clustering key چیست و چرا بر اسکیمای تو مسلط است؟

partition key هش می‌شود تا داده را روی نودها بگذارد و واحدِ هم‌مکانی و اتمیک‌بودن را تعریف می‌کند؛ clustering key سطرها را درونِ پارتیشن روی دیسک مرتب می‌کند. چون کوئریِ کارآمد باید کلیدِ پارتیشنِ کامل را بدهد، به‌ازای هر کوئری یک جدول طراحی می‌کنی («مدل‌سازیِ کوئری-اول») و غیرنرمال می‌سازی — کلیدها همان مدلِ داده‌اند.

۷. چه زمانی خواندن در کاساندرا قویاً سازگار است؟

وقتی R + W > RF. با RF=3، نوشتنِ QUORUM (W=2) + خواندنِ QUORUM (R=2) می‌دهد 4 > 3 → هر خواندن آخرین نوشتنِ commit‌شده را می‌بیند. CLهای پایین‌تر سازگاریِ نهایی و خواندنِ کهنهٔ محتمل می‌دهند؛ LOCAL_QUORUM تأخیر را درونِ یک دیتاسنتر نگه می‌دارد.

۸. (سخت) چه زمانی LWT در اسکایلا به کار می‌بری و هزینه‌اش چیست؟

وقتی به معنایِ compare-and-set نیاز داری — یکتایی (IF NOT EXISTS) یا آپدیتِ شرطی (IF balance = 100). Paxos را میانِ replicaهای یک پارتیشنِ واحد اجرا می‌کند (بدونِ تراکنشِ بین‌پارتیشنی)، سازگاریِ سریال می‌دهد، اما ~۳ تا ۴ رفت‌وبرگشت هزینه دارد — اغلب ۱۰ برابرِ نوشتنِ عادی. برای خواندنِ مقدارِ در حالِ پرواز باید با SERIAL بخوانی. فقط برای سطرهای نادری که نیاز دارند به کار ببر.

۹. چرا Elasticsearch از TF-IDF به BM25 کوچید؟

BM25 اشباعِ فراوانیِ ترم (بازدهِ نزولی با k1، پس keyword-stuffing دیگر کمک نمی‌کند) و نرمال‌سازیِ طولِ فیلدِ بهتر (b) را می‌افزاید. روی سندهای با طولِ متغیر مقاوم‌تر است و مرتبط‌سازیِ بهتری بدونِ رشدِ بی‌کرانِ TFِ کلاسیک تولید می‌کند.

۱۰. چرا نمی‌توان تعدادِ primary shard را بعد از ساختِ ایندکس تغییر داد؟

سندها با hash(_routing) % number_of_primary_shards مسیریابی می‌شوند. تغییرِ تعدادِ شارد، محلِ بایدیِ هر سند را عوض می‌کند و مسیریابی را بی‌اعتبار می‌کند. باید reindex کنی (یا از APIهای split/shrink که محدودند). در مقابل، تعدادِ replica به‌صورتِ زنده تغییرپذیر است.

۱۱. (کد / باگ را پیدا کن) این کدِ «cache-aside» زیر بار stampede می‌سازد. چرا؟
String v = redis.get(key);
if (v == null) {
    v = db.load(key);            // ۵۰۰۰ miss هم‌زمان همگی به DB می‌کوبند
    redis.set(key, v, 300);      // همه بعداً در همان ثانیه منقضی می‌شوند
}

دو باگ: (الف) بدونِ کشِ منفی → کلیدهای ناموجود همیشه miss می‌شوند (penetration)؛ (ب) TTLِ ثابتِ یکسان → همهٔ کپی‌ها با هم منقضی می‌شوند (avalanche)، و روی کلیدِ داغ همهٔ درخواست‌ها هم‌زمان بازسازی می‌کنند (breakdown). رفع: null را کش کن، TTL را jitter بده، و بازسازی را با mutex/singleflight به‌ازای هر کلید محافظت کن.

۱۲. (تله) `status: "active"` را به‌صورتِ فیلدِ `text` در Elasticsearch ذخیره می‌کنی و تجمیعِ `terms` چیزی/آشغال برمی‌گرداند. چرا؟

فیلدهای text تحلیل می‌شوند و (به‌طورِ پیش‌فرض) doc values ندارند، و فرمِ توکنایز‌شده‌شان آن چیزی نیست که تجمیع/فیلتر کنی. برای فیلترِ تطبیق‌دقیق، مرتب‌سازی و تجمیع به نوعِ keyword (یا زیرفیلدِ .keywordِ فیلدِ text با mappingِ پیش‌فرض) نیاز داری.

۱۳. (سخت) RDB در برابر AOF — یک Redisِ ۱۰۰ گیگابایتی داری و اسپایک‌های تأخیرِ دوره‌ای می‌بینی. چه خبر است؟

snapshotِ RDB و بازنویسیِ AOF فرایند را fork() می‌کنند؛ copy-on-write یعنی هر صفحه‌ای که والد بعداً تغییر دهد تکثیر می‌شود، بالقوه حافظه را دو برابر می‌کند و OS را روی خطاهای COW مشغول می‌کند → اسپایکِ تأخیر، و ریسکِ OOM اگر جای خالی نداشته باشی. کاهش: snapshot را خارج از پیک زمان‌بندی کن، ≥ ۲ برابرِ حافظهٔ اضافی داشته باش یا maxmemory محتاطانه، AOFِ everysec را ترجیح بده، و replica را برای بارِ snapshot در نظر بگیر.

۱۴. `allkeys-lru` دقیقاً چه می‌کند و آیا دقیق است؟

روی maxmemory، تقریباً کم‌استفاده‌ترین کلید را از میانِ همهٔ کلیدها حذف می‌کند. تقریبی است — Redis به‌تعدادِ maxmemory-samples کلید را نمونه‌برداری می‌کند و بهترین نامزد را حذف می‌کند نه اینکه همه‌چیز را اسکن کند، کمی دقت را با حذفِ O(1) معامله می‌کند. allkeys-lfu (بر پایهٔ فراوانی) معمولاً برای بارِ کشِ چوله بهتر است.

۱۵. (سخت) Pub/Sub در برابر Streams در Redis — کِی انتخابِ Pub/Sub داده را از دست‌تان می‌دهد؟

Pub/Sub بفرست-و-فراموش‌کن بدونِ پایداری است: اگر مشترک آفلاین یا کند باشد، پیام‌ها دور ریخته و هرگز بازپخش نمی‌شوند. هر نیاز به at-least-once، گروهِ مصرف‌کننده، ack یا بازپخش به Streams (XADD/XREADGROUP/XACK) نیاز دارد. انتخابِ Pub/Sub برای صفِ کار، هنگامِ قطعی به‌آرامی پیام‌ها را از دست می‌دهد.

جمع‌بندی
  • هیچ‌کدام رقیبِ هم و رقیبِ RDBMS نیستند — چهار ابزارِ تخصصی برای چهار دردِ متفاوت‌اند، و منبعِ حقیقت تقریباً همیشه Postgres/MySQL/Oracle می‌ماند.
  • Redis: انبارِ ساختار-دادهٔ درون‌حافظه و تک‌رشته‌ای. ساختارِ درست را انتخاب کن، الگوی کش را نام ببر (cache-aside پیش‌فرض)، سه بلا را بشناس (penetration/avalanche/breakdown)، بادوامش را با RDB+AOF بساز، و بدان که قفلِ توزیع‌شده‌اش best-effort است و برای درستی fencing token می‌خواهد.
  • ClickHouse: OLAPِ ستونی روی MergeTree. ORDER BY ترتیبِ فیزیکی + ایندکسِ sparse (یک mark در هر granuleِ ۸۱۹۲ سطری)، دسته‌ای درج کن نه سطر‌به‌سطر، و هرگز برای OLTP به کارش نبر.
  • ScyllaDB/Cassandra: ستون‌پهن، بر اساسِ کوئری مدل کن. partition key توزیع را تعیین می‌کند، clustering key ترتیب را؛ سازگاریِ قوی وقتی R + W > RF؛ و LWT/Paxos گران است، کم استفاده کن.
  • Elasticsearch: جستجو روی ایندکسِ معکوسِ Lucene. تفاوتِ text و keyword را بدان، BM25 (اشباعِ TF + نرمالِ طول) جای TF-IDF را از ES5 گرفت، تعدادِ primary shard ثابت است، و منبعِ حقیقت نیست.
  • مهم‌ترین سیگنالِ سنیوری: دانستنِ اینکه هر ابزار کجا اشتباه است.

منابع

Let's build one picture first. Imagine a professional kitchen. You have a versatile stove that cooks almost anything — that's your relational database (Postgres, MySQL, Oracle): the source of truth for everything. But next to it sit specialized tools: a flash deep-fryer, a screaming-hot pizza oven, a chilled display fridge. None is a "better stove"; each does one specific job ten times better. Redis, ClickHouse, ScyllaDB, and Elasticsearch are exactly those specialized tools beside the stove.

Roadmap for this lesson

In this chapter you'll learn four specialized datastores from the ground up:

  1. Redis — an in-memory data-structure store: caching, queues, locks, leaderboards.
  2. ClickHouse — a columnar analytical (OLAP) database on the MergeTree engine.
  3. ScyllaDB / Cassandra — a wide-column store with huge write throughput and tunable consistency.
  4. Elasticsearch — a text-search engine over an inverted index.

For each you'll learn three things: its data model (its shape), when it's right, and when it's a disaster. That third one is the biggest senior signal of all.

Part 0 — words you must know

Let's unpack a few terms from scratch so I never drop them cold later.

OLTP vs OLAP: the cashier vs the accountant

Picture a store. The cashier does thousands of tiny transactions a day: "sell this one item, decrement stock for this one product." Fast, single-row, precise. That's OLTP (Online Transaction Processing) — the everyday work of a relational database.

At month's end the accountant arrives and asks: "What was total sales per branch per month last year?" She doesn't change a row; she reads and sums millions of rows. That's OLAP (Online Analytical Processing) — ClickHouse's job.

The very database that's great for the cashier is slow for the accountant, and vice versa. This whole chapter revolves around that split.

  • ACID — the guarantee that a transaction either fully happens or not at all (atomic), data never corrupts (consistent), concurrent transactions don't tangle (isolated), and once committed nothing is lost (durable). It's the relational database's gold standard; all four tools in this chapter drop part of it to gain speed.
  • system of record — the one store that holds the "official, correct version" of data. A user's money, an order, an identity. If somewhere else and this store disagree, this store wins.
  • in-memory — data lives in RAM, not on disk. That's why Redis is fast, and why a power cut can lose data (unless you plan for it).
  • cardinality — the number of distinct values in a column. A "gender" column has low cardinality (few values); an "email" column has high cardinality (near-unique). This matters later.
  • idempotent — an operation that gives the same result whether run once or twice. "Set value to 10" is idempotent; "increment value by one" is not.
None of these compete with each other

The common junior mistake is thinking you must pick "one" of these. No. These four are point solutions to problems Postgres/MySQL solve badly at scale. In a real architecture all four can sit beside the same RDBMS, each curing one specific pain.

The big picture: four tools, four jobs

Frame this table in your mind; the whole chapter expands it:

Store Shape Sweet spot Do not use for
Redis In-memory key→value (rich values) Cache, sessions, rate limits, queues, leaderboards System of record for durable money-critical data
ClickHouse Columnar, append-mostly OLAP: analytics, aggregations over billions of rows OLTP: point updates, transactions, high-concurrency single-row reads
ScyllaDB / Cassandra Wide-column, partitioned Write-heavy, horizontally scaled, known access patterns Ad-hoc queries, joins, strong multi-key transactions
Elasticsearch Inverted index (documents) Full-text search, relevance, log analytics Source of truth; anything needing ACID

A recurring senior signal: knowing where each is wrong. In Ali's world — high-traffic Spring Boot APIs at Navashgaran, Redis caching at Neshan — the RDBMS stays the system of record; these four are accelerators bolted around it.


Redis

Data model — why it's not "just key-value"

A gym locker room, but each locker has a shape

A plain HashMap is like a row of identical lockers: each key is a box you drop a single plain string into. Redis is different: each locker can have a different shape — one is a queue, one is a ranked set, one is a field→value ledger. You pick not just the data but the right structure for it. That choice is the optimization itself.

Redis is an in-memory data-structure server whose command execution is single-threaded. Single-threaded means only one command runs at any instant — and because no locking is needed, every command is naturally atomic. The value behind a key is not just a string; it's a typed structure:

  • String — bytes up to 512 MB; also used as counters (INCR) and bitmaps.
  • Hash — field→value map; store an object without creating N separate keys.
  • List — linked list; LPUSH on one end and BRPOP on the other gives you a queue.
  • Set / Sorted Set (ZSet) — uniqueness. A ZSet is a skiplist + hash and does ranked inserts in O(log N) → exactly what leaderboards and time-ordered indexes want.
  • Stream — append-only log with consumer groups (a mini-Kafka inside Redis).
  • Plus HyperLogLog (cardinality in ~12 KB), Geo, Bitfields, and Pub/Sub.
Pick the right structure and half the work is done

A leaderboard on a ZSet is one ZREVRANGE. Simulating it with strings — read all, sort, write back every time — is a performance disaster. In Redis, "which structure?" is the most important decision.

Caching patterns — the name you say in an interview matters

The clerk and the back stockroom

The store (cache) has a front shelf and a big back stockroom (the database). When a customer wants an item, you check the front shelf first; if it's not there you fetch it from the back and put a copy on the front shelf so next time is fast. That's exactly cache-aside.

  • Cache-aside (lazy loading) — the app checks the cache itself; on a miss it reads the DB and populates the cache. Simple and resilient (cache down ≠ app down), but the first request is slow and data can be stale. This is the default and what Ali used at Neshan.
  • Read-through — the cache library transparently loads from the DB on a miss for you.
  • Write-through — a write goes to the cache and the DB synchronously. Cache always fresh, writes slower.
  • Write-behind (write-back) — write to the cache, flush to the DB asynchronously. Fast writes, risk of data loss on a crash.

Now see cache-aside in Java. Watch the comments — two subtle tricks live in there that I'll explain next:

// Cache-aside with Spring Data Redis, plus TTL and null-caching to stop stampede/penetration
public Product getProduct(long id) {
    String key = "product:" + id;
    String cached = redis.opsForValue().get(key);
    if (cached != null) {
        return cached.equals("__NULL__") ? null           // negative cache: kills cache penetration
                                          : deserialize(cached);
    }
    Product p = repository.findById(id).orElse(null);
    // jittered TTL avoids a synchronized avalanche (cache stampede)
    long ttl = 300 + ThreadLocalRandom.current().nextLong(60);
    redis.opsForValue().set(key,
            p == null ? "__NULL__" : serialize(p),
            Duration.ofSeconds(p == null ? 30 : ttl));
    return p;
}

The three classic cache disasters

You must know these three by name; interviewers love them.

The back door, the synchronized closing, and the rush on the one counter
  • Cache penetration: someone keeps asking for a product that doesn't exist. The front shelf is always empty, so every time you trek to the back and return empty-handed. An attacker can use this to crush the DB. Fix: cache the "not found" itself (__NULL__) or use a Bloom filter.
  • Cache avalanche: you stocked every shelf at once on New Year's Eve, so they all go empty at the same moment a year later and everything stampedes the stockroom simultaneously. Fix: jitter the TTLs so expirations spread out.
  • Cache breakdown (hot-key stampede): one wildly popular item runs out and a thousand customers want it at that instant; a thousand people go to the back to rebuild it at once. Fix: a mutex/singleflight so one thread rebuilds while the rest wait or serve stale data.

Now look at the code again: __NULL__ is the negative cache (fixes penetration), and ttl = 300 + random(60) is the jitter (fixes avalanche). These are production tricks, not decoration.

Eviction — when memory fills up

RAM isn't infinite. When usage hits maxmemory, the maxmemory-policy setting decides what gets thrown out:

  • noeviction (default) — nothing more is evicted; new writes error out.
  • allkeys-lru / allkeys-lfu — evict the least recently / least frequently used, across all keys.
  • volatile-lru / volatile-lfu / volatile-ttl — same, but only among keys that have a TTL.
  • allkeys-random / volatile-random — random.
Redis LRU is approximate, not exact

LRU means "least recently used (by time)" and LFU means "least frequently used (by count)." But Redis doesn't scan every key to find them exactly — it samples a few keys (maxmemory-samples) and evicts the best candidate among those. It trades a little accuracy for speed. For a pure cache, allkeys-lru is right; and since Redis 4, allkeys-lfu is usually better for skewed access (when a few keys are much hotter).

Persistence: RDB vs AOF

The family photo vs the diary

Redis is in-memory, but it has two ways to not reset to zero on restart:

  • RDB is like taking a family photo every few hours: a complete image of one moment. If something happens between two photos, you don't have those moments.
  • AOF is like keeping a diary: you jot down every action the instant you do it. By replaying the diary you can rebuild everything.
  • RDB — a periodic point-in-time binary snapshot. Redis fork()s itself for this (creates a child copy of the process) and uses copy-on-write. It's compact and restarts fast, but on a crash you lose everything since the last snapshot.
  • AOF — an append-only log of write commands. appendfsync everysec (default) loses at most ~1s; always is durable but slow. AOF is rewritten/compacted periodically.
  • Best practice: both — RDB for fast restore, AOF for durability. Redis 7 added multi-part AOF (a base RDB + an incremental AOF).
"Redis is just a cache so it's not durable" — this is wrong

Redis can be durable. But the real trap is elsewhere: that fork() for RDB or AOF-rewrite can, on a large dataset, nearly double memory (via copy-on-write) and cause latency spikes. So durability isn't free.

Pub/Sub and Streams

Live radio vs a mailbox

PUBLISH/SUBSCRIBE is like live radio: fire-and-forget. If your radio isn't on at that moment, that song is gone forever — not recorded, not replayed. Streams are like a mailbox: letters wait for you to pick them up, you get a receipt (XACK), and you can re-read them if needed.

So PUBLISH/SUBSCRIBE has no persistence and no delivery guarantee; an offline subscriber misses messages. For durable messaging use Streams (XADD/XREADGROUP) with consumer groups, acks (XACK), and replay.

Distributed locks — and what separates seniors from juniors

SET key value NX PX 30000 gives you a lock: set the value only if it doesn't exist (NX), with a 30-second expiry (PX). You must release it carefully too — with a Lua script that says "delete only if the token is yours":

// Safe release: check ownership and delete atomically — never a naive DEL
String lua = "if redis.call('get', KEYS[1]) == ARGV[1] " +
             "then return redis.call('del', KEYS[1]) else return 0 end";
redis.execute(new DefaultRedisScript<>(lua, Long.class), List.of(key), token);

But there are two subtleties that senior interviews put their finger on precisely:

  1. A single Redis is a SPOF and not linearizable across failover. SPOF means "single point of failure" — if it dies, everything dies. And because replication to a replica is asynchronous, a lock granted on the master can vanish if the master fails over before relaying it to the replica.
  2. Redlock (acquiring the lock from a majority of N independent masters) exists, but Martin Kleppmann's critique stands: Redlock relies on the assumption that clock drift and GC/pause stalls are bounded. Now suppose a process, due to a long GC pause, stalls past the TTL: its key expires and someone else grabs the lock, while the first one still thinks it holds it. Now two are in the critical section.
Fencing token: the only path to correctness

For real correctness you need a fencing token: a monotonically increasing number that goes up each time the lock is granted, which the protected resource itself checks, rejecting any request bearing an older token. With this, even if two think they hold the lock, only the one with the freshest token is allowed to write. The mental takeaway: a Redis lock is best-effort mutual exclusion, not a correctness guarantee. Use it to reduce duplicate work, not to guarantee it never happens.


ClickHouse

Data model, and why it's so fast

A row-wise library vs a column-wise library

A normal database lays data out row by row on disk: all of one customer's info together. If you want "the average age of all customers," you're forced to read everyone's whole record — name, address, phone, everything — just to reach the "age" column. ClickHouse flips it: it lays data out column by column — all ages together, all addresses together. Now for the average age you read only the "age" column file and nothing else.

ClickHouse is a columnar OLAP database. For an analytical query touching only 3 of 200 columns, it reads only those 3 columns' files — no wasted I/O. Second win: when all of a column's values are the same type, they compress brutally well (LZ4/ZSTD, and delta / double-delta for timestamps). Third win: the vectorized engine processes data in cache-friendly blocks. Those three together are why aggregations over billions of rows return in milliseconds.

The MergeTree family

MergeTree is the core engine; almost everything derives from it. See the code first, then we'll unpack it line by line:

CREATE TABLE events (
    event_date  Date,
    user_id     UInt64,
    event_type  LowCardinality(String),
    ts          DateTime,
    payload     String
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_date)          -- coarse partitions (for TTL/drops), NOT an index
ORDER BY (event_type, user_id, ts)         -- the SORTING key = physical order + sparse index
SETTINGS index_granularity = 8192;         -- rows per granule (mark)
A book index with one entry per 8192 pages

A normal B-tree is like an index with an entry for every page — precise, but if the book has billions of pages, the index itself becomes gigantic. ClickHouse builds an index with just one entry per 8192 pages; each 8192-row block is called a granule. This "sparse" index is so small it fits in RAM even for billions of rows. When you search, the index says "it's around this granule," and ClickHouse reads only that 8192-row block instead of the whole table.

Verified current behavior (ClickHouse docs, 2025):

  • ORDER BY is the sorting key: it defines the actual physical on-disk order of rows within each part.
  • The primary key is a sparse index, stored separately, with one entry per granule — and the default granule is 8192 rows. So the index has ~1 entry per 8192 rows. This is completely unlike a B-tree that indexes every row.
  • PRIMARY KEY can differ from ORDER BY only if it is a prefix of the ORDER BY tuple. That lets you keep the in-RAM index small while sorting by more columns.
  • Adaptive granularity is on by default via index_granularity_bytes (~10 MiB): a granule ends at 8192 rows or ~10 MiB, whichever comes first — protecting memory when rows are large.
  • Queries whose WHERE matches the sort-key prefix use the sparse index to skip whole granules; otherwise ClickHouse does a full scan (still fast, but no pruning).
Parts are like geological layers; merges compact them

Every INSERT creates a new immutable part — a small sorted mini-table, like a fresh sediment layer. A background process continually merges these layers into bigger ones — that's where the engine's name comes from: MergeTree. Here's the catch: if you do thousands of tiny inserts (row by row), you create thousands of thin layers and the merge process falls behind, giving you the "too many parts" error. That's why you must insert in large batches.

Key family members:

  • ReplacingMergeTree — dedups rows with the same sort key at merge time (eventually; not guaranteed until an OPTIMIZE ... FINAL or query-time FINAL).
  • SummingMergeTree / AggregatingMergeTree — pre-aggregate on merge.
  • CollapsingMergeTree / VersionedCollapsingMergeTree — model updates/deletes via +1/−1 sign rows.
  • ReplicatedMergeTree — replication via ZooKeeper/ClickHouse Keeper.

And skip indexes (minmax, set, bloom_filter) with their own GRANULARITY let you skip granules on columns that are not the sort key.

When ClickHouse is wrong

Don't use ClickHouse for OLTP — this is the classic mistake
  • No real UPDATE/DELETE semantics. Updates (ALTER TABLE ... UPDATE), called mutations, rewrite whole parts asynchronously — a heavy admin operation, not a transactional row edit.
  • No multi-statement ACID transactions, weak single-row consistency, eventual dedup.
  • Point lookups for one row by primary key are inefficient versus an OLTP store — it's built to scan ranges, not fetch one row.
  • Its workload is "a few heavy analytical queries," not "thousands of small concurrent queries."

The right architecture shape: OLTP in MySQL/Oracle (source of truth), then stream/ETL into ClickHouse for dashboards and analytics. Never point your write path at it for individual record edits.


ScyllaDB / Cassandra

ScyllaDB is a C++ rewrite of Cassandra: same CQL, same data model, but with a shard-per-core, "shard-aware" architecture and no JVM GC pauses. Everything I say about the data model applies to both.

The wide-column data model — model by query, not by entity

A fast-food kitchen that knows the menu in advance

In a normal restaurant (relational DB) you keep ingredients normalized and separate, assembling each order on the fly with joins. A high-traffic fast-food place doesn't work that way: the menu is fixed, so it keeps each item pre-made and packaged so an order goes out in a second. ScyllaDB is that: you know your queries in advance and build one table per query, even if data gets duplicated.

So you model by query, not by normalizing entities. The primary key has two parts:

CREATE TABLE messages_by_room (
    room_id    uuid,
    bucket     int,
    ts         timeuuid,
    user_id    uuid,
    body       text,
    PRIMARY KEY ((room_id, bucket), ts, user_id)  -- (partition key)(clustering keys)
) WITH CLUSTERING ORDER BY (ts DESC);
  • Partition key ((room_id, bucket)) — hashed to decide which node owns the data. All rows of a partition live together on the same replicas. This is the unit of distribution and of single-partition atomicity.
  • Clustering keys (ts, user_id) — define the sort order within a partition on disk, making range scans inside a partition cheap.

The golden rules:

  1. Every query must hit one partition by its full partition key. A query without it forces a scatter-gather across all nodes — its telltale sign is ALLOW FILTERING, a red flag in production.
  2. Avoid hot partitions (traffic skewed onto one partition) and unbounded partitions (a partition that grows forever) — that's why we added the bucket: to shard a busy room across several partitions.
  3. Denormalize: build one table per access pattern and accept duplicate data.

Tunable consistency

A vote among equal copies

In Cassandra/Scylla there's no single leader; every replica is equal and data is copied RF times (replication factor, e.g. 3 copies). For each query, you decide how many copies must agree — that's the consistency level (CL). It's like a vote: the more votes you demand, the surer you are but the slower you go.

  • Writes are sent to all replicas; the write succeeds once CL acknowledge (ONE, QUORUM, ALL, LOCAL_QUORUM…).
  • Reads contact CL replicas and reconcile the answers.
  • Strong consistency holds when R + W > RF. For example with RF=3, QUORUM reads (R=2) + QUORUM writes (W=2) gives 2+2 > 3, so any read is guaranteed to see the latest write. If that inequality doesn't hold, you get eventual consistency and may read stale data.
  • LOCAL_QUORUM keeps the vote within one datacenter for latency (in multi-DC deployments).

Consistency is repaired lazily: via read repair (during reads), hinted handoff (holding a write for a node that's temporarily down), and anti-entropy repair with nodetool repair.

Lightweight Transactions (LWT)

Default writes are last-write-wins with no conditionals. But sometimes you truly need compare-and-set — like "claim this username only if nobody has it." Here you use LWT with IF:

INSERT INTO users (username, id) VALUES ('ali', ...) IF NOT EXISTS;
UPDATE accounts SET balance = 90 WHERE id = ? IF balance = 100;
LWT is powerful but costly — verified (ScyllaDB docs, 2025)
  • LWT uses Paxos for consensus among the replicas of a single partition — all conditions must target the same partition; there are no cross-partition transactions.
  • It gives serial consistency, but costs roughly three-to-four round trips (prepare/read/propose/commit) versus one for a normal write → often an order of magnitude slower. Use sparingly.
  • To read the latest value mid-transaction you must read at SERIAL/LOCAL_SERIAL consistency; a normal QUORUM read may miss an in-flight Paxos value.

When it's wrong

No joins, no ad-hoc WHERE on arbitrary columns, no strong multi-partition transactions, and it's painful for evolving/unknown query patterns. Rule of thumb: if you can't enumerate your queries up front, this is the wrong database.


Elasticsearch

The inverted index — the heart of search

The index at the back of a book

At the back of a textbook is the index: for each important word, the list of pages it appears on. You don't flip through the whole book looking for "migration"; you go to the index and jump straight to the pages. An inverted index is exactly that: a mapping of term → list of documents containing that term.

A relational index maps a row to its column values; an inverted index does the opposite — each term → list of documents (postings) containing it, with positions and frequencies. That's what makes "find every document containing migration" roughly O(1) instead of a full scan. Elasticsearch itself is a distributed layer over Apache Lucene; Lucene owns the index structures (segments, postings, doc values).

Analyzers

A text-shredding factory

Before text enters the index, it passes through a production line: character filters → tokenizer → token filters. For instance "The Migrations!" goes in, gets lowercased, split on non-letters, stopwords (like "the") stripped, and stemmed down to [migrat]. The crucial point: the same production line must run at query time, or your terms won't match the stored ones.

That's why you have two field types: the unanalyzed keyword field (exact, not tokenized) for filtering/aggregations/sorting, and the text field for full-text search. A common bug: expecting an exact match or aggregation on a text field — when you need the .keyword sub-field.

Relevance and BM25

Scoring the results fairly

When you search "database migration," Elasticsearch must decide which document ranks higher. It combines three simple intuitions: (1) a document with the word more often is better — but not without limit; (2) a rare word matters more than a common one; (3) a match in a short title is worth more than the same match buried in a long, wordy body. The model that blends these is called BM25.

Since Elasticsearch 5 / Lucene 6, the default relevance model is BM25 (replacing classic TF-IDF). BM25 scores a document for a query by combining:

  • Term frequency (TF) — but with saturation: the 10th occurrence of a word adds far less than the 2nd (parameter k1).
  • Inverse document frequency (IDF) — rare terms weigh more than common ones.
  • Field-length normalization — a match in a short title beats the same match buried in a long body (parameter b).

That saturation + length norm is exactly why BM25 beats naive TF-IDF (keyword stuffing no longer scores without bound), and a classic senior question is "why did they switch."

Shards and replicas

A library with several branches and several copies per branch

A big index can't fit on one machine, so you split it into pieces (primary shards) and put each on a machine — that's horizontal scale and parallelism. Then you make several copies (replicas) of each piece for high availability and more read throughput.

  • An index is split into primary shards (each a self-contained Lucene index) — the unit of horizontal scale and parallelism. Primary shard count is fixed at creation (you must reindex/split to change it).
  • Each primary has replica shards — copies for HA and read throughput; replica count is changeable live.
  • A document routes to a shard by hash(routing) % number_of_primary_shards — this is why the primary count is immutable.
  • Lucene segments are immutable; writes first go to an in-memory buffer + a translog (durability), are then refreshed into new searchable segments (default every 1s → "near real-time"), and background-merged. Deletes are tombstones reclaimed on merge.
Oversharding: too many small shards

Too many small shards wastes heap and cluster state. Aim for shards in the tens-of-GB range, not hundreds of few-MB shards.

When it's wrong

Elasticsearch is not a system of record: no ACID transactions, no joins (only limited parent/child and nested), and it's near-real-time, not immediately consistent. Use it as a search/analytics index fed from your primary DB (or Kafka/CDC), never as the sole store of truth for critical data.


Choosing between them (the real interview question)

If you remember only one thing from this chapter, make it this decision table:

  • Need sub-ms reads of hot data / sessions / counters / locks → Redis.
  • Need GROUP BY over billions of rows for dashboards → ClickHouse.
  • Need massive write throughput with predictable, partition-scoped queries and no single point of failure → ScyllaDB/Cassandra.
  • Need "search by words," typo tolerance, relevance ranking, log exploration → Elasticsearch.
  • Your source of truth for transactional, relational, money-critical data almost always stays Postgres/MySQL/Oracle — and these four are specialized read/scale layers around it, kept in sync by ETL, CDC, or dual writes.

Interview Questions

1. Redis is single-threaded — how does it serve 100k+ ops/s, and what breaks that?

Command execution is single-threaded, so operations are atomic and there's no locking overhead; the work is CPU-cheap and in-memory. It scales via I/O multiplexing (epoll) and, since Redis 6, multithreaded I/O for socket read/write (parsing stays single-threaded). What breaks it: a single O(N) command (KEYS *, big SMEMBERS, SORT) blocks everything. Use SCAN and avoid large blocking commands.

2. (Hard) Your Redis distributed lock occasionally lets two workers into the critical section. Why?

Several reasons: (a) master failover before the lock replicated to the replica; (b) the holder's process paused (GC/CPU starvation) past the TTL, so the key expired and another client acquired it while the first thinks it holds it; (c) a naive DEL deleting someone else's lock after your TTL expired. Fixes: unique token + Lua compare-and-delete, and for correctness a fencing token the resource validates. Per Kleppmann, Redlock is not safe for correctness under unbounded pauses — it's best-effort.

3. Difference between `ORDER BY` and `PRIMARY KEY` in ClickHouse MergeTree?

ORDER BY is the sorting key: physical on-disk row order + the basis of the sparse primary index. PRIMARY KEY, if specified separately, must be a prefix of ORDER BY; it lets you shrink the in-RAM index while still sorting by extra columns. The index stores one mark per granule (8192 rows default), not per row.

4. (Hard) You insert 5,000 rows/sec one-by-one into ClickHouse and it falls over. Why, and the fix?

Every insert creates a new immutable part; the background merger can't keep up → "too many parts" and stalls. ClickHouse is built for batched inserts. Fix: buffer and insert in large blocks (tens of thousands of rows), use async inserts, or a Buffer table / Kafka engine. This is the single most common ClickHouse production mistake.

5. Why is ClickHouse fast for OLAP but wrong for OLTP?

Columnar storage reads only referenced columns and compresses them heavily; the vectorized engine scans ranges via a sparse index. But there are no transactional row updates (mutations rewrite parts asynchronously), no ACID, and point single-row lookups are inefficient. It's built to scan and aggregate, not to edit and fetch one row.

6. In Cassandra/Scylla, what's the difference between partition key and clustering key, and why does it dominate your schema?

The partition key is hashed to place data on nodes and defines the unit of co-location and atomicity; the clustering key sorts rows within a partition on disk. Because efficient queries must specify the full partition key, you design a table per query ("query-first modeling") and denormalize — the keys are the data model.

7. When is a Cassandra read strongly consistent?

When R + W > RF. With RF=3, QUORUM writes (W=2) + QUORUM reads (R=2) give 4 > 3 → any read sees the latest committed write. Lower CLs give eventual consistency and possible stale reads; LOCAL_QUORUM keeps latency within one datacenter.

8. (Hard) When would you use LWT in Scylla, and what's the cost?

When you need compare-and-set semantics — uniqueness (IF NOT EXISTS) or conditional update (IF balance = 100). It runs Paxos among the replicas of a single partition (no cross-partition transactions), gives serial consistency, but costs ~3–4 round trips — often 10× a normal write. To read an in-flight value you must read at SERIAL. Use it only for the rare rows that need it.

9. Why did Elasticsearch switch from TF-IDF to BM25?

BM25 adds term-frequency saturation (diminishing returns via k1, so keyword stuffing stops helping) and better field-length normalization (b). It's more robust on documents of varying length and produces better relevance without the unbounded TF growth of classic TF-IDF.

10. Why can't you change the number of primary shards after index creation?

Documents route by hash(_routing) % number_of_primary_shards. Changing the shard count changes where every document should live, invalidating routing. You must reindex (or use the split/shrink APIs, which are constrained). Replica count, by contrast, is changeable live.

11. (Code / find-the-bug) This "cache-aside" code causes a stampede under load. Why?
String v = redis.get(key);
if (v == null) {
    v = db.load(key);            // 5,000 concurrent misses all hit the DB
    redis.set(key, v, 300);      // all expire at the same second later
}

Two bugs: (a) no negative caching → nonexistent keys always miss (penetration); (b) identical fixed TTL → all copies expire together (avalanche), and on a hot key all requests rebuild simultaneously (breakdown). Fix: cache nulls, jitter the TTL, and guard rebuild with a per-key mutex/singleflight.

12. (Gotcha) You store `status: "active"` as a `text` field in Elasticsearch and a `terms` aggregation returns nothing/garbage. Why?

text fields are analyzed and (by default) don't have doc values, and the tokenized form isn't what you aggregate/filter on. You need a keyword type (or the .keyword sub-field of a text field with the default mapping) for exact-match filters, sorting, and aggregations.

13. (Hard) RDB vs AOF — you have a 100 GB Redis and see periodic latency spikes. What's happening?

RDB snapshotting and AOF rewrite fork() the process; copy-on-write means every page the parent then modifies gets duplicated, potentially doubling memory and causing the OS to spend time on COW faults → latency spikes, and OOM risk if you lack headroom. Mitigate: schedule snapshots off-peak, ensure ≥ 2× memory headroom or use maxmemory conservatively, prefer AOF everysec, and consider replicas taking the snapshot load.

14. What does `allkeys-lru` actually do, and is it exact?

On maxmemory, evict the approximately least-recently-used key across all keys. It's approximate — Redis samples maxmemory-samples keys and evicts the best candidate rather than scanning everything, trading a little accuracy for O(1) eviction. allkeys-lfu (frequency-based) is usually better for skewed cache workloads.

15. (Hard) Pub/Sub vs Streams in Redis — when does choosing Pub/Sub lose you data?

Pub/Sub is fire-and-forget with no persistence: if a subscriber is offline or slow, messages are dropped and never replayed. Any at-least-once requirement, consumer groups, acknowledgements, or replay needs Streams (XADD/XREADGROUP/XACK). Choosing Pub/Sub for a job queue silently loses messages on disconnects.

In a nutshell
  • None competes with the others or with the RDBMS — they're four specialized tools for four different pains, and the source of truth almost always stays Postgres/MySQL/Oracle.
  • Redis: an in-memory, single-threaded data-structure store. Pick the right structure, name your cache pattern (cache-aside is the default), know the three disasters (penetration/avalanche/breakdown), make it durable with RDB+AOF, and know its distributed lock is best-effort and needs a fencing token for correctness.
  • ClickHouse: columnar OLAP on MergeTree. ORDER BY is physical order + sparse index (one mark per 8192-row granule), insert in batches not row-by-row, and never use it for OLTP.
  • ScyllaDB/Cassandra: wide-column, model by query. The partition key sets distribution, the clustering key sets order; strong consistency when R + W > RF; and LWT/Paxos is costly, so use it sparingly.
  • Elasticsearch: search over Lucene's inverted index. Know text vs keyword, BM25 (TF saturation + length norm) replaced TF-IDF since ES5, primary shard count is fixed, and it's not a system of record.
  • The biggest senior signal: knowing where each tool is wrong.

Sources