Market Data Pipeline

Centralized Market Data Service

Technical breakdown of the centralized market_data_service.py daemon — WebSocket aggregation, subscription management, reference counting, snapshot persistence, and Redis publication.

⏱️ 6 min read📊 Level: Intermediate

In a multi-user, multi-bot environment, having every bot open a dedicated WebSocket connection to an exchange (like Binance or Bybit) for the same symbol would quickly trigger API rate limits and overload the server. DepthSight solves this with a centralized daemon: the MarketDataService (market_data_service.py, ~446 lines).


Centralized Aggregation Architecture

The MarketDataService acts as the single master connection manager for all live market feeds. It owns one DataConsumer per exchange (initialized in direct mode), tracks subscriber IDs per stream via reference counting, and publishes every data update to Redis for fan-out to all worker processes.

Rendering diagram...

Class Structure

AttributeTypePurpose
consumersDict[str, DataConsumer]One DataConsumer per exchange ID (e.g., "binance")
_stream_subscribersDict[str, Set[str]]Reference-counted subscriber sets per stream_key
_stream_specsDict[str, Dict]Stream specification metadata
_stop_eventsyncio.EventGraceful shutdown signal

Exchange Consumer Management (_get_consumer)

The _get_consumer() method (lines 65–95) creates singleton consumers per exchange ID:

Sources:

Each consumer is initialized with market_data_mode="direct", meaning it opens raw WebSocket connections to the exchange. A publish_callback is wired to _publish_market_payload, so every incoming data update is forwarded to Redis in real-time.


Subscription Command Loop

Main Loop (run(), lines 145–170)

The service listens on MARKET_DATA_REDIS_COMMAND_CHANNEL for subscription commands:

Sources:

If no messages arrive for 30 seconds (configurable watchdog), the service automatically reconnects the Redis pubsub connection via _reconnect_pubsub() (lines 172–184).

Command Dispatch (_handle_command, lines 186–193)

Sources:

Subscribe Handler (_handle_subscribe, lines 195–251)

  1. Extracts subscriber_id, required_metrics, needs_companion_orderbook.
  2. Iterates stream specs from _command_stream_specs() (lines 413–430), which supports both a batch format (stream_keys list) and a single stream_key shorthand.
  3. Reference Counting: Adds the subscriber to _stream_subscribers[stream_key]. Only calls consumer.ensure_subscription() when the subscriber set transitions from empty to non-empty — guaranteeing exactly one exchange WebSocket connection per unique stream_key, regardless of how many bot workers subscribe.
  4. If metrics are required, updates _required_metrics under consumer._metrics_lock.
  5. Calls _recalculate_kline_indicators() for kline streams with metrics.
  6. Writes a warm snapshot via _write_snapshot_for_spec() after subscribing.

Unsubscribe Handler (_handle_unsubscribe, lines 252–281)

Removes the subscriber from _stream_subscribers[stream_key]. Only calls consumer.remove_subscription() when the subscriber set becomes empty (ref count reaches zero), preventing unnecessary WebSocket reconnections.


Publication Pipeline (_publish_market_payload)

When a WebSocket payload is received by any consumer's callback, it enters the publication pipeline (lines 283–315):

Sources:

The event channel format: depthsight:market_data:events:{exchange}:{symbol}:{data_type}.


Snapshot System (_write_snapshot_for_spec)

To prevent newly launched bots from waiting for the next candle to initialize their indicators, the service continuously writes warm snapshots (lines 317–411):

Data TypeSourceSnapshot Content
kline_1m, kline_5m, etc._global_kline_cacheFull deque of recent candles
ggTrade_global_agg_trade_dequesRecent trade records
depthconsumer.get_latest_depth()Current L2 order book snapshot
open_interestconsumer.get_open_interest_history()OI time series

Snapshots are stored in Redis with key pattern depthsight:market_data:snapshot:{stream_key} and a configurable TTL (default: 3600 seconds):

Sources:

Connection Lifecycle

Startup (start(), lines 97–118)

  1. Creates aiohttp.ClientSession for HTTP connections.
  2. Connects to Redis.
  3. Initializes the default "binance" consumer via _get_consumer("binance").
  4. Subscribes to MARKET_DATA_REDIS_COMMAND_CHANNEL.

Shutdown (stop(), lines 120–143)

  1. Sets _stop_event.
  2. Unsubscribes and closes the pubsub connection.
  3. Iterates all consumers and calls consumer.stop().
  4. Clears the consumers dict.
  5. Closes Redis and aiohttp session.

Signal Handling (lines 437–441)

Sources:

Reference Counting Architecture

The service uses a two-level reference counting system:

LevelStructureScope
Service Level_stream_subscribers: Dict[str, Set[str]]Across all bot workers
Consumer LevelDataConsumer._global_ws_registryWithin a single Python process

A stream is active only when _stream_subscribers[stream_key] is non-empty. When the last subscriber leaves, the consumer's WebSocket is closed, and resources are freed. This design allows DepthSight to serve thousands of concurrent bot workers with a minimal number of exchange connections.