Market Data Pipeline

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.

⏱️ 5 min read📊 Level: Intermediate

The DataConsumer (bot_module/data_consumer.py, ~4,090 lines) is the component responsible for feeding market data (candles, order books, trades) into the trading runloops of active bots. It is the lowest-level bridge between exchange data streams and strategy evaluation.


Dual Ingestion Modes

A DataConsumer instance can operate in one of two modes, determined at construction (lines 236–258):

Direct Mode ("direct")

Used in single-user setups or local development:

  • The consumer opens raw WebSocket connections to the exchange via CCXT Pro.
  • Global Connection Registry: To prevent duplicate connections within the same Python process, active streams are registered in _global_ws_registry with reference counting.
  • If multiple strategy threads request the same stream, they share the same subscription task.

Redis Mode ("redis", "redis_pubsub", "pubsub")

Used in multi-process production deployments and SaaS installations:

  • The consumer does not connect to the internet. Instead, it subscribes to Redis Pub/Sub channels broadcast by the centralized MarketDataService.
  • Upon startup, it downloads a warm snapshot from Redis to populate its local data queues immediately.
  • Subscription requests are sent to the MARKET_DATA_REDIS_COMMAND_CHANNEL for the service to process.

Global Connection Registry and Reference Counting

Within a single process, the consumer maintains a global registry to share WebSocket connections across all consumers and strategies:

Subscription Registration (ensure_subscription)

When a stream is requested:

  1. Check if task_and_client_key exists in _global_ws_registry.
  2. If yes — increment ref_count, add consumer ID.
  3. If no — spawn a new asyncio.Task for the WebSocket loop, register with ref_count=1.

Unsubscription (remove_subscription)

When ref_count reaches zero, the WebSocket task is cancelled and the registry entry is deleted. Otherwise, the socket stays open for other consumers.


Shared Memory Caches

Within a single process hosting multiple user strategy tasks, DataConsumer maintains shared global caches protected by asynchronous locks:

Sources:
CacheKeyValuePurpose
_global_kline_cachesymbol:timeframe:exchange:marketdeque of candle tuplesRaw candle history (maxlen 5000)
_global_kline_df_cacheSame keypd.DataFrameDataFrame snapshot of kline cache
_global_depth_cachesymbol:exchange:marketOrder book dictLatest L2 snapshot
_global_agg_trade_dequessymbol:exchange:marketdeque of trade recordsRecent aggregate trades
_global_active_pairssymbolDict with price, indicators, metricsCentral runtime state store

Memory Savings: If multiple user threads trade ETHUSDT, they share a single _global_kline_cache deque. Thread Safety: The _global_cache_lock ensures indicator updates and queue mutations are atomic, preventing race conditions during high-frequency volatility spikes.


Real-Time Indicator Processing Pipeline

Kline Indicator Recalculation

Triggered on each closed candle:

  1. Reads _required_metrics for the symbol.
  2. Fetches kline history DataFrame via get_kline_history().
  3. Iterates each required indicator:
    • pandas_ta indicators: Called dynamically via getattr(kline_df.ta, kind)(**params).
    • Custom indicators: NATR_30, RELATIVE_VOLUME, IS_VOLUME_SPIKE computed via utility functions.
  4. Writes results to both _global_active_pairs and local _active_pairs.

Tape Metrics Recalculation

Triggered on each aggTrade:

  1. Retrieves the aggTrade deque.
  2. Prunes trades older than agg_trade_maxlen_seconds (300s).
  3. Calculates volume, delta, and count across windows [5, 10, 30, 60, 120] seconds.
  4. Updates _global_active_pairs with fresh tape metrics.

Event Broadcasting

When a kline closes, the consumer pushes {"type": "CANDLE_CLOSE"} events to ALL registered queues in _global_event_queues[stream_key]. For aggTrades, it pushes {"type": "TICK"} events. This fan-out design ensures every strategy receives every relevant data point.


Warm Snapshot Restoration

When a consumer starts in Redis mode, it follows a specific boot sequence to avoid cold-start latency:

1. Snapshot Request

After subscribing to Redis channels, the consumer requests a warm snapshot from the MarketDataService.

2. Wait + Load

Polls up to 5 seconds (MARKET_DATA_REDIS_SNAPSHOT_WAIT_SECONDS) for the snapshot to become available on Redis.

3. Snapshot Application

Depending on data type:

  • Klines: Clears and fills _global_kline_cache and _global_kline_df_cache. Sets last_price from the last candle's close.
  • aggTrades: Clears and fills _global_agg_trade_deques. Sets last_price from the last trade.
  • Depth: Writes directly into _latest_depth_cache.
  • Open Interest: Converts rows to DataFrame and stores in _open_interest_cache.

This ensures the strategy controller receives fully warmed indicator queues on its first evaluation cycle.


Data Flow Summary

The complete path from exchange stream to strategy evaluation:

Exchange WebSocket -> CCXT Pro Client -> _update_local_cache()
  -> _global_kline_cache / _global_agg_trade_deques
  -> _recalculate_kline_indicators() / _recalculate_tape_metrics()
  -> _global_active_pairs (current price, indicators)
  -> _publish_to_redis (if callback configured)
  -> Broadcast to _global_event_queues[stream_key]
  -> Strategy Controller's event_queue
  -> Strategy Instance (process_klines)
  -> TradingController._handle_signal()