"""
[신규 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()
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"
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)
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("📊 팟/컨테이너 유닛 레벨 시계열 데이터 롤업 집계 중...")
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()
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] 생성 완료.")
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()