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")
).reset_index().fillna(0)
for col in ["mem_request_max", "mem_limit_max", "mem_usage_p95"]:
df_pod[col] = df_pod[col] / (1024**3)
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")
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)
total_pods_count = len(df_pod)
top_30p_row_limit = max(1, int(total_pods_count * 0.30))
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)
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)
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()