diff --git a/daily_summary.py b/daily_summary.py index 6c999aa..057d432 100644 --- a/daily_summary.py +++ b/daily_summary.py @@ -3,7 +3,7 @@ # ========================================== # solar_logs 데이터를 집계하여 daily_stats 테이블에 저장 -from datetime import datetime, timedelta +from datetime import datetime, timedelta, timezone try: from dotenv import load_dotenv @@ -33,7 +33,8 @@ def calculate_daily_stats(date_str: str = None): date_str: 집계 대상 날짜 (YYYY-MM-DD). 미지정 시 오늘. """ if date_str is None: - date_str = datetime.now().strftime('%Y-%m-%d') + kst = timezone(timedelta(hours=9)) + date_str = datetime.now(kst).strftime('%Y-%m-%d') print(f"\n📊 [일일 통계 집계] {date_str}") print("-" * 60) @@ -46,15 +47,18 @@ def calculate_daily_stats(date_str: str = None): # 1. 용량 정보 조회 capacities = get_plant_capacities(client) - # 2. 해당일 로그 조회 - start_dt = f"{date_str}T00:00:00" - end_dt = f"{date_str}T23:59:59" + # 2. 해당일 로그 조회 (KST 날짜 범위를 UTC로 변환하여 쿼리) + kst = timezone(timedelta(hours=9)) + start_kst = datetime.strptime(f"{date_str} 00:00:00", "%Y-%m-%d %H:%M:%S").replace(tzinfo=kst) + end_kst = datetime.strptime(f"{date_str} 23:59:59", "%Y-%m-%d %H:%M:%S").replace(tzinfo=kst) + start_utc = start_kst.astimezone(timezone.utc).isoformat() + end_utc = end_kst.astimezone(timezone.utc).isoformat() try: result = client.table("solar_logs") \ .select("plant_id, current_kw, today_kwh, created_at") \ - .gte("created_at", start_dt) \ - .lte("created_at", end_dt) \ + .gte("created_at", start_utc) \ + .lte("created_at", end_utc) \ .order("created_at", desc=False) \ .execute() @@ -72,8 +76,8 @@ def calculate_daily_stats(date_str: str = None): stats_list = [] for plant_id, group in df.groupby('plant_id'): - # 마지막 로그의 today_kwh - total_generation = group['today_kwh'].iloc[-1] if len(group) > 0 else 0 + # 당일 최댓값 today_kwh 사용 (마지막 로그가 아닌 최댓값으로 중간 리셋 보호) + total_generation = group['today_kwh'].max() if len(group) > 0 else 0 # 최대 출력 peak_kw = group['current_kw'].max() if len(group) > 0 else 0 diff --git a/database.py b/database.py index 62fbb2b..42939cf 100644 --- a/database.py +++ b/database.py @@ -100,40 +100,69 @@ def save_to_supabase(data_list): # [보호 로직] # 1. 오류 상태(크롤링 실패)인 경우 daily_stats 갱신 금지 # 2. today_kwh == 0인 경우 daily_stats 갱신 금지 (새벽 0 값으로 하루치 덮어쓰기 방지) - daily_records = [] - kst_date_str = datetime.now(timezone(timedelta(hours=9))).strftime("%Y-%m-%d") - - for item in data_list: - plant_id = item.get('id', '') - if not plant_id: - continue - - status = item.get('status', '') - is_error = '오류' in status - today_val = float(item.get('today', 0)) + # 3. 야간 시간대(21:00~06:00 KST) daily_stats 갱신 금지 (일몰 이후 잔류값 보호) + # 4. DB에 이미 저장된 값보다 작은 경우 갱신 금지 (최댓값 보호) + kst = timezone(timedelta(hours=9)) + kst_now_dt = datetime.now(kst) + kst_date_str = kst_now_dt.strftime("%Y-%m-%d") + kst_hour = kst_now_dt.hour - # 오류 상태이거나 today_kwh가 0이면 daily_stats 갱신 건너뜀 - if is_error: - print(f" ⚠️ [{plant_id}] 오류 상태 → daily_stats 갱신 건너뜀") - continue - if today_val == 0: - print(f" ⚠️ [{plant_id}] today_kwh=0 → daily_stats 갱신 건너뜀 (새벽/야간 추정)") - continue - - daily_records.append({ - "plant_id": plant_id, - "date": kst_date_str, - "total_generation": today_val, - "created_at": kst_now - # updated_at은 자동으로 NOW()로 설정됨 (DB 기본값) - }) - - if daily_records: + # 야간 시간대 차단 (21:00 ~ 익일 06:00 KST) + is_night = kst_hour >= 21 or kst_hour < 6 + if is_night: + print(f" ⚠️ [야간 차단] KST {kst_hour:02d}시 → daily_stats 갱신 건너뜀 (일몰 후 잔류값 보호)") + else: + daily_records = [] + + # 기존 daily_stats 값 조회 (MAX 보호용) try: - stats_result = client.table("daily_stats").upsert(daily_records, on_conflict="plant_id, date").execute() - print(f"✅ [DB] daily_stats 업데이트 완료: {len(daily_records)}건") + existing_res = client.table("daily_stats") \ + .select("plant_id, total_generation") \ + .eq("date", kst_date_str) \ + .execute() + existing_map = {row['plant_id']: float(row.get('total_generation') or 0) + for row in existing_res.data} except Exception as e: - print(f"⚠️ [DB] daily_stats 업데이트 실패: {e}") + print(f" ⚠️ [DB] 기존 daily_stats 조회 실패: {e}") + existing_map = {} + + for item in data_list: + plant_id = item.get('id', '') + if not plant_id: + continue + + status = item.get('status', '') + is_error = '오류' in status + today_val = float(item.get('today', 0)) + + # 오류 상태이거나 today_kwh가 0이면 daily_stats 갱신 건너뜀 + if is_error: + print(f" ⚠️ [{plant_id}] 오류 상태 → daily_stats 갱신 건너뜀") + continue + if today_val == 0: + print(f" ⚠️ [{plant_id}] today_kwh=0 → daily_stats 갱신 건너뜀 (새벽/야간 추정)") + continue + + # [MAX 보호] 기존 값보다 작으면 갱신 건너뜀 + existing_val = existing_map.get(plant_id, 0) + if today_val <= existing_val: + print(f" ⚠️ [{plant_id}] 신규({today_val:.1f}) ≤ 기존({existing_val:.1f}) → daily_stats 갱신 건너뜀 (최댓값 보호)") + continue + + daily_records.append({ + "plant_id": plant_id, + "date": kst_date_str, + "total_generation": today_val, + "created_at": kst_now + # updated_at은 자동으로 NOW()로 설정됨 (DB 기본값) + }) + + if daily_records: + try: + stats_result = client.table("daily_stats").upsert(daily_records, on_conflict="plant_id, date").execute() + print(f"✅ [DB] daily_stats 업데이트 완료: {len(daily_records)}건") + except Exception as e: + print(f"⚠️ [DB] daily_stats 업데이트 실패: {e}") for r in records: print(f" → {r['plant_id']}: {r['current_kw']} kW / {r['today_kwh']} kWh") diff --git a/main.py b/main.py index be750c6..73a12fa 100644 --- a/main.py +++ b/main.py @@ -3,7 +3,7 @@ # ========================================== import re -from datetime import datetime +from datetime import datetime, timezone, timedelta # 환경 변수 로드 (최상단에서 실행) try: @@ -162,14 +162,42 @@ def integrated_monitoring(save_to_db=True, company_filter=None, force_run=False) return total_results + +def run_daily_close(force=False): + """ + 일일 마감 집계 실행 (KST 21:00~21:10 자동 트리거 또는 force=True) + solar_logs 데이터를 집계하여 daily_stats에 당일 최종값을 확정합니다. + """ + kst = timezone(timedelta(hours=9)) + kst_now = datetime.now(kst) + kst_hour = kst_now.hour + kst_minute = kst_now.minute + + is_close_window = (kst_hour == 21 and kst_minute < 10) + + if force or is_close_window: + date_str = kst_now.strftime('%Y-%m-%d') + print(f"\n🌙 [KST {kst_hour:02d}:{kst_minute:02d}] 일일 마감 집계 트리거 → {date_str} 통계 확정") + try: + from daily_summary import calculate_daily_stats + calculate_daily_stats(date_str) + except Exception as e: + print(f" ❌ 마감 집계 실패: {e}") + else: + print(f" ℹ️ 마감 집계 스킵 (KST {kst_hour:02d}시 {kst_minute:02d}분, 대상 시간대 아님)") + if __name__ == "__main__": import sys - + # 인자 처리: --force 옵션으로 스케줄러 무시 force_run = '--force' in sys.argv or '-f' in sys.argv - + force_close = '--close' in sys.argv # 마감 집계 강제 실행 + if force_run: print("⚡ [강제 실행 모드] 스케줄러 무시하고 모든 사이트 크롤링") - + integrated_monitoring(save_to_db=True, force_run=force_run) + # 마감 집계: 21:00~21:10 KST 자동 실행 또는 --close 옵션 + run_daily_close(force=force_close) +