26J29f2

Young-Kyoo Kim·2026년 6월 29일
"""
[신규 2단계] 프로메테우스 시계열 통합, 워크로드 매핑 및 Core-Hours 기반 파레토 빌드 엔진
실행: python step2_pipeline.py
"""

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

RAW_DIR = Path("./data/raw")
MERGED_DIR = Path("./data/merged")
MERGED_DIR.mkdir(parents=True, exist_ok=True)

def classify_workload(pod_name):
    """[검증 완료] 팟 실명 정규식 패턴 매핑 매크로"""
    p = str(pod_name).lower()
    if "spark" in p or "-exec-" in p or "-driver" in p:
        if "executor" in p or "-exec-" in p: return "SPARK_EXECUTOR"
        if "driver" in p or "-driver" in p: return "SPARK_DRIVER"
        return "SPARK_SYSTEM"
    if "airflow" in p:
        if "worker" in p: return "AIRFLOW_WORKER"
        if "scheduler" in p: return "AIRFLOW_SCHEDULER"
        if "dag-processor" in p: return "AIRFLOW_DAG_PROC"
        if "triggerer" in p: return "AIRFLOW_TRIGGERER"
        return "AIRFLOW_SYSTEM"
    if "starrocks" in p or "-be-" in p or "-fe-" in p or "-cn-" in p:
        if "-be" in p: return "STARROCKS_BE"
        if "-fe" in p: return "STARROCKS_FE"
        if "-cn" in p: return "STARROCKS_CN"
        return "STARROCKS_SYSTEM"
    if "postgres" in p or "pgpool" in p:
        if "backup" in p: return "POSTGRES_BACKUP"
        if "pgpool" in p: return "POSTGRES_POOL"
        return "POSTGRESQL"
    if "jupyter" in p or "notebook" in p: return "JUPYTERLAB"
    if "hyperset" in p or "superset" in p:
        if "worker" in p: return "HYPERSET_WORKER"
        return "HYPERSET"
    if "minio" in p or "aistor" in p: return "MINIO_STORAGE"
    if "cilium" in p: return "CILIUM_CNI"
    return "GENERAL_APPS"

def run_enrich_and_pareto_pipeline():
    raw_files = list(RAW_DIR.glob("prom_raw_*.parquet"))
    if not raw_files:
        print("❌ 에러: data/raw/ 폴더에 prom_raw_*.parquet 원천 파일이 없습니다.")
        return

    print(f"📦 총 {len(raw_files)}개의 시계열 청크 파일 병합 시작...")
    df_list = [pd.read_parquet(f) for f in raw_files]
    df_raw = pd.concat(df_list, ignore_index=True)
    
    # 날짜 필드 문자열 추출 (YYYY-MM-DD)
    df_raw["date"] = df_raw["timestamp"].dt.strftime("%Y-%m-%d")
    
    print("⚙️ 기술 워크로드 도메인 분류 레이블 주입 중...")
    df_raw["workload_type"] = df_raw["pod"].apply(classify_workload)

    print("📊 팟/컨테이너 유닛 레벨 시계열 데이터 롤업 집계 중...")
    # 💡 핵심: 1분 주기로 수집되었으므로, 그룹 내 데이터 카운트(size)가 곧 실제 가동 시간(분)이 됨!
    df_pod = df_raw.groupby(["date", "cluster", "namespace", "workload_type", "node", "pod", "container"]).agg(
        minutes_running      = ("timestamp", "size"), 
        cpu_request_max      = ("cpu_request", "max"),
        cpu_usage_p95        = ("cpu_usage", lambda x: x.quantile(0.95)),
        mem_request_max      = ("mem_request", "max"),
        mem_usage_p95        = ("mem_usage", lambda x: x.quantile(0.95))
    ).reset_index()

    # 원천 메트릭 바이트 단위를 GB 스케일로 정규화
    df_pod["mem_request_max"] = df_pod["mem_request_max"] / (1024**3)
    df_pod["mem_usage_p95"] = df_pod["mem_usage_p95"] / (1024**3)

    print("⚙️ 시간 가중치 고정밀 [Core-Hours / GB-Hours] 누수 정량화 연산 중...")
    df_pod["cpu_allocated_core_hours"] = df_pod["cpu_request_max"] * (df_pod["minutes_running"] / 60.0)
    df_pod["cpu_usage_core_hours"]     = df_pod["cpu_usage_p95"] * (df_pod["minutes_running"] / 60.0)
    df_pod["cpu_waste_core_hours"]     = (df_pod["cpu_allocated_core_hours"] - df_pod["cpu_usage_core_hours"]).clip(lower=0)

    df_pod["mem_allocated_gb_hours"]   = df_pod["mem_request_max"] * (df_pod["minutes_running"] / 60.0)
    df_pod["mem_usage_gb_hours"]       = df_pod["mem_usage_p95"] * (df_pod["minutes_running"] / 60.0)
    df_pod["mem_waste_gb_hours"]       = (df_pod["mem_allocated_gb_hours"] - df_pod["mem_usage_gb_hours"]).clip(lower=0)

    # 공급 부족 및 리스크 팟 판정 수식 (양수 스펙 트래킹)
    df_pod["cpu_shortage_cores"] = (df_pod["cpu_usage_p95"] - df_pod["cpu_request_max"]).clip(lower=0)
    df_pod["status"] = np.where(df_pod["cpu_shortage_cores"] > 0.5, "⚠️ Request부족", 
                        np.where(df_pod["cpu_waste_core_hours"] > 10, "📉 과다할당", "✅ 최적화완료"))

    # 마스터 가공 원부 세이브
    df_pod.to_parquet(MERGED_DIR / "enriched_fixed_7d.parquet", index=False)
    print("💾 1차 마스터 롤업 원부 [enriched_fixed_7d.parquet] 생성 완료.")

    # ── 📈 2차 Namespace 단위 거버넌스 파레토 원부 계산 ──
    print("📈 Namespace별 파레토 누적 기여도 산출 중...")
    df_ns = df_pod.groupby("namespace").agg(
        minutes_running_sum = ("minutes_running", "sum"),
        container_cnt       = ("container", "count"),
        total_allocated_core_hours = ("cpu_allocated_core_hours", "sum"),
        total_waste_core_hours     = ("cpu_waste_core_hours", "sum")
    ).reset_index().sort_values(by="total_waste_core_hours", ascending=False).reset_index(drop=True)

    global_total_waste = 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 * 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("💾 2차 거버넌스 파레토 원부 [pareto_fixed_ns.parquet] 생성 완료.\n")

if __name__ == "__main__":
    run_enrich_and_pareto_pipeline()

0개의 댓글