26J29k5

Young-Kyoo Kim·2026년 6월 29일
"""
[신규 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)
    
    # na=False로 노드명 결측치 크래시 방어, case=False로 호스트네임 대소문자 무시 안전성 확보
    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():
    # .parquet 패턴으로 완결판(.parquet) 및 실시간 가동판(_active.parquet) 일괄 흡수
    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)
    
    # ── 🛡️ [추가 반영 1] 최초 병합 직후 메모리 세이빙을 위한 노드 정규식 입구 컷 필터 집행 ──
    df_raw = filter_target_nodes(df_raw)
    
    if df_raw.empty:
        print("⚠️ [안내] 필터 조건에 부합하는 노드가 없어 후처리 단계를 조기 종료합니다.")
        return

    # 날짜 문자열 정규화 (YYYY-MM-DD)
    df_raw["date"] = df_raw["timestamp"].dt.strftime("%Y-%m-%d")
    
    # ── 🛡️ [추가 반영 2] 타노스 환경 대응용 노드명 기반 클러스터 소속 동적 재할당 ──
    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분 주기 카운트 역산)...")
    # 1분 단위 시계열의 크기(size)가 곧 실제 연속 기동 분 수가 됨 (0분 뭉개짐 버그 해결)
    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)

    # 바이트 단위를 기가바이트(GB) 규격으로 변환
    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)

    # ── 🛡️ [추가 반영 3] 6번 탭 및 4번 거버넌스 위반군 탭을 위한 스키마 플래그 주입 ──
    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] 빌드 완료.")

    # ── 📈 2차 거버넌스 타겟팅 전용 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__":
    print("=== [2단계 후처리] 데이터 병합 정제 및 거버넌스 분석 파이프라인 가동 ===")
    run_enrich_and_pareto_pipeline()

0개의 댓글