BifrostPipe – Unified Market-Data Bus
High-performance C++ Kafka-centric data pipeline streaming equity and option tick data from Alpaca Pro WebSockets with real-time enrichment and technical analysis
Status: completed · 2024-11-15
Overview
BifrostPipe is a Kafka-centric data pipeline written in C++ that streams equity and option tick data from Alpaca's Pro Market-Data WebSockets into your local or cloud stack, with real-time enrichment and technical analysis capabilities.
Technologies
C++, Kafka, libwebsockets, librdkafka, TA-lib, Avro C++, Redis, Docker, REST API, Prometheus
- Latency
- <1ms
- Throughput
- 10K+ TPS
- Cache Hit Rate
- 95%+
- Uptime
- 99.9%
BifrostPipe – Unified Market-Data Bus
BifrostPipe is a Kafka-centric data pipeline written in C++ that streams equity and option tick data from Alpaca's Pro Market-Data WebSockets into your local or cloud stack. It enriches stock data with fundamentals from Financial Modeling Prep (FMP), computes technical indicators using TA-lib, and writes clean, versioned Avro records to downstream Kafka topics for analytics, machine learning, and trading bots.
Built for Arch Linux power users, containerized with Docker, and optimized for sub-millisecond throughput.
Key Features
🔄 Dynamic Subscription Management
- REST API Subscriptions: Subscribe to specific tickers for stocks and options via REST API
- Kafka Message Subscriptions: Manage subscriptions dynamically via Kafka messages for high-volume systems
- Custom Kafka Topics: Specify custom Kafka topics per subscription or use default topics
- User-Based Management: Multi-user subscription management with user isolation
📊 Technical Analysis Integration
- TA-lib Integration: Compute technical indicators for stocks and options with configurable defaults
- Per-Subscription Configuration: Specify custom indicators per subscription
- Real-time Computation: On-the-fly technical analysis with sub-millisecond latency
- Historical Data Support: Initialize indicators requiring historical data
🚀 Real-time Data Enrichment
- FMP Fundamentals: On-the-fly joins with Financial Modeling Prep fundamentals
- Market Data Enhancement: Sector, market-cap, float, insider-ownership enrichment
- Redis Caching: High-performance caching for enrichment data
- Configurable TTL: Customizable cache expiration policies
⚡ High-Performance Architecture
- C++ Implementation: Leverages libwebsockets, librdkafka, and TA-lib for low-latency processing
- Sub-millisecond Throughput: Optimized for high-frequency trading requirements
- Concurrent Processing: Multi-threaded architecture for maximum performance
- Memory Efficiency: Optimized memory usage for continuous operation
🔧 Configuration & Management
- YAML Configuration: Config-driven setup for Kafka, Alpaca, FMP, and TA-lib settings
- Environment Variables: Secure configuration management
- Docker Support: Containerized deployment with Docker Compose
- Prometheus Metrics: Built-in observability with metrics exporter
Architecture Overview
graph TD
A[Alpaca Stock WebSocket] --> B[Stock WS Client]
C[Alpaca Options WebSocket] --> D[Option WS Client]
E[REST API Server] --> F[Subscription Manager]
G[Kafka Subscription Manager] --> F
B --> H[TA-lib Processor]
D --> H
F --> H
H --> I[Avro Serializer]
I --> J[Kafka Producer]
J --> K[Stock Ticks Topic]
J --> L[Option Ticks Topic]
J --> M[Custom Topics]
N[FMP API] --> O[Redis Cache]
O --> P[Enrichment Processor]
K --> P
L --> P
P --> Q[Enriched Data Topic]
R[Prometheus Metrics] --> S[Grafana Dashboard]
H --> R
J --> R
Data Processing Pipeline
graph TD
A[WebSocket Data] --> B{Data Type}
B -->|Stock| C[Stock Processor]
B -->|Option| D[Option Processor]
C --> E[TA-lib Analysis]
D --> E
E --> F[Avro Serialization]
F --> G[Kafka Producer]
G --> H[Kafka Topics]
I[FMP Cache] --> J[Enrichment Engine]
H --> J
J --> K[Enriched Topic]
L[Subscription Manager] --> M[Topic Routing]
M --> G
Technical Implementation
WebSocket Clients
- Dual WebSocket Architecture: Separate clients for stock and option data streams
- Authentication Management: Secure authentication with Alpaca Pro Market-Data
- Reconnection Logic: Automatic reconnection with exponential backoff
- Data Validation: Real-time validation of incoming JSON/MsgPack frames
Kafka Integration
- librdkafka: High-performance Kafka client library
- Avro Serialization: Versioned schema support for data evolution
- Topic Management: Dynamic topic creation and partition management
- Batch Processing: Configurable batching for optimal throughput
REST API Management
POST /subscribe - Subscribe to ticker (stock/option)
POST /unsubscribe - Unsubscribe from ticker
GET /subscriptions - List active subscriptions
POST /api/option-chain-subscribe - Subscribe to option chains
GET /api/options/expirations/{symbol} - Get available expirations
GET /api/options/strikes/{symbol}/{expiration} - Get strike prices
GET /api/user/{userId}/subscriptions - List user subscriptions
Kafka Subscription Management
- Message-Driven Subscriptions: Manage subscriptions via Kafka messages
- High-Volume Support: Designed for automated and integrated systems
- Subscription Topic: Dedicated topic for subscription management messages
- Real-time Processing: Immediate subscription updates without API calls
Performance & Observability
Metrics & Monitoring
- Prometheus Integration: Built-in metrics exporter for observability
- Key Metrics: Ticks/second, WebSocket latency, enrichment hit-rate
- Grafana Dashboards: Pre-configured dashboards for monitoring
- Custom Alerts: Configurable alerting for performance thresholds
Production Optimization
- Memory Management: Optimized memory allocation and garbage collection
- Connection Pooling: Efficient WebSocket and HTTP connection management
- Caching Strategy: Redis-based caching for FMP enrichment data
- Scaling Support: Horizontal scaling with Kafka partitioning
Option Chain Management
Advanced Option Features
- Strike Range Subscriptions: Subscribe to multiple options within strike ranges
- Expiration Management: List and manage option expiration dates
- Chain Tagging: Group option subscriptions with custom tags
- Bulk Operations: Efficient bulk subscription and unsubscription
Greeks and Analytics
- Real-time Greeks: Option Greeks computation and streaming
- Chain Analysis: Complete option chain data processing
- Volatility Tracking: Implied volatility calculation and monitoring
- Risk Metrics: Delta, gamma, theta, vega calculations
Deployment & Configuration
Prerequisites
- Arch Linux Environment: Optimized for Arch Linux power users
- Dependencies: libwebsockets, librdkafka, TA-lib, Avro C++, cpprestsdk
- External Services: Kafka cluster, Redis instance, Alpaca Pro subscription
- API Access: Financial Modeling Prep API key
Configuration Management
- YAML Configuration: Centralized configuration in bifrost.yml
- Environment Variables: Secure credential management
- Docker Support: Complete containerization with Docker Compose
- Kubernetes Ready: Helm charts for Kubernetes deployment
This high-performance data pipeline enables real-time market data processing with enterprise-grade reliability and sub-millisecond latency for demanding trading applications.