For the past few months, I've been dedicating my nights and weekends to a side project that combines my passions for machine learning and financial markets. As an ML/AI engineer by day, I wanted to challenge myself with something that would push my technical boundaries.
MarketPulseAI is an open-source, near real-time stock market analysis system that combines two powerful perspectives:
- Stock Market Prediction using price data and deep learning algorithms
- Market Sentiment Analysis from social media and financial news
Think of it as having two eyes on the market: one watching actual prices and trading patterns, and another watching what people are saying about stocks on social media and in the news. The system then integrates these signals to provide a holistic view of potential market movements.
- Microservices Architecture: Modular design with loosely coupled services
- Near Real-Time Processing: Data flows through the system with minimal latency
- Dual Analysis Approach: Combines price prediction with sentiment analysis
- Comprehensive Data Validation: Ensures high-quality data reaches prediction models
- Interactive Dashboards: Real-time visualization of market data and sentiment
- Robust Monitoring: Complete observability of system health and prediction accuracy
MarketPulseAI follows a modern data engineering pattern with clear separation of concerns:
-
Market Data Pipeline
- Connects to stock market feeds via Kafka Connect
- Collects real-time price data, volumes, and market indicators
- Validates and routes data to appropriate Kafka topics
-
Sentiment Analysis Pipeline
- Monitors Twitter, Reddit, and financial news via APIs
- Collects and filters relevant posts and articles
- Validates content before processing
-
Market Data Analysis
- Spark Streaming jobs process real-time market data
- Feature engineering extracts technical indicators
- Deep learning models predict short-term price movements
-
Sentiment Analysis
- Text preprocessing and cleaning
- NLP or Transformer models determine sentiment polarity and intensity
- Aggregation of sentiment across different sources
-
Online Features (Redis)
- Low-latency access to real-time features
- Optimized caching with custom serialization
- Configurable TTL for data freshness
-
Historical Features (Cassandra)
- Distributed storage for training data
- Time-series optimized schema
- Efficient querying for model training
- Combines predictions from market and sentiment models
- Weighted ensemble approach for final predictions
- Continuous evaluation and re-weighting based on performance
- FastAPI service with WebSocket support
- Streamlit dashboards for interactive visualization
- Real-time updates of predictions and sentiment scores
- Prometheus metrics collection
- Grafana dashboards for system visualization
- InfluxDB for time-series performance metrics
- Alerting system for anomaly detection
| Component | Technologies |
|---|---|
| Data Ingestion | Apache Kafka, Kafka Connect |
| Processing | Apache Spark Streaming |
| Storage | Redis, Apache Cassandra |
| ML & Analytics | PyTorch, TensorFlow, spaCy, NLTK |
| API & Serving | FastAPI, WebSockets |
| Visualization | Streamlit, Plotly |
| Deployment | Docker, Kubernetes |
| Monitoring | Prometheus, Grafana, InfluxDB |
I wanted to challenge myself with an ambitious, end-to-end application that would push my limits across multiple domains. Here's my thinking behind the approach:
โข Comprehensive Learning: Building everything from data ingestion to visualization gave me a holistic view of the entire ML pipeline in production
โข Technical Depth: I dove deep into configuring each technology - tuning Kafka partitioning for optimal throughput, optimizing Spark executor memory allocation, and implementing custom serialization for Redis caching
โข New Territories: This project pushed me to explore technologies I rarely use, particularly in platform monitoring and observability (Prometheus metric collection, Grafana dashboard configuration, and alert management)
โข Open Source First: I committed to using open-source technologies throughout the stack to keep the project accessible and modifiable
โข Local to Cloud Path: I designed everything to run locally first with Docker Compose, with a clear migration path to cloud services once the fundamentals are solid
โข Community Driven: After countless discussions with data engineers and professionals on Reddit and Discord, I incorporated many of their suggestions and best practices
โข Practical Approach: I deliberately chose near real-time processing over true streaming (no Apache Flink) because it provides an excellent balance between performance and complexity for this use case
โข Proof of Concept: I wanted to demonstrate that even smaller companies can implement sophisticated data pipelines without massive infrastructure investments
โข Skill Stretching: If I can successfully push ML into production for real-time analysis, other ML deployment scenarios will become significantly easier by comparison
โข Finance Education: This project doubled as an incredible learning journey into financial markets, technical analysis, and trading psychology
The most fascinating discovery has been seeing how social sentiment sometimes predicts price movements before they appear in market data!
- Add more data sources (options flow, institutional trading patterns)
- Implement transformer-based models for better time-series forecasting
- Add support for cryptocurrency markets
- Improve Transformer models for more nuanced sentiment analysis
- Optimize performance across the entire pipeline
- Add cloud deployment templates (AWS, GCP, Azure)
- Add user authentication and multi-user support
Contributions are welcome! Please feel free to submit a Pull Request.
- Fork the repository
- Create your feature branch (
git checkout -b feature/amazing-feature) - Commit your changes (
git commit -m 'Add some amazing feature') - Push to the branch (
git push origin feature/amazing-feature) - Open a Pull Request
At the end, the project should look like this:
stock-market-analysis-system/
โ
โโโ .env # Environment variables (gitignored)
โโโ .gitignore # Git ignore rules
โโโ README.md # Project documentation
โโโ requirements.txt # Python dependencies
โโโ docker-compose.yml # Docker services configuration
โโโ Makefile # Common commands for development
โ
โโโ config/ # Configuration files
โ โโโ kafka/ # Kafka configuration
โ โโโ spark/ # Spark configuration
โ โโโ redis/ # Redis configuration
โ โโโ database/ # Database schema and migrations
โ โโโ logging/ # Logging configuration
โ โโโ apis/ # API configuration files
โ โโโ newsapi.yaml # NewsAPI configuration
โ โโโ other_apis.yaml # Other API configurations
โ
โโโ data/ # Data storage (gitignored)
โ โโโ raw/ # Raw data from APIs
โ โ โโโ market/ # Raw market data
โ โ โโโ social/ # Raw social media data
โ โ โโโ news/ # Raw news data from NewsAPI
โ โโโ processed/ # Processed data
โ โ โโโ market/ # Processed market data
โ โ โโโ social/ # Processed social media data
โ โ โโโ news/ # Processed news data
โ โโโ models/ # Saved ML models
โ โโโ cache/ # Cached data
โ
โโโ notebooks/ # Jupyter notebooks for experimentation
โ โโโ market_data_exploration/ # Market data analysis
โ โโโ sentiment_analysis/ # Sentiment analysis exploration
โ โโโ news_data_exploration/ # News data analysis notebooks
โ โโโ model_development/ # ML model development
โ โโโ visualization/ # Visualization experiments
โ
โโโ docs/ # Documentation
โ โโโ architecture/ # System architecture docs
โ โโโ api/ # API documentation
โ โ โโโ newsapi/ # NewsAPI integration docs
โ โโโ user_guide/ # User documentation
โ โโโ development/ # Development guidelines
โ
โโโ scripts/ # Utility scripts
โ โโโ setup/ # Setup scripts
โ โโโ data_collection/ # Data collection scripts
โ โ โโโ market/ # Market data collection scripts
โ โ โโโ social/ # Social media collection scripts
โ โ โโโ news/ # News data collection scripts
โ โ โโโ newsapi_fetcher.py # NewsAPI data fetcher
โ โ โโโ news_backfill.py # Historical news data backfill
โ โโโ backup/ # Backup scripts
โ โโโ maintenance/ # Maintenance scripts
โ
โโโ tests/ # Test suite
โ โโโ unit/ # Unit tests
โ โ โโโ market_data/ # Market data tests
โ โ โโโ sentiment/ # Sentiment analysis tests
โ โ โโโ news/ # News data collection tests
โ โ โโโ prediction/ # Prediction engine tests
โ โ โโโ api/ # API tests
โ โโโ integration/ # Integration tests
โ โ โโโ news_pipeline_test.py # News data pipeline integration test
โ โโโ performance/ # Performance tests
โ โโโ fixtures/ # Test fixtures
โ โโโ news/ # News data test fixtures
โ
โโโ monitoring/ # System monitoring
โ โโโ health_checks/ # Health check scripts
โ โ โโโ news_api_health.py # NewsAPI health check
โ โโโ metrics/ # Metrics collection
โ โ โโโ news_metrics.py # News data collection metrics
โ โโโ alerts/ # Alert configuration
โ โโโ dashboards/ # Monitoring dashboards
โ โโโ news_dashboard.json # News data monitoring dashboard
โ
โโโ src/ # Source code
โ โโโ data_collection/ # Data collection modules
โ โ โโโ __init__.py
โ โ โโโ market_data/ # Market data collection
โ โ โ โโโ __init__.py
โ โ โ โโโ collectors/ # Data source collectors
โ โ โ โโโ parsers/ # Data parsers
โ โ โ โโโ validation/ # Data validation
โ โ โ
โ โ โโโ social_media/ # Social media collection
โ โ โ โโโ __init__.py
โ โ โ โโโ twitter/ # Twitter data collection
โ โ โ โโโ reddit/ # Reddit data collection
โ โ โ
โ โ โ
โ โ โโโ news/ # News data collection (NEW)
โ โ โโโ __init__.py
โ โ โโโ collectors/ # News data collectors
โ โ โ โโโ __init__.py
โ โ โ โโโ newsapi_collector.py # NewsAPI collector
โ โ โ โโโ base_collector.py # Base news collector class
โ โ โโโ validation/ # News data validation
โ โ โโโ __init__.py
โ โ โโโ news_validator.py # News data validator
โ โ
โ โโโ data_processing/ # Data processing modules
โ โ โโโ __init__.py
โ โ โโโ market_data/ # Market data processing
โ โ โ โโโ __init__.py
โ โ โ โโโ cleaning/ # Data cleaning
โ โ โ โโโ feature_eng/ # Feature engineering
โ โ โ โโโ indicators/ # Technical indicators
โ โ โ
โ โ โโโ sentiment/ # Sentiment processing
โ โ โ โโโ __init__.py
โ โ โ โโโ preprocessing/ # Text preprocessing
โ โ โ โโโ analysis/ # Sentiment analysis
โ โ โ โโโ aggregation/ # Sentiment aggregation
โ โ โ
โ โ โโโ news/ # News data processing (NEW)
โ โ โโโ __init__.py
โ โ โโโ cleaning/ # News cleaning
โ โ โ โโโ __init__.py
โ โ โ โโโ news_cleaner.py # News text cleaner
โ โ โโโ enrichment/ # News data enrichment
โ โ โ โโโ __init__.py
โ โ โ โโโ company_matcher.py # Match companies to news
โ โ โ โโโ topic_classifier.py # News topic classifier
โ โ โโโ aggregation/ # News data aggregation
โ โ โโโ __init__.py
โ โ โโโ news_aggregator.py # News aggregation by symbol/sector
โ โ
โ โโโ storage/ # Data storage modules
โ โ โโโ __init__.py
โ โ โโโ database/ # Database operations
โ โ โ โโโ __init__.py
โ โ โ โโโ market_db.py # Market data DB operations
โ โ โ โโโ social_db.py # Social media DB operations
โ โ โ โโโ news_db.py # News data DB operations (NEW)
โ โ โโโ streaming/ # Streaming data handlers
โ โ โโโ cache/ # Caching operations
โ โ โโโ models/ # Model storage
โ โ
โ โโโ models/ # Machine learning models
โ โ โโโ __init__.py
โ โ โโโ market_prediction/ # Price prediction models
โ โ โ โโโ __init__.py
โ โ โ โโโ feature_selection/# Feature selection
โ โ โ โโโ training/ # Model training
โ โ โ โโโ evaluation/ # Model evaluation
โ โ โ โโโ prediction/ # Prediction generation
โ โ โ
โ โ โโโ sentiment_models/ # Sentiment models
โ โ โ โโโ __init__.py
โ โ โ โโโ classification/ # Sentiment classification
โ โ โ โ โโโ __init__.py
โ โ โ โ โโโ social_classifier.py # Social media classifier
โ โ โ โ โโโ news_classifier.py # News content classifier (NEW)
โ โ โ โโโ evaluation/ # Sentiment model evaluation
โ โ โ
โ โ โโโ combined_models/ # Combined prediction models
โ โ โโโ __init__.py
โ โ โโโ feature_fusion/ # Feature combination
โ โ โโโ ensemble/ # Ensemble models
โ โ
โ โโโ streaming/ # Streaming data processing
โ โ โโโ __init__.py
โ โ โโโ kafka/ # Kafka producers/consumers
โ โ โ โโโ __init__.py
โ โ โ โโโ producers/ # Kafka producers
โ โ โ โ โโโ __init__.py
โ โ โ โ โโโ market_producer.py # Market data producer
โ โ โ โ โโโ social_producer.py # Social media producer
โ โ โ โ โโโ news_producer.py # News data producer (NEW)
โ โ โ โโโ consumers/ # Kafka consumers
โ โ โ โโโ __init__.py
โ โ โ โโโ market_consumer.py # Market data consumer
โ โ โ โโโ social_consumer.py # Social media consumer
โ โ โ โโโ news_consumer.py # News data consumer (NEW)
โ โ โ
โ โ โโโ spark/ # Spark streaming jobs
โ โ โโโ __init__.py
โ โ โโโ news_processing_job.py # News processing Spark job (NEW)
โ โ
โ โโโ prediction_engine/ # Prediction engine
โ โ โโโ __init__.py
โ โ โโโ scheduler/ # Prediction scheduling
โ โ โโโ executor/ # Prediction execution
โ โ โโโ evaluation/ # Real-time evaluation
โ โ
โ โโโ api/ # API layer
โ โ โโโ __init__.py
โ โ โโโ routes/ # API routes
โ โ โ โโโ __init__.py
โ โ โ โโโ market_routes.py # Market data routes
โ โ โ โโโ social_routes.py # Social media routes
โ โ โ โโโ news_routes.py # News data routes (NEW)
โ โ โโโ serializers/ # Data serializers
โ โ โ โโโ __init__.py
โ โ โ โโโ news_serializer.py # News data serializer (NEW)
โ โ โโโ auth/ # Authentication
โ โ โโโ middleware/ # API middleware
โ โ
โ โโโ dashboard/ # Web dashboard
โ โ โโโ backend/ # Dashboard backend
โ โ โ โโโ __init__.py
โ โ โ โโโ server.py # Web server
โ โ โ โโโ websockets/ # WebSocket handlers
โ โ โ
โ โ โโโ frontend/ # Dashboard frontend
โ โ โโโ public/ # Public assets
โ โ โโโ src/ # Frontend source code
โ โ โ โโโ components/ # UI components
โ โ โ โ โโโ news/ # News components (NEW)
โ โ โ โ โโโ NewsCard.js # News item card
โ โ โ โ โโโ NewsStream.js # News stream component
โ โ โ โ โโโ NewsAnalytics.js # News analytics component
โ โ โ โโโ pages/ # Page components
โ โ โ โ โโโ NewsPage.js # News analysis page (NEW)
โ โ โ โโโ services/ # API services
โ โ โ โ โโโ newsService.js # News API service (NEW)
โ โ โ โโโ hooks/ # Custom hooks
โ โ โ โโโ utils/ # Utility functions
โ โ โ โโโ App.js # Main app component
โ โ โ
โ โ โโโ package.json # Frontend dependencies
โ โ โโโ README.md # Frontend documentation
โ โ
โ โโโ utils/ # Utility modules
โ โ โโโ __init__.py
โ โ โโโ logging/ # Logging utilities
โ โ โโโ config/ # Configuration utilities
โ โ โโโ validation/ # Data validation
โ โ โโโ datetime/ # Date/time utilities
โ โ โโโ metrics/ # Metrics utilities
โ โ โโโ api/ # API utilities (NEW)
โ โ โโโ __init__.py
โ โ โโโ rate_limiter.py # API rate limiting
โ โ โโโ request_retry.py # Request retry logic
โ โ
โ โโโ alerts/ # Alert system
โ โโโ __init__.py
โ โโโ triggers/ # Alert triggers
โ โโโ notifications/ # Notification delivery
โ โโโ templates/ # Alert templates
โ
โโโ dockerfiles/ # Dockerfile for each service
โโโ api/ # API service Dockerfile
โโโ data_collection/ # Data collection Dockerfile
โโโ processing/ # Processing Dockerfile
โโโ prediction/ # Prediction Dockerfile
โโโ news_collector/ # News collector Dockerfile (NEW)
โโโ dashboard/ # Dashboard Dockerfile
- Update frequency: Near real-time (seconds, not milliseconds)
- Capacity:
- Handles millions of market data points per second
- Processes thousands of social media posts
- Analyzes hundreds of news articles
- Prediction latency: < 5 seconds from data ingestion to prediction
- Typical accuracy metrics: Details in performance documentation
This project is licensed under the MIT License - see the LICENSE file for details.
- Apache Kafka
- Apache Spark
- Redis
- Apache Cassandra
- FastAPI
- Streamlit
- PyTorch
- TensorFlow
- Docker
- Kubernetes
- Prometheus
- Grafana
If you find MarketPulseAI useful, please consider starring the repository to help others discover it!
