26J29e2

Young-Kyoo Kim·2026년 6월 28일
"""
[2.5단계-최종 패치] 스키마 정규화(Node/Date 에일리어스 매핑) 및 Core-Hours 왜곡 보정 엔진
실행: python step2_5_fix_pareto.py
"""

import pandas as pd
import numpy as np
from pathlib import Path

# 공통 프로젝트 경로 설정 구성원 로드
from config import MERGED_DIR

def fix_opencost_time_and_recalculate_pareto():
    src_path = MERGED_DIR / "enriched_7d.parquet"
    if not src_path.exists():
        print(f"❌ 원천 enriched_7d.parquet 파일이 없습니다. 경로를 확인하세요: {src_path}")
        return

    # 1. 데이터 로드
    df = pd.read_parquet(src_path)
    print(f"📊 원천 데이터 로드 완료 (총 행수: {len(df):,}개)")
    print(f"🔍 현재 데이터프레임 보유 컬럼 목록: {list(df.columns)}")

    # ── [🛡️ 패치 A] 'node' 관련 KeyError 방어 및 에일리어스 매핑 ──
    if "node" not in df.columns:
        node_aliases = ["node_name", "nodename", "host", "hostname", "properties.node", "kubernetes_node"]
        found_node_col = None
        for alias in node_aliases:
            if alias in df.columns:
                found_node_col = alias
                break
        
        if found_node_col:
            df["node"] = df[found_node_col].astype(str)
            print(f"  -> ✅ 포착된 '{found_node_col}' 컬럼을 표준 'node' 필드로 변환 매핑했습니다.")
        else:
            df["node"] = "unknown-node-slot"
            print("  -> ⚠️ [경고] 노드 관련 컬럼을 찾지 못해 'unknown-node-slot'으로 임시 대체합니다. (step1/2 파싱 확인 필요)")

    # ── [🛡️ 패치 B] 'date' 관련 KeyError 방어 및 오픈코스트 window 변환 ──
    if "date" not in df.columns:
        date_aliases = ["timestamp", "time", "window", "day", "start_time", "start"]
        found_date_col = None
        for alias in date_aliases:
            if alias in df.columns:
                found_date_col = alias
                break
        
        if found_date_col:
            if found_date_col == "window":
                # OpenCost 'window' 포맷 예시: "2026-06-28T00:00:00Z/2026-06-29T00:00:00Z" -> 앞자리 날짜만 정제
                df["date"] = df["window"].astype(str).str.split("T").str[0]
                print("  -> ✅ OpenCost 'window' 시계열 필드에서 표준 'date' 문자열을 추출했습니다.")
            else:
                df["date"] = df[found_date_col].astype(str).str.split(" ").str[0]
                print(f"  -> ✅ 포착된 '{found_date_col}' 컬럼을 표준 'date' 필드로 변환 매핑했습니다.")
        else:
            df["date"] = "2026-06-29"
            print("  -> ⚠️ [경고] 날짜 관련 컬럼을 찾지 못해 '2026-06-29' 기본값으로 강제 주입합니다.")

    # ── [🛡️ 패치 C] minutes_running 컬럼 검증 및 초기화 ──
    if "minutes_running" not in df.columns:
        if "minutes" in df.columns:
            df["minutes_running"] = df["minutes"].astype(float)
        else:
            df["minutes_running"] = 0.0

    # 가동 시간 마비 상태 복구 (정적 10,080분 / 동적 5분 기조 주입)
    if (df["minutes_running"] == 0).all() or df["minutes_running"].isna().all():
        df['workload_type'] = df['workload_type'].str.upper()
        df["minutes_running"] = np.where(
            df['workload_type'].isin(["STARROCKS", "POSTGRESQL", "JUPYTERLAB", "AIRFLOW"]),
            10080.0,
            np.where((df["minutes_running"] == 0) | (df["minutes_running"].isna()), 5.0, df["minutes_running"])
        )
        print("✅ 가동 시간(minutes_running) 정상화 처리를 완료했습니다.")

    # ── [핵심 연산] Core-Hours 기반 가중치 변환 ──
    print("⚙️ 시간 가중치(Core-Hours) 연산 적용 중...")
    df["cpu_request_max"] = df["cpu_request_max"].astype(float)
    df["cpu_usage_p95"] = df["cpu_usage_p95"].astype(float)
    
    df["cpu_allocated_core_hours"] = df["cpu_request_max"] * (df["minutes_running"] / 60.0)
    df["cpu_usage_core_hours"] = df["cpu_usage_p95"] * (df["minutes_running"] / 60.0)
    df["cpu_waste_core_hours"] = (df["cpu_allocated_core_hours"] - df["cpu_usage_core_hours"]).clip(lower=0)

    if "mem_request_max" in df.columns and "mem_usage_p95" in df.columns:
        df["mem_allocated_gb_hours"] = df["mem_request_max"].astype(float) * (df["minutes_running"] / 60.0)
        df["mem_usage_gb_hours"] = df["mem_usage_p95"].astype(float) * (df["minutes_running"] / 60.0)

    # 파이프라인 무결성을 보장하는 수정된 파일 포맷 저장
    df.to_parquet(MERGED_DIR / "enriched_fixed_7d.parquet", index=False)
    print("💾 정규화가 완료된 마스터 파티션 데이터가 enriched_fixed_7d.parquet에 병합되었습니다.")
    
    # ── Namespace별 파레토 분석 맵 세이브 ──
    df_ns = df.groupby("namespace").agg(
        minutes_running_sum   = ("minutes_running", "sum"),
        container_cnt         = ("container", "count"),
        total_request_cores   = ("cpu_request_max", "sum"),
        total_waste_cores_raw = ("cpu_waste_cores", "sum"),
        total_allocated_core_hours = ("cpu_allocated_core_hours", "sum"),
        total_waste_core_hours     = ("cpu_waste_core_hours", "sum")
    ).reset_index()
    
    df_ns = df_ns.sort_values(by="total_waste_core_hours", ascending=False).reset_index(drop=True)
    global_total_waste_core_hours = df_ns["total_waste_core_hours"].sum() if df_ns["total_waste_core_hours"].sum() > 0 else 0.1
    df_ns["waste_share_pct"] = (df_ns["total_waste_core_hours"] / global_total_waste_core_hours * 100).round(2)
    df_ns["waste_cumsum_pct"] = df_ns["waste_share_pct"].cumsum().round(2)
    df_ns["is_top_80_percent_offender"] = df_ns["waste_cumsum_pct"] <= 80.5

    df_ns.to_parquet(MERGED_DIR / "pareto_fixed_ns.parquet", index=False)
    print("🚀 [성공] 스키마 변환 및 파레토 최적화 서열 매핑이 완료되었습니다.")

if __name__ == "__main__":
    fix_opencost_time_and_recalculate_pareto()

0개의 댓글