26J29f5

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 derive_cluster_from_node(node_name):
    """
    📌 [사내 노드 명명 규칙 기반 클러스터 자동 판별 함수]
    노드 호스트네임의 패턴을 분석하여 정확한 클러스터 명칭을 반환합니다.
    사내 컨벤션(예: icdlhk8s-prod-wk01, cls-stage-node02 등)에 맞게 조건을 수정하세요.
    """
    n = str(node_name).lower()
    
    # 예시 1: 노드 이름에 특정 클러스터 접두사나 키워드가 포함된 경우
    if "icdlh" in n or "prod" in n:
        return "prod-data-lakehouse"
    elif "stage" in n or "stg" in n:
        return "stage-lakehouse"
    elif "dev" in n:
        return "dev-sandbox"
    
    # 예시 2: 도메인 형태로 분리 가능한 경우 (지정 구문이 있을 때)
    # 예: node-01.cluster-a.internal -> cluster-a 추출
    if "." in n:
        parts = n.split(".")
        if len(parts) > 1 and "cluster" in parts[1]:
            return parts[1]

    # 매칭되는 규칙이 없을 경우 기본 백업값
    return "unclassified-onprem-cluster"

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("🌐 노드 물리 명명 규칙(Hostname) 분석 기반 클러스터(Cluster) 소속 식별 중...")
    df_raw["cluster"] = df_raw["node"].apply(derive_cluster_from_node)
    
    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개의 댓글