System Design শেখো
সব কেস স্টাডি

Top-K / Trending System ডিজাইন

15 মিনিটadvanced
এক নজরে
  • বিশাল event stream থেকে real-time-এ সবচেয়ে জনপ্রিয় K-টি item (top hashtag, top video) বের করা।
  • exact counting স্কেলে ব্যয়বহুল, তাই Count-Min Sketch ও HyperLogLog-এর মতো approximate counting দিয়ে memory বাঁচানো হয়।
  • stream দিয়ে দ্রুত আনুমানিক ফল, batch দিয়ে নির্ভুল ফল — Lambda architecture-এ দুটোই মেলানো হয়।

Twitter-এর trending hashtag বা YouTube-এর "top trending videos" — এদের পিছনে একটা মজার সমস্যা আছে: প্রতি সেকেন্ডে লাখো event আসছে, আর তার মধ্য থেকে real-time-এ সবচেয়ে জনপ্রিয় K-টি item বের করতে হবে। সব কিছু exact গুনতে গেলে memory আর compute দুটোই ধসে পড়ে। চলো ধাপে ধাপে ডিজাইন করি।

১. সমস্যা বোঝা (Requirements)

Functional requirements:

  • নির্দিষ্ট time window-এর জন্য top-K item ফেরত দিতে হবে (যেমন last 1 hour-এর top 10 hashtag)।
  • একাধিক window সাপোর্ট — last 5 min, last 1 hour, last 24 hours।
  • real-time-এর কাছাকাছি আপডেট (কয়েক সেকেন্ড দেরি গ্রহণযোগ্য)।

Non-functional requirements:

  • Scalability — প্রতি সেকেন্ডে লাখো-কোটি event।
  • Low latency — top-K query দ্রুত (মিলিসেকেন্ড)।
  • Memory efficiency — কোটি কোটি unique item-এর জন্য exact counter রাখা অসম্ভব।
  • Approximation চলবে — trending list-এ ৩ নম্বর আর ৪ নম্বরে সামান্য এদিক-ওদিক হলে কেউ লক্ষ্য করবে না।
সহজ উদাহরণ

ভাবো একটা বিশাল ক্রিকেট স্টেডিয়ামে কোন স্লোগান সবচেয়ে বেশি হচ্ছে তা জানতে চাও। প্রতিটা দর্শকের প্রতিটা কথা আলাদা করে গোনা অসম্ভব। তার বদলে কয়েকজন volunteer স্ট্যান্ডের আলাদা আলাদা অংশে দাঁড়িয়ে আনুমানিক গোনে — "এদিকে বাংলাদেশ বেশি, ওদিকেও বাংলাদেশ"। নিখুঁত না হলেও কোনটা top সেটা ঠিকঠাক বেরিয়ে আসে। Count-Min Sketch ঠিক এভাবেই আনুমানিক গোনে।

২. স্কেল আন্দাজ (Estimation)

event rate (tweet/view)     = 1 million/sec (peak)
unique item (hashtag/video) = 100 million+

exact counting memory:
  100M item × (key ~50 B + counter 8 B) ≈ 5.8 GB শুধু counter-এ
  multiple window × multiple node = কয়েকশো GB → ব্যবহারিক নয়

Count-Min Sketch memory:
  width 2^20, depth 5, 4-byte cell
  = 5 × 1M × 4 B ≈ 20 MB মাত্র — যেকোনো unique সংখ্যার জন্য fixed!
  => 100M item হোক বা 1B, memory একই

query QPS (top-K পড়া)       ≈ মাঝারি (UI refresh) — cache করা যায়

মূল শিক্ষা: exact counting-এ memory unique-item সংখ্যার সাথে বাড়ে, কিন্তু sketch-এ memory fixed — এটাই approximate counting-এর মূল জয়।

৩. API ডিজাইন

POST /v1/events            # internal ingestion
     { item_id, type, timestamp }     # type = hashtag/video

GET  /v1/topk
     ?window=1h            # 5m | 1h | 24h
     &k=10
     &category=hashtag
→ 200
  {
    "window": "1h",
    "items": [
      { "item": "#Eid", "approx_count": 482000, "rank": 1 },
      { "item": "#Padma", "approx_count": 311000, "rank": 2 }
    ]
  }

event ingestion সাধারণত API না হয়ে সরাসরি Kafka-র মতো message queue-তে যায়; উপরের /events শুধু conceptual।

৪. ডেটা মডেল

স্টোরবিষয়বস্তুউদ্দেশ্য
Kafka topicraw event streamingestion buffer
Count-Min Sketch (in-memory)প্রতি window-এর approximate countদ্রুত heavy hitter
Min-Heap (size K)বর্তমান top-K itemদ্রুত top-K পড়া
Batch store (HDFS/S3 + warehouse)raw event logনির্ভুল recompute
Result cache (Redis)precomputed top-K per windowUI query

Choice: real-time path-এ in-memory probabilistic structure (CMS + heap), নির্ভুলতার জন্য batch path-এ raw log। ফলাফল Redis-এ cache করে query সস্তা রাখা হয়।

৫. হাই-লেভেল ডিজাইন

components:

  • Ingestion → Kafka (partitioned by item hash)
  • Stream Processor (Flink/Spark Streaming) — CMS update + heap maintain
  • Aggregator — partition-wise top-K merge করে global top-K
  • Batch Processor — raw log থেকে নির্ভুল top-K (দৈনিক/ঘণ্টায়)
  • Result Cache (Redis) + Query Service

Flow:

  1. event Kafka-তে যায়, item hash দিয়ে partition হয় (একই item সবসময় একই partition-এ)।
  2. প্রতিটি stream processor তার partition-এর জন্য Count-Min Sketch update করে এবং একটি local min-heap (size K) maintain করে।
  3. Aggregator সব partition-এর local top-K নিয়ে merge করে global top-K বের করে।
  4. ফল Redis-এ লেখা হয়; Query Service সেখান থেকে পড়ে।
  5. সমান্তরালে batch layer raw log-এ নির্ভুল হিসাব করে real-time ফলকে সংশোধন করে।

৬. গভীরে (Deep Dive)

Count-Min Sketch — heavy hitters

Count-Min Sketch হলো একটি 2D array (depth rows × width columns) এবং depthটি ভিন্ন hash function। কোনো item দেখলে প্রতিটি row-এ তার hash অনুযায়ী একটি cell +1 হয়। count জানতে চাইলে সব row-এর cell-গুলোর minimum নাও।

add(item):
  for i in 0..depth:
     sketch[i][ hash_i(item) % width ] += 1

count(item):
  return min over i of sketch[i][ hash_i(item) % width ]

কেন min? কারণ hash collision-এ count শুধু বাড়তে পারে (অন্য item একই cell-এ পড়লে), কখনো কমে না — তাই সব estimate-এর minimum-টাই সত্যের সবচেয়ে কাছাকাছি (overestimate, কখনো underestimate নয়)। জনপ্রিয় item (heavy hitter)-এর count বড়, তাই সামান্য collision-error তাদের ranking-এ প্রায় প্রভাব ফেলে না — ঠিক এই কারণেই trending-এর জন্য এটা পারফেক্ট।

Top-K বের করা: CMS + Min-Heap

শুধু CMS দিয়ে "কোন item-গুলো top" তা জানা যায় না (CMS শুধু count দেয়, item list রাখে না)। তাই সাথে size-K একটা min-heap রাখি:

  • নতুন item-এর CMS count নাও।
  • heap-এ K-এর কম থাকলে যোগ করো।
  • নয়তো heap-এর সবচেয়ে ছোট (root) count-এর সাথে তুলনা করো; বড় হলে root সরিয়ে নতুনটা ঢোকাও।

heap-এ সবসময় বর্তমান top-K candidate থাকে, top-K পড়া O(K)।

সাবধান

Count-Min Sketch কখনো count overestimate করে (hash collision-এর কারণে), কিন্তু underestimate করে না। তাই rare item-কে এটা বড় দেখাতে পারে। trending-এর জন্য সমস্যা নয় (আমরা heavy hitter-ই চাই), কিন্তু কেউ যদি এই count দিয়ে billing বা exact reporting করতে চায় — সর্বনাশ। approximate count-কে কখনো exact বলে চালিয়ে দিও না।

HyperLogLog — unique count

কখনো আমরা "কতবার" নয়, "কতজন unique ইউজার" জানতে চাই (যেমন কতজন distinct user video দেখেছে)। এর জন্য HyperLogLog। এটি প্রতিটি element-কে hash করে hash-এর leading zero-র সর্বোচ্চ সংখ্যা থেকে cardinality আনুমানিক করে — মাত্র কয়েক KB memory-তে কোটি কোটি unique element গোনে (~2% error)। distinct counting-এ এটা আর exact set রাখার মধ্যে memory-র পার্থক্য আকাশ-পাতাল।

Sliding Window

"last 1 hour"-এর trending মানে পুরোনো count ভুলে যেতে হবে। দুই কৌশল:

  1. Bucketed window: সময়কে ছোট bucket-এ ভাগ করি (যেমন প্রতি ১ মিনিটে একটি CMS)। "last 1 hour" = গত ৬০টি bucket-এর CMS যোগফল। নতুন bucket এলে সবচেয়ে পুরোনো bucket ফেলে দিই (count সরে যায়)।
  2. Decay: পুরোনো count-কে সময়ের সাথে exponentially কমিয়ে দিই, যাতে সাম্প্রতিক জিনিস বেশি ওজন পায়।

bucketed approach সহজ ও পরিষ্কার window দেয়, তাই সাধারণত এটাই বেছে নেওয়া হয়।

Batch + Stream: Lambda Architecture

stream layer দ্রুত কিন্তু আনুমানিক (CMS-এর error সহ); batch layer ধীর কিন্তু নির্ভুল (raw log scan করে exact top-K)। Lambda Architecture-এ দুটোই চলে — UI-তে সাধারণত stream-এর fresh result দেখানো হয়, আর batch periodically সেই result সংশোধন করে দেয়। এতে real-time freshness ও long-term accuracy দুটোই পাওয়া যায়।

৭. বটলনেক ও স্কেলিং

  • Partitioning: item hash দিয়ে Kafka partition — একই item একই processor-এ যায়, ফলে count বিভক্ত হয় না। তবে একটি viral hashtag এক partition-কে hot করে দিতে পারে; সমাধান — সেই key sub-partition করে aggregator-এ merge করা।
  • Aggregation overhead: প্রতিটি partition local top-K (যেমন top-2K) পাঠালে aggregator merge করে global top-K বের করে — সব item পাঠানোর দরকার নেই।
  • Memory: CMS-এর width/depth বাড়ালে accuracy বাড়ে কিন্তু memory বাড়ে — trade-off tune করতে হয়।
  • Multiple windows: প্রতিটি window-এর জন্য আলাদা bucketed sketch; পুরোনো bucket evict করে memory নিয়ন্ত্রণ।
  • Query scaling: top-K result Redis-এ cache, short TTL — query path DB/stream থেকে decouple।
  • Hot key skew: viral event-এ একটি item-এ চাপ; sub-partition + load-aware routing দিয়ে balance।

৮. সারসংক্ষেপ

Top-K/trending system-এর মূল অন্তর্দৃষ্টি — exact counting স্কেলে অসম্ভব, কিন্তু approximation-এ trade হারানো accuracy trending-এ কেউ টের পায় না। Count-Min Sketch fixed memory-তে heavy hitter গোনে, HyperLogLog unique count দেয়, min-heap top-K maintain করে, আর bucketed sliding window পুরোনো ডেটা ভুলিয়ে দেয়। সবশেষে Lambda architecture stream-এর গতি ও batch-এর নির্ভুলতা মিলিয়ে দেয়।

টিপস

Interview-এ প্রথমে "প্রতিটি item-এ exact counter + global heap" নিষ্পাপ সমাধান দাও, তারপর memory bottleneck দেখিয়ে Count-Min Sketch-এ evolve করো, এবং স্পষ্ট করে বলো এটি approximate (overestimate, never underestimate)। এই trade-off সচেতনতাই senior signal।

মিনি কুইজ

1. প্রতিটি hashtag-এর জন্য exact counter রাখার মূল সমস্যা কী, যখন কোটি কোটি ইউনিক item থাকে?

2. HyperLogLog কী আনুমানিক হিসাব করে?

3. Lambda architecture-এ stream layer-এর প্রধান ভূমিকা কী?