Deepvest/ProfilingAI
0
1# src/core/data_fetcher.py2 3import yfinance as yf4import finnhub5import pandas as pd6import numpy as np7from typing import List, Dict, Optional8from datetime import datetime, timedelta9import aiohttp10import asyncio11import json12from bs4 import BeautifulSoup13import requests14from newspaper import Article15import tweepy16import os17from dotenv import load_dotenv18 19load_dotenv(os.path.join(os.path.dirname(__file__), 'Deepvest_system', 'src', 'data_API.env'))20 21 22class DataFetcher:23 """Gestionnaire centralisé de récupération de données"""24 25 def __init__(self):26 self.finnhub_client = finnhub.Client(api_key="cu9bhbhr01qnf5nmldv0cu9bhbhr01qnf5nmldvg")27 self.twitter_auth = tweepy.OAuthHandler(28 "HaynwOOFiEd2Ty3m6Xu7sPhrd",29 "nuu7y46JtjB7qbYnhK79w86AqR9PC2maDzl2qDpJZLceQeEe5Y"30 )31 self.twitter_auth.set_access_token(32 "1871112111419486208-taeYKwppobslPmDa4NaDMCqCEG9qoa",33 "N2izETc3KGe6tPE8yc1rvKTo20lnZewhxWw19HhzlebRp"34 )35 self.twitter_api = tweepy.API(self.twitter_auth)36 37 async def fetch_market_data(self, symbols: List[str], period: str = "2y") -> pd.DataFrame:38 """Récupération des données de marché"""39 data = pd.DataFrame()40 41 for symbol in symbols:42 try:43 # Données de base avec yfinance44 yf_data = yf.download(symbol, period=period)45 46 # Données supplémentaires avec Finnhub47 finnhub_data = self._get_finnhub_data(symbol)48 49 # Combiner les données50 combined_data = pd.concat([yf_data, finnhub_data], axis=1)51 data[symbol] = combined_data['Close']52 53 except Exception as e:54 print(f"Erreur lors de la récupération des données pour {symbol}: {e}")55 56 return data57 58 async def fetch_news_data(self, symbols: List[str], days: int = 7) -> List[Dict]:59 """Récupération des nouvelles financières"""60 news_data = []61 62 async with aiohttp.ClientSession() as session:63 for symbol in symbols:64 try:65 # Nouvelles de Finnhub66 finnhub_news = self.finnhub_client.company_news(67 symbol, 68 _from=(datetime.now() - timedelta(days=days)).strftime('%Y-%m-%d'),69 to=datetime.now().strftime('%Y-%m-%d')70 )71 72 # Articles de presse avec newspaper3k73 for news in finnhub_news:74 try:75 article = Article(news['url'])76 article.download()77 article.parse()78 79 news_data.append({80 'symbol': symbol,81 'title': news['headline'],82 'content': article.text,83 'source': news['source'],84 'url': news['url'],85 'datetime': datetime.fromtimestamp(news['datetime']),86 'sentiment': news.get('sentiment')87 })88 except Exception as e:89 print(f"Erreur lors de l'analyse de l'article: {e}")90 91 except Exception as e:92 print(f"Erreur lors de la récupération des nouvelles pour {symbol}: {e}")93 94 return news_data95 96 async def fetch_alternative_data(self, symbols: List[str]) -> Dict:97 """Récupération des données alternatives"""98 return {99 'social_media': await self._fetch_social_media_data(symbols),100 'web_traffic': await self._fetch_web_traffic_data(symbols),101 'satellite': await self._fetch_satellite_data(symbols)102 }103 104 async def _fetch_social_media_data(self, symbols: List[str]) -> List[Dict]:105 """Récupération des données des réseaux sociaux"""106 social_data = []107 108 for symbol in symbols:109 try:110 # Tweets111 tweets = self.twitter_api.search_tweets(112 q=f"${symbol}",113 lang="en",114 count=100,115 tweet_mode="extended"116 )117 118 for tweet in tweets:119 social_data.append({120 'platform': 'twitter',121 'text': tweet.full_text,122 'timestamp': tweet.created_at,123 'engagement': tweet.favorite_count + tweet.retweet_count,124 'symbol': symbol125 })126 127 # Reddit (exemple avec PRAW)128 # Reddit data collection here...129 130 except Exception as e:131 print(f"Erreur lors de la récupération des données sociales pour {symbol}: {e}")132 133 return social_data134 135 async def _fetch_web_traffic_data(self, symbols: List[str]) -> pd.DataFrame:136 """Récupération des données de trafic web"""137 traffic_data = pd.DataFrame()138 139 for symbol in symbols:140 try:141 # Similerweb API (si disponible)142 company_domain = self._get_company_domain(symbol)143 if company_domain:144 traffic_metrics = await self._get_similerweb_metrics(company_domain)145 traffic_data[symbol] = traffic_metrics146 147 except Exception as e:148 print(f"Erreur lors de la récupération du trafic web pour {symbol}: {e}")149 150 return traffic_data151 152 async def _fetch_satellite_data(self, symbols: List[str]) -> Dict:153 """Récupération des données satellite"""154 satellite_data = {}155 156 for symbol in symbols:157 try:158 # Données RS Data (exemple)159 company_locations = self._get_company_locations(symbol)160 if company_locations:161 satellite_metrics = await self._get_satellite_metrics(company_locations)162 satellite_data[symbol] = satellite_metrics163 164 except Exception as e:165 print(f"Erreur lors de la récupération des données satellite pour {symbol}: {e}")166 167 return satellite_data168 169 def _get_company_domain(self, symbol: str) -> Optional[str]:170 """Récupère le domaine principal de l'entreprise"""171 try:172 company_profile = self.finnhub_client.company_profile2(symbol=symbol)173 return company_profile.get('weburl', '').replace('www.', '').replace('https://', '')174 except:175 return None176 177 async def _get_similerweb_metrics(self, domain: str) -> pd.Series:178 """Récupère les métriques de trafic web via Similerweb"""179 # Implémentez ici la logique de récupération des données Similerweb180 # Nécessite un compte Similerweb API181 pass182 183 def _get_company_locations(self, symbol: str) -> List[Dict]:184 """Récupère les emplacements principaux de l'entreprise"""185 try:186 # Implémenter la logique de récupération des emplacements187 pass188 except:189 return []190 191 async def _get_satellite_metrics(self, locations: List[Dict]) -> Dict:192 """Récupère les métriques satellite pour les emplacements donnés"""193 # Implémenter la logique de récupération des données satellite194 pass195 196class SatelliteDataAnalyzer:197 """Analyse des données satellite pour indicateurs économiques"""198 199 async def get_retail_activity(self, locations: List[Dict]) -> Dict[str, float]:200 """Analyse l'activité de vente au détail via le trafic de parkings"""201 retail_metrics = {}202 203 for location in locations:204 try:205 # Récupération du comptage de voitures206 parking_count = await self._analyze_parking_occupancy(207 lat=location['latitude'],208 lon=location['longitude'],209 radius=500 # mètres210 )211 212 retail_metrics[location['name']] = {213 'parking_occupancy': parking_count / location['total_spaces'],214 'customer_traffic': self._estimate_customer_traffic(parking_count),215 'yoy_change': self._calculate_yoy_change(location['id'], parking_count)216 }217 except Exception as e:218 print(f"Erreur analyse retail pour {location['name']}: {e}")219 220 return retail_metrics221 222 async def get_industrial_activity(self, locations: List[Dict]) -> Dict[str, float]:223 """Analyse l'activité industrielle via imagerie thermique et pollution"""224 industrial_metrics = {}225 226 for location in locations:227 try:228 # Analyse de l'activité thermique229 heat_signature = await self._analyze_thermal_activity(230 lat=location['latitude'],231 lon=location['longitude']232 )233 234 # Analyse des émissions235 emissions = await self._analyze_emissions(236 lat=location['latitude'],237 lon=location['longitude']238 )239 240 industrial_metrics[location['name']] = {241 'activity_level': heat_signature['activity_score'],242 'emissions_level': emissions['level'],243 'operational_status': heat_signature['operational_status']244 }245 except Exception as e:246 print(f"Erreur analyse industrielle pour {location['name']}: {e}")247 248 return industrial_metrics249 250 async def get_shipping_activity(self, ports: List[Dict]) -> Dict[str, float]:251 """Analyse l'activité maritime et logistique"""252 shipping_metrics = {}253 254 for port in ports:255 try:256 # Analyse du trafic maritime257 vessel_count = await self._analyze_vessel_traffic(258 lat=port['latitude'],259 lon=port['longitude'],260 radius=5000 # mètres261 )262 263 # Analyse des conteneurs264 container_count = await self._analyze_container_stacks(265 lat=port['latitude'],266 lon=port['longitude']267 )268 269 shipping_metrics[port['name']] = {270 'vessel_occupancy': vessel_count / port['capacity'],271 'container_volume': container_count,272 'port_congestion': self._calculate_congestion_score(vessel_count, port['capacity'])273 }274 except Exception as e:275 print(f"Erreur analyse maritime pour {port['name']}: {e}")276 277 return shipping_metrics278 279class SocialMediaAnalyzer:280 """Analyse avancée des médias sociaux pour le sentiment des investisseurs"""281 282 async def analyze_investment_sentiment(self, symbols: List[str]) -> Dict[str, Dict]:283 """Analyse complète du sentiment des investisseurs"""284 sentiment_data = {}285 286 for symbol in symbols:287 try:288 # Reddit (r/wallstreetbets, r/stocks, etc.)289 reddit_sentiment = await self._analyze_reddit_sentiment(symbol)290 291 # StockTwits292 stocktwits_sentiment = await self._analyze_stocktwits_sentiment(symbol)293 294 # Twitter $cashtags295 twitter_sentiment = await self._analyze_twitter_sentiment(symbol)296 297 # Agrégation et normalisation298 sentiment_data[symbol] = {299 'overall_sentiment': self._aggregate_sentiment_scores([300 reddit_sentiment['score'],301 stocktwits_sentiment['score'],302 twitter_sentiment['score']303 ]),304 'sentiment_momentum': self._calculate_sentiment_momentum(symbol),305 'retail_interest': self._gauge_retail_interest(306 reddit_sentiment['volume'],307 stocktwits_sentiment['volume'],308 twitter_sentiment['volume']309 ),310 'institutional_hints': self._detect_institutional_activity(symbol)311 }312 313 except Exception as e:314 print(f"Erreur analyse sentiment pour {symbol}: {e}")315 316 return sentiment_data317 318 def _aggregate_sentiment_scores(self, scores: List[float]) -> float:319 """Agrège les scores de sentiment avec pondération"""320 weights = [0.4, 0.3, 0.3] # Pondération par importance de la source321 return sum(score * weight for score, weight in zip(scores, weights))322 323 def _calculate_sentiment_momentum(self, symbol: str) -> float:324 """Calcule la dynamique du sentiment"""325 # Implémentation de la logique de momentum326 pass327 328 def _gauge_retail_interest(self, *volumes: int) -> float:329 """Évalue l'intérêt des investisseurs particuliers"""330 return sum(volumes) / len(volumes) if volumes else 0331 332 def _detect_institutional_activity(self, symbol: str) -> Dict:333 """Détecte les signes d'activité institutionnelle"""334 # Analyse des ordres importants, des dark pools, etc.335 pass