""" FRED (Federal Reserve Economic Data) Service FRED API 호출, 캐싱, 일일 1000개 제한 관리 """ import logging import httpx from datetime import datetime, timedelta from typing import Dict, List, Optional, Tuple from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy import select, func, and_, desc from app.models.fred_data import FredSeries, FredObservation, FredApiUsage, FredCacheStats from app.core.config import settings logger = logging.getLogger(__name__) class FredService: """FRED API 서비스 with intelligent caching and daily limit management""" def __init__(self): self.api_key = "2b12c4c62a7e9d9002d746dad7bfd147" self.base_url = "https://api.stlouisfed.org/fred" self.daily_limit = 1000 self.cache_duration_hours = 24 # 24시간 캐시 async def _check_daily_limit(self, db: AsyncSession) -> Tuple[bool, int, int]: """ 일일 API 사용량 확인 Returns: (can_make_request, used_today, remaining) """ today = datetime.now().strftime('%Y-%m-%d') # 오늘의 API 사용량 조회 result = await db.execute( select(func.count(FredApiUsage.id)) .where(and_( FredApiUsage.date == today, FredApiUsage.success == True )) ) used_today = result.scalar() or 0 remaining = self.daily_limit - used_today can_make_request = remaining > 0 logger.debug(f"📊 FRED API usage today: {used_today}/{self.daily_limit} (remaining: {remaining})") return can_make_request, used_today, remaining async def _log_api_usage( self, db: AsyncSession, endpoint: str, series_id: Optional[str] = None, request_params: Optional[Dict] = None, success: bool = True, response_size: int = 0 ): """API 사용량 로깅""" today = datetime.now().strftime('%Y-%m-%d') usage_log = FredApiUsage( date=today, endpoint=endpoint, series_id=series_id, request_params=request_params or {}, success=success, response_size=response_size ) db.add(usage_log) await db.commit() logger.info(f"📝 FRED API call logged: {endpoint} {'✅' if success else '❌'}") async def _is_cache_valid(self, cached_at: datetime, cache_hours: int = 24) -> bool: """캐시 유효성 확인""" if not cached_at: return False expiry_time = cached_at + timedelta(hours=cache_hours) return datetime.now() < expiry_time async def _make_fred_request(self, endpoint: str, params: Dict) -> Optional[Dict]: """FRED API 요청 실행""" url = f"{self.base_url}/{endpoint}" params['api_key'] = self.api_key params['file_type'] = 'json' try: async with httpx.AsyncClient(timeout=30.0) as client: response = await client.get(url, params=params) response.raise_for_status() data = response.json() logger.info(f"✅ FRED API success: {endpoint} -> {len(data.get('seriess', data.get('observations', [])))} records") return data except httpx.HTTPError as e: logger.error(f"❌ FRED API error: {endpoint} -> {e}") return None except Exception as e: logger.error(f"❌ FRED API unexpected error: {endpoint} -> {e}") return None async def get_series_info(self, db: AsyncSession, series_id: str, force_refresh: bool = False) -> Optional[Dict]: """ FRED 시리즈 정보 조회 (캐시 우선) Args: db: Database session series_id: FRED series ID (e.g., "GDP") force_refresh: 캐시 무시하고 API 호출 Returns: Series information dict or None """ try: # 1. 캐시에서 조회 (force_refresh가 아닌 경우) if not force_refresh: cached_series = await db.execute( select(FredSeries).where(FredSeries.id == series_id) ) cached = cached_series.scalar_one_or_none() if cached and await self._is_cache_valid(cached.cached_at): logger.info(f"📦 FRED series cache hit: {series_id}") return { 'id': cached.id, 'title': cached.title, 'units': cached.units, 'frequency': cached.frequency, 'last_updated': cached.last_updated.isoformat() if cached.last_updated else None, 'cached': True, 'cached_at': cached.cached_at.isoformat() } # 2. API 호출 가능 여부 확인 can_call, used, remaining = await self._check_daily_limit(db) if not can_call: logger.warning(f"🚫 FRED API daily limit reached: {used}/{self.daily_limit}") # 캐시된 데이터라도 반환 if not force_refresh: cached_series = await db.execute( select(FredSeries).where(FredSeries.id == series_id) ) cached = cached_series.scalar_one_or_none() if cached: return { 'id': cached.id, 'title': cached.title, 'units': cached.units, 'frequency': cached.frequency, 'last_updated': cached.last_updated.isoformat() if cached.last_updated else None, 'cached': True, 'cache_expired': True, 'api_limit_reached': True } return None # 3. FRED API 호출 params = {'series_id': series_id} response_data = await self._make_fred_request('series', params) if not response_data or 'seriess' not in response_data: await self._log_api_usage(db, 'series', series_id, params, False, 0) return None # 4. 응답 처리 및 캐시 저장 series_data = response_data['seriess'][0] if response_data['seriess'] else None if not series_data: await self._log_api_usage(db, 'series', series_id, params, False, 0) return None # 5. 데이터베이스에 저장/업데이트 cached_series = await db.execute( select(FredSeries).where(FredSeries.id == series_id) ) existing = cached_series.scalar_one_or_none() if existing: # 업데이트 existing.title = series_data.get('title') existing.units = series_data.get('units') existing.units_short = series_data.get('units_short') existing.frequency = series_data.get('frequency') existing.frequency_short = series_data.get('frequency_short') existing.seasonal_adjustment = series_data.get('seasonal_adjustment') existing.seasonal_adjustment_short = series_data.get('seasonal_adjustment_short') existing.last_updated = datetime.fromisoformat(series_data['last_updated'].replace('-05', '')) if series_data.get('last_updated') else None existing.popularity = series_data.get('popularity', 0) existing.notes = series_data.get('notes') existing.cached_at = datetime.now() existing.cache_expires_at = datetime.now() + timedelta(hours=self.cache_duration_hours) existing.fred_metadata = series_data else: # 새로 생성 new_series = FredSeries( id=series_data['id'], title=series_data.get('title'), units=series_data.get('units'), units_short=series_data.get('units_short'), frequency=series_data.get('frequency'), frequency_short=series_data.get('frequency_short'), seasonal_adjustment=series_data.get('seasonal_adjustment'), seasonal_adjustment_short=series_data.get('seasonal_adjustment_short'), last_updated=datetime.fromisoformat(series_data['last_updated'].replace('-05', '')) if series_data.get('last_updated') else None, popularity=series_data.get('popularity', 0), notes=series_data.get('notes'), cached_at=datetime.now(), cache_expires_at=datetime.now() + timedelta(hours=self.cache_duration_hours), fred_metadata=series_data ) db.add(new_series) await db.commit() await self._log_api_usage(db, 'series', series_id, params, True, 1) logger.info(f"✅ FRED series cached: {series_id} - {series_data.get('title', 'Unknown')[:50]}") return { 'id': series_data['id'], 'title': series_data.get('title'), 'units': series_data.get('units'), 'frequency': series_data.get('frequency'), 'last_updated': series_data.get('last_updated'), 'cached': False, 'api_calls_remaining': remaining - 1 } except Exception as e: logger.error(f"❌ Error in get_series_info: {e}") return None async def get_series_observations( self, db: AsyncSession, series_id: str, start_date: Optional[str] = None, end_date: Optional[str] = None, limit: Optional[int] = None, force_refresh: bool = False ) -> Optional[Dict]: """ FRED 시리즈 관측값 조회 (캐시 우선) Args: db: Database session series_id: FRED series ID start_date: YYYY-MM-DD format end_date: YYYY-MM-DD format limit: 최대 반환 개수 force_refresh: 캐시 무시하고 API 호출 Returns: Observations data dict or None """ try: # 1. 캐시에서 조회 (force_refresh가 아닌 경우) if not force_refresh: query = select(FredObservation).where(FredObservation.series_id == series_id) if start_date: query = query.where(FredObservation.date >= start_date) if end_date: query = query.where(FredObservation.date <= end_date) query = query.order_by(desc(FredObservation.date)) if limit: query = query.limit(limit) cached_obs = await db.execute(query) cached_data = cached_obs.scalars().all() # 캐시가 있고 최근 데이터인지 확인 if cached_data and await self._is_cache_valid(cached_data[0].cached_at): logger.info(f"📦 FRED observations cache hit: {series_id} ({len(cached_data)} records)") return { 'series_id': series_id, 'observations': [ { 'date': obs.date, 'value': obs.value, 'realtime_start': obs.realtime_start, 'realtime_end': obs.realtime_end } for obs in cached_data ], 'count': len(cached_data), 'cached': True, 'cached_at': cached_data[0].cached_at.isoformat() if cached_data else None } # 2. API 호출 가능 여부 확인 can_call, used, remaining = await self._check_daily_limit(db) if not can_call: logger.warning(f"🚫 FRED API daily limit reached: {used}/{self.daily_limit}") return None # 3. FRED API 호출 params = {'series_id': series_id} if start_date: params['observation_start'] = start_date if end_date: params['observation_end'] = end_date if limit: params['limit'] = str(limit) response_data = await self._make_fred_request('series/observations', params) if not response_data or 'observations' not in response_data: await self._log_api_usage(db, 'observations', series_id, params, False, 0) return None observations = response_data['observations'] # 4. 새 데이터만 추가 (기존 데이터는 유지 - 영구 저장) existing_dates = set() existing_query = await db.execute( select(FredObservation.date).where(FredObservation.series_id == series_id) ) existing_dates = {row[0] for row in existing_query.fetchall()} # 5. 새로운 관측값만 저장 (중복 방지) new_observations_count = 0 for obs_data in observations: obs_date = obs_data['date'] if obs_date not in existing_dates: new_obs = FredObservation( series_id=series_id, date=obs_date, value=obs_data['value'], realtime_start=obs_data.get('realtime_start'), realtime_end=obs_data.get('realtime_end'), cached_at=datetime.now() ) db.add(new_obs) new_observations_count += 1 else: # 기존 데이터의 cached_at 업데이트 (최신성 표시) update_query = await db.execute( select(FredObservation).where( and_( FredObservation.series_id == series_id, FredObservation.date == obs_date ) ) ) existing_obs = update_query.scalar_one_or_none() if existing_obs: existing_obs.cached_at = datetime.now() await db.commit() await self._log_api_usage(db, 'observations', series_id, params, True, len(observations)) total_in_db = len(existing_dates) + new_observations_count logger.info(f"✅ FRED observations processed: {series_id} ({new_observations_count} new, {total_in_db} total in DB)") return { 'series_id': series_id, 'observations': observations, 'count': len(observations), 'cached': False, 'api_calls_remaining': remaining - 1 } except Exception as e: logger.error(f"❌ Error in get_series_observations: {e}") return None async def get_api_usage_stats(self, db: AsyncSession, days: int = 7) -> Dict: """API 사용량 통계 조회""" try: end_date = datetime.now() start_date = end_date - timedelta(days=days) # 기간별 사용량 조회 - 간단한 방식으로 변경 usage_query = await db.execute( select( FredApiUsage.date, func.count(FredApiUsage.id).label('total_calls'), func.count(FredApiUsage.id).filter(FredApiUsage.success == True).label('successful_calls'), func.sum(FredApiUsage.response_size).label('total_records') ) .where(FredApiUsage.date >= start_date.strftime('%Y-%m-%d')) .group_by(FredApiUsage.date) .order_by(desc(FredApiUsage.date)) ) daily_stats = [] for row in usage_query.fetchall(): daily_stats.append({ 'date': row.date, 'total_calls': row.total_calls or 0, 'successful_calls': row.successful_calls or 0, 'total_records': row.total_records or 0, 'success_rate': (row.successful_calls / row.total_calls * 100) if row.total_calls > 0 else 0 }) # 오늘의 사용량 today = datetime.now().strftime('%Y-%m-%d') can_call, used_today, remaining = await self._check_daily_limit(db) # 캐시 통계 cache_stats_query = await db.execute( select( func.count(FredSeries.id).label('cached_series'), func.count(FredObservation.id).label('cached_observations') ) ) cache_row = cache_stats_query.fetchone() return { 'success': True, 'data': { 'daily_limit': self.daily_limit, 'used_today': used_today, 'remaining_today': remaining, 'usage_percentage': (used_today / self.daily_limit * 100), 'can_make_requests': can_call, 'daily_stats': daily_stats, 'cache_stats': { 'cached_series': cache_row.cached_series or 0, 'cached_observations': cache_row.cached_observations or 0, 'cache_duration_hours': self.cache_duration_hours } } } except Exception as e: logger.error(f"❌ Error getting API usage stats: {e}") return {'success': False, 'error': str(e)} # 싱글톤 인스턴스 fred_service = FredService()