A professional-grade web scraping solution with dynamic proxy rotation and asynchronous data processing, designed for high-concurrency financial data collection.
- π Dynamic Proxy Rotator: Intelligent proxy selection based on health scores with automatic recovery
- β‘ Asyncio + ThreadPoolExecutor: Optimal combination for I/O-bound and CPU-bound tasks
- π Thread-Safe Buffer: RLock-protected storage with MD5-based deduplication
- β JSON Schema Validation: Comprehensive validation for financial data structures
- π Zero Data Loss: Guaranteed data integrity under 100+ concurrent requests
- π― Smart Retry Logic: Automatic failover with alternative proxies
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
β FinancialDataScraper β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ€
β ββββββββββββββββββββ ββββββββββββββββββββββββββββββββ β
β β DynamicProxy β β ThreadSafeBuffer β β
β β Rotator β β - RLock Protection β β
β β - Health Scores β β - MD5 Deduplication β β
β β - Weighted Pick β β - Statistics Tracking β β
β ββββββββββββββββββββ ββββββββββββββββββββββββββββββββ β
β β
β ββββββββββββββββββββββββββββββββββββββββββββββββββββββββ β
β β asyncio Event Loop β β
β β (I/O-bound: HTTP Requests via aiohttp) β β
β ββββββββββββββββββββββββββββββββββββββββββββββββββββββββ β
β β¬οΈ β
β ββββββββββββββββββββββββββββββββββββββββββββββββββββββββ β
β β ThreadPoolExecutor β β
β β (CPU-bound: JSON Schema Validation) β β
β ββββββββββββββββββββββββββββββββββββββββββββββββββββββββ β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
# Clone the repository
git clone https://github.com/yourusername/multi-threaded-web-scraper.git
cd multi-threaded-web-scraper
# Install dependencies
pip install aiohttpimport asyncio
from web_scraper import FinancialDataScraper, FINANCIAL_ENDPOINTS, PROXY_POOL
async def main():
scraper = FinancialDataScraper(
endpoints=FINANCIAL_ENDPOINTS,
proxies=PROXY_POOL,
max_workers=10
)
# Run with 100 concurrent requests
stats = await scraper.run(request_count=100)
print(f"Stored: {stats['stored']}")
print(f"Duplicates: {stats['duplicates']}")
if __name__ == "__main__":
asyncio.run(main())# Define your own endpoints
endpoints = [
"https://api.example.com/stock/AAPL",
"https://api.example.com/stock/GOOGL",
"https://api.example.com/stock/MSFT"
]
# Define your proxy pool
proxies = [
"http://proxy1.example.com:8080",
"http://proxy2.example.com:8080",
"http://proxy3.example.com:8080"
]
scraper = FinancialDataScraper(
endpoints=endpoints,
proxies=proxies,
max_workers=20 # Adjust thread pool size
)Test the scraper with 100 concurrent requests:
python benchmark_test.py======================================================================
COMPREHENSIVE BENCHMARK TEST
======================================================================
Test Configuration:
- Concurrent Requests: 100
- Financial Endpoints: 5
- Expected Records: 500
============================================================
BENCHMARK RESULTS
============================================================
Number of Requests: 100
Total Expected Records: 500
Total Stored: 500
Data Loss: 0
Success Rate: 100.00%
Elapsed Time: 0.280 seconds
Requests/Second: 1788.84
============================================================
β
SUCCESS: Zero data loss confirmed!
multi-threaded-web-scraper/
βββ web_scraper.py # Main scraper implementation
βββ benchmark_test.py # Comprehensive benchmark tests
βββ README.md # This file
βββ LICENSE # MIT License
βββ requirements.txt # Python dependencies
Intelligent proxy management with health-based scoring:
- Health Score: Tracks success/failure ratio (0.0 - 1.0)
- Weighted Selection: Probabilistic selection favoring healthy proxies
- Automatic Recovery: Failed proxies can recover over time
@dataclass
class ProxyHealth:
url: str
health_score: float = 1.0
failures: int = 0
def record_success(self):
self.health_score = min(1.0, self.health_score + 0.1)
def record_failure(self):
self.health_score = max(0.0, self.health_score - 0.2)Thread-safe data storage with deduplication:
- RLock Protection: Prevents race conditions
- MD5 Hashing: Efficient duplicate detection
- Statistics Tracking: Monitors received, stored, and lost data
Main orchestrator combining asyncio and threading:
- Async HTTP Requests: Non-blocking I/O with aiohttp
- CPU-bound Validation: Offloaded to ThreadPoolExecutor
- Retry Logic: Automatic failover with alternative proxies
| Metric | Value |
|---|---|
| Concurrent Requests | 100+ |
| Data Loss | 0% |
| Throughput | ~1800 req/sec |
| Execution Time | < 1 second |
| Memory Efficiency | O(n) with deduplication |
The implementation ensures thread safety through:
- RLock for buffer operations
- Atomic operations for proxy health updates
- Proper executor shutdown on completion
- Session management with async context managers
MIT License - see LICENSE for details.
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 amazing feature') - Push to the branch (
git push origin feature/amazing-feature) - Open a Pull Request
For issues and questions, please open an issue on GitHub.
Built with β€οΈ using Python, asyncio, and ThreadPoolExecutor