26J29g1

Young-Kyoo Kim·2026년 6월 29일
queries = {
        "cpu_request": f"kube_pod_container_resource_requests{selector.replace('container!=', 'resource=\"cpu\", container!=')}",
        "cpu_limit":   f"kube_pod_container_resource_limits{selector.replace('container!=', 'resource=\"cpu\", container!=')}", # ◀ 추가
        "cpu_usage":   f"rate(container_cpu_usage_seconds_total{selector}[1m])",
        "mem_request": f"kube_pod_container_resource_requests{selector.replace('container!=', 'resource=\"memory\", container!=')}",
        "mem_limit":   f"kube_pod_container_resource_limits{selector.replace('container!=', 'resource=\"memory\", container!=')}", # ◀ 추가
        "mem_usage":   f"container_memory_working_set_bytes{selector}",
        "oom_event":   f"kube_pod_container_status_terminated_reason{selector.replace('container!=', 'reason=\"OOMKilled\", container!=')}" # ◀ 추가
    }
    
    
===

"""
[신규 2단계 - 거버넌스 강화판] OOMKilled 및 Request/Limit 위반 검출 엔진
실행: 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 derive_cluster_from_node(node_name):
    n = str(node_name).lower()
    if "icdlh" in n or "prod" in n: return "prod-data-lakehouse"
    elif "stage" in n or "stg" in n: return "stage-lakehouse"
    return "dev-sandbox"

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"
        return "AIRFLOW_SYSTEM"
    if "starrocks" in p or "-be-" in p or "-fe-" in p: return "STARROCKS_BE" if "-be" in p else "STARROCKS_FE"
    if "postgres" in p or "pgpool" in p: return "POSTGRESQL"
    if "jupyter" in p or "notebook" in p: return "JUPYTERLAB"
    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/ 폴더에 원천 파일이 없습니다.")
        return

    df_raw = pd.concat([pd.read_parquet(f) for f in raw_files], ignore_index=True)
    df_raw["date"] = df_raw["timestamp"].dt.strftime("%Y-%m-%d")
    df_raw["cluster"] = df_raw["node"].apply(derive_cluster_from_node)
    df_raw["workload_type"] = df_raw["pod"].apply(classify_workload)

    # 🛡️ 방어 코드: 신규 징수 컬럼이 없을 경우를 대비한 기본값 핸들링
    for col in ["cpu_limit", "mem_limit", "oom_event"]:
        if col not in df_raw.columns:
            df_raw[col] = 0.0

    print("📊 팟/컨테이너 유닛 레벨 시계열 데이터 가공 롤업 중...")
    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_limit_max        = ("cpu_limit", "max"),        # ◀ 추가
        cpu_usage_p95        = ("cpu_usage", lambda x: x.quantile(0.95)),
        mem_request_max      = ("mem_request", "max"),
        mem_limit_max        = ("mem_limit", "max"),        # ◀ 추가
        mem_usage_p95        = ("mem_usage", lambda x: x.quantile(0.95)),
        oom_strike_sum       = ("oom_event", "sum")         # ◀ 추가 (0보다 크면 해당 기간 OOM 발생)
    ).reset_index().fillna(0)

    # 바이트 단위를 GB 스케일로 정규화
    for col in ["mem_request_max", "mem_limit_max", "mem_usage_p95"]:
        df_pod[col] = df_pod[col] / (1024**3)

    # 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["is_oom_killed"] = df_pod["oom_strike_sum"] > 0
    df_pod["has_no_request"] = (df_pod["cpu_request_max"] == 0) | (df_pod["mem_request_max"] == 0)
    df_pod["has_no_limit"] = (df_pod["cpu_limit_max"] == 0) | (df_pod["mem_limit_max"] == 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["is_oom_killed"], "💥 OOM장애발생",
                        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)
    
    # 파레토 생성 로직 생략 (기존 스크립트와 동일)
    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.to_parquet(MERGED_DIR / "pareto_fixed_ns.parquet", index=False)
    print("✅ [2단계 완료] OOM 및 거버넌스 원천 플래그 빌드 완수.")

if __name__ == "__main__":
    run_enrich_and_pareto_pipeline()


===

"""
[신규 6단계 - 동적 30% 반영판] 4대 거버넌스 탭 통합 마스터 엑셀 빌더
실행: python step6_excel_builder.py
"""

import pandas as pd
from pathlib import Path

MERGED_DIR = Path("./data/merged")
OUT_DIR = Path("./data/output")
OUT_DIR.mkdir(parents=True, exist_ok=True)

def build_master_excel_report():
    df_pod = pd.read_parquet(MERGED_DIR / "enriched_fixed_7d.parquet")
    df_ns = pd.read_parquet(MERGED_DIR / "pareto_fixed_ns.parquet")

    # ── 📝 탭 0. 전사 요약 데이터프레임 빌드 ──
    summary_data = {
        "인프라 거버넌스 총괄 지표": [
            "관측된 총 고유 컨테이너 수 (개)",
            "자원 최적화 권고 대상 범위 (전체 Pod의 30% 규격 수량)",
            "클러스터 내 리소스 미설정 위반 팟 수 (개)",
            "해당 주간 OOMKilled 장애 발생 팟 수 (개)"
        ],
        "실측 통계치": [
            len(df_pod),
            int(len(df_pod) * 0.30),
            len(df_pod[df_pod["has_no_request"] | df_pod["has_no_limit"]]),
            len(df_pod[df_pod["is_oom_killed"] == True])
        ]
    }
    df_summary = pd.DataFrame(summary_data)

    # ── 📝 탭 2. 💡 [동적 연산] 전체 Pod 수량의 상위 30% 추출 레이어 ──
    total_pods_count = len(df_pod)
    top_30p_row_limit = max(1, int(total_pods_count * 0.30)) # 최소 1행 방어선
    
    df_waste_top30p = df_pod.sort_values(by="cpu_waste_core_hours", ascending=False).head(top_30p_row_limit)[
        ["date", "cluster", "namespace", "workload_type", "pod", "container", "minutes_running", "cpu_request_max", "cpu_usage_p95", "cpu_waste_core_hours", "mem_waste_gb_hours"]
    ].reset_index(drop=True)

    # ── 📝 탭 3. 💥 [부활] 자원 부족 및 OOMKilled 장애 발생군 명세 ──
    df_shortage = df_pod[(df_pod["cpu_shortage_cores"] > 0) | (df_pod["is_oom_killed"] == True)].sort_values(
        by=["is_oom_killed", "cpu_shortage_cores"], ascending=[False, False]
    )[
        ["date", "cluster", "namespace", "workload_type", "pod", "container", "status", "cpu_request_max", "cpu_usage_p95", "mem_limit_max", "mem_usage_p95"]
    ].reset_index(drop=True)

    # ── 📝 탭 4. 🚨 [신규 구상] Request / Limit 규격 미설정 위반군 전수 리스트업 ──
    df_violators = df_pod[df_pod["has_no_request"] | df_pod["has_no_limit"]].sort_values(by="minutes_running", ascending=False)[
        ["date", "cluster", "namespace", "workload_type", "pod", "container", "has_no_request", "has_no_limit", "cpu_request_max", "cpu_limit_max", "mem_request_max", "mem_limit_max"]
    ].reset_index(drop=True)

    # 엑셀 시트 디스크 세이브 파이프라인 집행
    excel_path = OUT_DIR / "workload_analysis_master_fixed.xlsx"
    with pd.ExcelWriter(excel_path, engine="openpyxl") as writer:
        df_summary.to_excel(writer, sheet_name="0. 전사종합요약", index=False)
        df_ns.to_excel(writer, sheet_name="1. 파레토분석_NS", index=False)
        df_waste_top30p.to_excel(writer, sheet_name="2. 컨테이너낭비Top30P", index=False)
        df_shortage.to_excel(writer, sheet_name="3. 자원부족및OOM장애군", index=False)
        df_violators.to_excel(writer, sheet_name="4. 리소스미설정위반군", index=False)

    print(f"🚀 [마스터 엑셀 완공] 전체 Pod의 30%인 {top_30p_row_limit}개 행을 포함하여 파일 생성을 완료했습니다.")
    print(f" -> 경로: {excel_path}")

if __name__ == "__main__":
    build_master_excel_report()

0개의 댓글