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

Real-time Market Data Feed ডিজাইন

14 মিনিটadvanced
এক নজরে
  • 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, qtyhistorical/replay, charting
audit_logseq_no, payload, tsregulatory 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-র যাত্রা:

  1. Matching engine প্রতিটা order book পরিবর্তনে একটা sequenced update event বের করে।
  2. Feed handler / normalizer raw event নিয়ে একটা compact wire format (incremental message) বানায়, seq_no বসায়।
  3. Internal multicast bus: এই incremental stream একটা multicast group-এ একবার publish হয়। সব downstream consumer (gateway, OMS, surveillance) একই group থেকে পড়ে।
  4. Snapshot service: একটা component multicast stream consume করে in-memory order book maintain করে, এবং নিয়মিত interval-এ full snapshot publish করে আরেকটা group-এ।
  5. Edge gateways (WebSocket fan-out): এরা multicast থেকে পড়ে, লাখো client-এর কাছে unicast WebSocket-এ পাঠায়। এখানেই conflation ও per-client queue থাকে।
  6. 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 > 1000 delta 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 উপায় কোনটি?