Real-time Market Data Feed ডিজাইন
- ●Market data feed হলো price update গুলো millions of subscriber-এর কাছে কম latency-তে fan-out করার সিস্টেম।
- ●একবার snapshot পাঠিয়ে তারপর incremental update দিলে bandwidth বাঁচে; sequence number দিয়ে gap ধরা হয়।
- ●Slow client-দের জন্য conflation করা হয়—পুরনো price ফেলে শুধু সর্বশেষটা পাঠানো হয়।
স্টক এক্সচেঞ্জের matching engine যখন একটা trade করে বা order book বদলায়, তখন সেই খবর সঙ্গে সঙ্গে লাখো ট্রেডার, অ্যালগো বট আর মোবাইল অ্যাপের কাছে পৌঁছাতে হয়। এটাই market data feed—একটা বিশাল broadcast সমস্যা, যেখানে latency, ordering আর scale তিনটাই একসাথে সামলাতে হয়। চলো ডিজাইন করি।
১. সমস্যা বোঝা (Requirements)
Functional requirements:
- প্রতিটা symbol-এর জন্য real-time price update (last trade, best bid/ask, order book depth) পাঠানো।
- Client subscribe/unsubscribe করতে পারবে নির্দিষ্ট symbol-এ।
- নতুন client কানেক্ট করলে current state (snapshot) পাবে, তারপর live update।
Non-functional requirements:
- Scale: লাখো concurrent subscriber, সেকেন্ডে লাখো update।
- Latency: update generate হওয়া থেকে client-এ পৌঁছানো পর্যন্ত কয়েক ms (internal feed-এ microsecond)।
- Ordering: update গুলো client-এ একই ক্রমে পৌঁছাতে হবে; নাহলে stale price দেখা যাবে।
- No silent gap: কোনো update হারালে client যেন বুঝতে পারে (gap detection)।
- Fairness: কোনো client যেন systematically আগে data না পায়।
ভাবো একটা ক্রিকেট ম্যাচের live স্কোর। প্রতিবার পুরো scorecard পাঠানোর বদলে কমেন্টেটর শুধু পরিবর্তনটা বলেন—"চার!", "উইকেট!"। এটাই incremental update। আর কেউ মাঝপথে টিভি অন করলে আগে পুরো স্কোর (snapshot) দেখিয়ে দেওয়া হয়, তারপর ball-by-ball আপডেট।
২. স্কেল আন্দাজ (Estimation)
Update generation:
Symbols = 5,000
Updates/sec/symbol = 200 (peak)
Total updates/sec = 1,000,000
Fan-out:
Subscribers = 2,000,000 (mobile + web + algo)
Avg symbols/subscriber = 20
Naive unicast load = 1M updates * (matching subscribers)
→ বিলিয়ন messages/sec (অসম্ভব!)
→ তাই multicast (internal) + conflation (edge) লাগবেই
Bandwidth (incremental message ~ 40 bytes):
1M/sec * 40 B = 40 MB/sec per stream (multicast, একবার)
Full order book snapshot ~ 4 KB → কম, শুধু join-এর সময়
শিক্ষা: যদি প্রতি client-কে আলাদা copy পাঠাও, সংখ্যা বিস্ফোরিত হয়। তাই internal-এ one-to-many multicast আর edge-এ per-client conflation—এই দুই কৌশলই scaling-এর চাবি।
৩. API ডিজাইন
# Client-facing (WebSocket)
WS CONNECT /v1/marketdata
→ SUBSCRIBE { symbols: ["GP", "BEXIMCO"], depth: 5 }
→ UNSUBSCRIBE { symbols: ["GP"] }
# Server → Client messages
SNAPSHOT { symbol, seq_no, bids:[...], asks:[...], last_trade }
DELTA { symbol, seq_no, side, price, qty, action: ADD|UPDATE|DELETE }
HEARTBEAT { ts, last_seq_no } # gap detect + liveness
GAP_FILL { symbol, from_seq, to_seq } # client recovery অনুরোধ করতে পারে
# Internal feed (UDP multicast)
group 239.1.1.10:5000 → incremental stream
group 239.1.1.11:5001 → periodic snapshot stream
প্রতিটা message-এ seq_no থাকা বাধ্যতামূলক—এটাই ordering আর gap detection-এর ভিত্তি।
৪. ডেটা মডেল
Market data বেশিরভাগই in-memory ও streaming, কিন্তু কিছু state রাখতে হয়:
| Store | মূল ফিল্ড | উদ্দেশ্য |
|---|---|---|
book_state (in-mem) | symbol, bids, asks, last_seq_no | প্রতিটা symbol-এর current snapshot তৈরির উৎস |
subscription (in-mem) | session_id, symbols[], conn | কে কোন symbol-এ subscribed |
tick_history (TSDB) | symbol, ts, price, qty | historical/replay, charting |
audit_log | seq_no, payload, ts | regulatory record |
কেন এখানে strong consistency-র চেয়ে ordering জরুরি? Market data-তে আমরা পুরো ডাটাবেস-style consistency চাই না—আমরা চাই প্রতিটা client একই ordered stream পায় এবং কোনো gap থাকলে টের পায়। তাই canonical "truth" হলো sequenced stream নিজেই; in-memory book_state শুধু সেই stream থেকে derived। Historical data রাখা হয় একটা time-series database-এ যেখানে read-heavy charting query দ্রুত চলে।
৫. হাই-লেভেল ডিজাইন
ধাপে ধাপে data-র যাত্রা:
- Matching engine প্রতিটা order book পরিবর্তনে একটা sequenced update event বের করে।
- Feed handler / normalizer raw event নিয়ে একটা compact wire format (incremental message) বানায়,
seq_noবসায়। - Internal multicast bus: এই incremental stream একটা multicast group-এ একবার publish হয়। সব downstream consumer (gateway, OMS, surveillance) একই group থেকে পড়ে।
- Snapshot service: একটা component multicast stream consume করে in-memory order book maintain করে, এবং নিয়মিত interval-এ full snapshot publish করে আরেকটা group-এ।
- Edge gateways (WebSocket fan-out): এরা multicast থেকে পড়ে, লাখো client-এর কাছে unicast WebSocket-এ পাঠায়। এখানেই conflation ও per-client queue থাকে।
- Client: join করার সময় আগে latest snapshot নেয়, তার
seq_noমনে রাখে, তারপর delta apply করে। delta-রseq_no-তে gap পেলে gap-fill চায় বা snapshot পুনরায় নেয়।
৬. গভীরে (Deep Dive)
Snapshot + Incremental ও gap recovery
একজন client কখন কীভাবে synced হয়?
- Join: snapshot নাও (ধরো
seq_no = 1000)। এর মানে এই snapshot 1000 পর্যন্ত সব update-এর ফল। - এরপর শুধু
seq_no > 1000delta apply করো। 1001, 1002, 1003... ক্রমানুসারে। - যদি 1002-এর পর সরাসরি 1004 আসে—gap! 1003 হারিয়েছে। তখন দুটো অপশন: (ক) gap-fill request করে শুধু 1003 আবার নাও, অথবা (খ) আবার fresh snapshot নাও আর সেখান থেকে শুরু করো।
এই sequence number-ই হলো নিঃশব্দ data corruption থেকে বাঁচার মূল রক্ষাকবচ।
Conflation: slow consumer সমস্যা
সবচেয়ে কঠিন বাস্তবতা: একটা mobile client বা ধীর নেটওয়ার্কের user হয়তো সেকেন্ডে ৫০০টা update সামলাতে পারবে না, যেখানে engine ১০০০টা পাঠাচ্ছে। কী করব?
- Buffer করে রাখা? Memory ফুলে যাবে, আর client পেছনে পড়ে stale data দেখবে—খারাপ।
- Conflation: যেহেতু price data-তে শুধু সর্বশেষ মান গুরুত্বপূর্ণ, gateway প্রতি symbol-এর জন্য শুধু latest state রাখে। client যখন পাঠানোর জন্য প্রস্তুত, তখন তার কাছে coalesced (একত্রিত) current state যায়, মাঝের সব intermediate tick বাদ।
queue for slow client (GP symbol):
incoming: 99.1, 99.3, 99.0, 99.5 (4টা update)
conflated send: 99.5 (শুধু latest, 1টা)
দ্রুত algo trader-দের জন্য conflation বন্ধ রাখা হয় (তারা প্রতিটা tick চায়), কিন্তু retail UI-র জন্য এটা vital।
Conflation মানে কিছু intermediate price client কখনোই দেখবে না। এটা retail price display-র জন্য ঠিক আছে, কিন্তু কখনোই trade execution-এর সিদ্ধান্ত conflated feed-এ নিও না। কেউ যদি conflated retail feed দেখে high-frequency trading করতে চায়, সে পুরনো (stale) দামে অর্ডার দিয়ে লোকসান করবে—কারণ সে মাঝের movement মিস করেছে।
Multicast vs Unicast ও fan-out architecture
- Internal (datacenter): UDP multicast। এক packet switch একবার পাঠালেই সব subscribed receiver পায়—network-এ N গুণ duplication নেই। কিন্তু UDP unreliable, তাই sequence number + আলাদা retransmission/snapshot stream লাগে।
- External (internet/mobile): multicast internet-এ চলে না, তাই edge gateway-তে এসে এটা per-client unicast (WebSocket over TCP) হয়ে যায়। এই edge layer-ই horizontally scale করে—আরও client এলে আরও gateway যোগ করো।
৭. বটলনেক ও স্কেলিং
- Fan-out explosion: আসল bottleneck হলো edge gateway-র সংখ্যা ও CPU। সমাধান: stateless gateway, সহজেই add করা যায়; এবং hot symbol-এর জন্য shared conflated state।
- Ordering vs latency: strict per-symbol ordering রাখতেই হবে, কিন্তু symbol গুলোর মধ্যে কোনো ordering লাগে না—তাই symbol দিয়ে partition করে parallel করা যায়।
- Slow client backpressure: conflation + bounded queue। queue overflow হলে client-কে "resync needed" পাঠিয়ে snapshot reload করানো হয়।
- Geographic distribution: ঢাকা, চট্টগ্রাম, বা বিদেশে edge POP বসিয়ে latency কমানো; কিন্তু সাবধান—এতে আবার fairness প্রশ্ন আসে (কে আগে পায়)।
৮. সারসংক্ষেপ
Market data feed একটা broadcast ও backpressure সমস্যা। মূল চারটা ধারণা: (১) snapshot + incremental delta দিয়ে bandwidth বাঁচানো, (২) sequence number দিয়ে ordering ও gap detection, (৩) internal multicast দিয়ে efficient one-to-many fan-out, আর (৪) edge-এ conflation দিয়ে slow client সামলানো।
ইন্টারভিউতে সবচেয়ে impressive পয়েন্ট হলো slow consumer-এর handling। বেশিরভাগ candidate শুধু "publish করব" বলে থেমে যায়; তুমি যদি conflation, bounded queue আর gap-recovery নিয়ে কথা বলো, তবেই বোঝা যায় তুমি real-time system-এর কঠিন বাস্তবতা জানো।
মিনি কুইজ
1. একটা slow consumer যদি update-এর গতি সামলাতে না পারে, তখন conflation কী করে?
2. Snapshot + incremental মডেল কেন ব্যবহার করা হয়?
3. লাখো subscriber-কে একই data exchange-এর internal network-এ পাঠানোর সবচেয়ে efficient উপায় কোনটি?