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.
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.
Class Structure
| Attribute | Type | Purpose |
|---|---|---|
| consumers | Dict[str, DataConsumer] | One DataConsumer per exchange ID (e.g., "binance") |
| _stream_subscribers | Dict[str, Set[str]] | Reference-counted subscriber sets per stream_key |
| _stream_specs | Dict[str, Dict] | Stream specification metadata |
| _stop_event | syncio.Event | Graceful 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)
- Extracts subscriber_id, required_metrics, needs_companion_orderbook.
- 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.
- 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.
- If metrics are required, updates _required_metrics under consumer._metrics_lock.
- Calls _recalculate_kline_indicators() for kline streams with metrics.
- 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 Type | Source | Snapshot Content |
|---|---|---|
| kline_1m, kline_5m, etc. | _global_kline_cache | Full deque of recent candles |
| ggTrade | _global_agg_trade_deques | Recent trade records |
| depth | consumer.get_latest_depth() | Current L2 order book snapshot |
| open_interest | consumer.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):
Connection Lifecycle
Startup (start(), lines 97–118)
- Creates aiohttp.ClientSession for HTTP connections.
- Connects to Redis.
- Initializes the default "binance" consumer via _get_consumer("binance").
- Subscribes to MARKET_DATA_REDIS_COMMAND_CHANNEL.
Shutdown (stop(), lines 120–143)
- Sets _stop_event.
- Unsubscribes and closes the pubsub connection.
- Iterates all consumers and calls consumer.stop().
- Clears the consumers dict.
- Closes Redis and aiohttp session.
Signal Handling (lines 437–441)
Sources:Reference Counting Architecture
The service uses a two-level reference counting system:
| Level | Structure | Scope |
|---|---|---|
| Service Level | _stream_subscribers: Dict[str, Set[str]] | Across all bot workers |
| Consumer Level | DataConsumer._global_ws_registry | Within 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.
Genetic Strategy Optimization
Deep dive into the DEAP-based genetic evolution algorithm, parameter mutation, crossover operations, fitness functions, and Out-of-Sample validation inside DepthSight.
Redis Fan-Out and Data Consumer
Detailed analysis of the DataConsumer runtime — dual ingestion modes, global connection registry, shared memory caches, real-time indicator processing pipeline, and warm snapshot restoration.