"""
[신규 2단계 - 최종 통합판] 노드 정규식 필터링, 클러스터 역추적 및 거버넌스 인덱싱 엔진
실행:
(Linux 환경 예시) export NODE_PREFIX_PATTERN=r"name1wk\d+"
python step2_pipeline.py
"""
import os
import re
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 filter_target_nodes(df):
"""
⚙️ [인프라 거버넌스 가드] 환경변수 정규식을 활용한 고속 입구 컷 필터
NODE_PREFIX_PATTERN 환경변수 규칙에 맞는 노드 시계열만 남기고 하방 연산 오버헤드 방어
"""
node_pattern = os.getenv("NODE_PREFIX_PATTERN", ".*")
print(f"🌐 [노드 정규식 필터] 주입된 환경변수 패턴: '{node_pattern}'")
try:
re.compile(node_pattern)
except re.error:
print(f"🚨 [경고] 환경변수 정규식 문법에 오류가 포착되었습니다: '{node_pattern}' -> 필터를 무풍지대(.*)로 폴백합니다.")
node_pattern = ".*"
initial_rows = len(df)
df_filtered = df[df["node"].str.contains(node_pattern, regex=True, na=False, case=False)].reset_index(drop=True)
filtered_rows = len(df_filtered)
dropped_rows = initial_rows - filtered_rows
print(f" -> ✂️ 필터링 결과: {initial_rows:,}행 중 {dropped_rows:,}행 노이즈 제거 -> {filtered_rows:,}행 생존")
return df_filtered
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"
elif "dev" in n:
return "dev-sandbox"
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/ 폴더에 수집된 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 = filter_target_nodes(df_raw)
if df_raw.empty:
print("⚠️ [안내] 필터 조건에 부합하는 노드가 없어 후처리 단계를 조기 종료합니다.")
return
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("📊 팟/컨테이너 실명 기반 타임시리즈 롤업 집계 가동 (1분 주기 카운트 역산)...")
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)
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["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)
print("💾 1차 필터링 무결성 마스터 원부 [enriched_fixed_7d.parquet] 빌드 완료.")
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__":
print("=== [2단계 후처리] 데이터 병합 정제 및 거버넌스 분석 파이프라인 가동 ===")
run_enrich_and_pareto_pipeline()