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.
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_registrywith 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_CHANNELfor 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:
- Check if
task_and_client_keyexists in_global_ws_registry. - If yes — increment
ref_count, add consumer ID. - If no — spawn a new
asyncio.Taskfor the WebSocket loop, register withref_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:
| Cache | Key | Value | Purpose |
|---|---|---|---|
_global_kline_cache | symbol:timeframe:exchange:market | deque of candle tuples | Raw candle history (maxlen 5000) |
_global_kline_df_cache | Same key | pd.DataFrame | DataFrame snapshot of kline cache |
_global_depth_cache | symbol:exchange:market | Order book dict | Latest L2 snapshot |
_global_agg_trade_deques | symbol:exchange:market | deque of trade records | Recent aggregate trades |
_global_active_pairs | symbol | Dict with price, indicators, metrics | Central 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:
- Reads
_required_metricsfor the symbol. - Fetches kline history DataFrame via
get_kline_history(). - Iterates each required indicator:
- pandas_ta indicators: Called dynamically via
getattr(kline_df.ta, kind)(**params). - Custom indicators:
NATR_30,RELATIVE_VOLUME,IS_VOLUME_SPIKEcomputed via utility functions.
- pandas_ta indicators: Called dynamically via
- Writes results to both
_global_active_pairsand local_active_pairs.
Tape Metrics Recalculation
Triggered on each aggTrade:
- Retrieves the aggTrade deque.
- Prunes trades older than
agg_trade_maxlen_seconds(300s). - Calculates volume, delta, and count across windows [5, 10, 30, 60, 120] seconds.
- Updates
_global_active_pairswith 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_cacheand_global_kline_df_cache. Setslast_pricefrom the last candle's close. - aggTrades: Clears and fills
_global_agg_trade_deques. Setslast_pricefrom 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()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.
FastAPI REST Routes
Comprehensive overview of the FastAPI application, modular router architecture, all 22+ route modules, the Redis command bus, security middleware, and rate limiting.