fix: daily stats calculation logic and daily close task
This commit is contained in:
parent
c885bab007
commit
c9ed91d885
|
|
@ -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
|
||||
|
|
|
|||
89
database.py
89
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")
|
||||
# 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
|
||||
|
||||
for item in data_list:
|
||||
plant_id = item.get('id', '')
|
||||
if not plant_id:
|
||||
continue
|
||||
# 야간 시간대 차단 (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 = []
|
||||
|
||||
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
|
||||
|
||||
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:
|
||||
# 기존 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")
|
||||
|
|
|
|||
30
main.py
30
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)
|
||||
|
||||
|
|
|
|||
Loading…
Reference in New Issue
Block a user