Top-K / Trending System ডিজাইন
- ●বিশাল 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 topic | raw event stream | ingestion 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 window | UI 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:
- event Kafka-তে যায়, item hash দিয়ে partition হয় (একই item সবসময় একই partition-এ)।
- প্রতিটি stream processor তার partition-এর জন্য Count-Min Sketch update করে এবং একটি local min-heap (size K) maintain করে।
- Aggregator সব partition-এর local top-K নিয়ে merge করে global top-K বের করে।
- ফল Redis-এ লেখা হয়; Query Service সেখান থেকে পড়ে।
- সমান্তরালে 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 ভুলে যেতে হবে। দুই কৌশল:
- Bucketed window: সময়কে ছোট bucket-এ ভাগ করি (যেমন প্রতি ১ মিনিটে একটি CMS)। "last 1 hour" = গত ৬০টি bucket-এর CMS যোগফল। নতুন bucket এলে সবচেয়ে পুরোনো bucket ফেলে দিই (count সরে যায়)।
- 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-এর প্রধান ভূমিকা কী?