Team Ai
Apppublic

Deepvest/ProfilingAI

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
data_fetcher.py335 linesDownload Raw Back to core
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