26J29k4

Young-Kyoo Kim·2026년 6월 29일
"""
[1단계 - 정각 정렬 및 액티브 분리형] Prometheus/Thanos FinOps 원천 데이터 무결성 증분 수집 엔진
실행: python step1_prom_fetch.py --days 1 --step 1m
"""

import os
import argparse
import requests
import pandas as pd
from datetime import datetime, timedelta, timezone
from pathlib import Path

THANOS_QUERY_URL = os.getenv("THANOS_QUERY_URL", "http://thanos-query.internal.zone:9090")
RAW_DATA_DIR = Path("./data/raw")
RAW_DATA_DIR.mkdir(parents=True, exist_ok=True)

def parse_arguments():
    parser = argparse.ArgumentParser(description="FinOps Prometheus Raw Data Fetcher")
    parser.add_argument("--url", type=str, default=THANOS_QUERY_URL, help="Thanos/Prometheus Endpoint URL")
    parser.add_argument("--days", type=int, default=None, help="최근 며칠간의 데이터를 수집할지 정의 (정수)")
    parser.add_argument("--start-date", type=str, default=None, help="수집 시작 날짜 (YYYY-MM-DD)")
    parser.add_argument("--end-date", type=str, default=None, help="수집 종료 날짜 (YYYY-MM-DD)")
    parser.add_argument("--step", type=str, default="1m", help="수집 해상도 주기 (예: 1m, 5m, 15m)")
    parser.add_argument("--cluster", type=str, default=None, help="특정 클러스터 필터 키워드")
    parser.add_argument("--namespace", type=str, default=None, help="특정 네임스페이스 필터 (Regex 지원)")
    parser.add_argument("--force", action="store_true", help="기존 파일을 무시하고 전수 강제 재수집")
    return parser.parse_args()

def build_promql_selectors(cluster, namespace):
    selectors = ['container!=""']
    if cluster: selectors.append(f'cluster="{cluster}"')
    if namespace:
        if "|" in namespace or ".*" in namespace or "^" in namespace:
            selectors.append(f'namespace=~"{namespace}"')
        else:
            selectors.append(f'namespace="{namespace}"')
    return "{" + ", ".join(selectors) + "}"

def fetch_range_chunk(url, query, start_ts, end_ts, step):
    endpoint = f"{url}/api/v1/query_range"
    params = {"query": query, "start": start_ts, "end": end_ts, "step": step}
    try:
        response = requests.get(endpoint, params=params, timeout=120)
        if response.status_code != 200: return None
        res_json = response.json()
        if res_json.get("status") != "success": return None
        return res_json["data"]["result"]
    except Exception:
        return None

def process_matrix_to_df(matrix_result, metric_name):
    if not matrix_result: return pd.DataFrame()
    records = []
    for item in matrix_result:
        metric_labels = item.get("metric", {})
        cluster = metric_labels.get("cluster", "default")
        namespace = metric_labels.get("namespace", "unknown")
        pod = metric_labels.get("pod", "unknown")
        container = metric_labels.get("container", "unknown")
        node = metric_labels.get("node", metric_labels.get("instance", "unknown")).split(":")[0]
        
        for val_pair in item.get("values", []):
            records.append({
                "timestamp": pd.to_datetime(int(val_pair[0]), unit='s'),
                "cluster": cluster, "namespace": namespace, "pod": pod, "container": container, "node": node,
                metric_name: float(val_pair[1])
            })
    return pd.DataFrame(records)

def main():
    args = parse_arguments()
    
    now_utc = datetime.now(timezone.utc)
    
    # ── 🛡️ [패치 1] 실행 시각에 구애받지 않도록 00시 00분 정각 정렬 (0분 시작) ──
    if args.days:
        start_dt = (now_utc - timedelta(days=args.days)).replace(hour=0, minute=0, second=0, microsecond=0)
        end_dt = now_utc
    elif args.start_date and args.end_date:
        start_dt = datetime.strptime(args.start_date, "%Y-%m-%d").replace(tzinfo=timezone.utc)
        end_dt = (datetime.strptime(args.end_date, "%Y-%m-%d") + timedelta(days=1)).replace(tzinfo=timezone.utc)
    else:
        print("❌ 에러: (--days) 또는 (--start-date 및 --end-date) 중 하나는 필수입니다.")
        return

    selector = build_promql_selectors(args.cluster, args.namespace)
    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!=')}"
    }

    print(f"🚀 FinOps 무결성 증분 수집 시작 -> URL: {args.url} | 해상도: {args.step}")
    print(f"📅 정렬된 수집 범위 (UTC): {start_dt.strftime('%Y-%m-%d %H:%M:%S')} ~ {end_dt.strftime('%Y-%m-%d %H:%M:%S')}\n")

    chunk_delta = timedelta(hours=6)
    current_start = start_dt

    while current_start < end_dt:
        current_end = min(current_start + chunk_delta, end_dt)
        chunk_str = current_start.strftime("%Y%m%d_%H")
        
        # ── 🛡️ [패치 2] 마감된 청크 vs 진행 중인 실시간 청크 분리 매커니즘 ──
        # 현재 루프의 6시간 범위가 완전히 과거 시점인 경우에만 '마감 완료'로 판단
        is_completed_chunk = (current_start + chunk_delta) <= now_utc
        
        if is_completed_chunk:
            out_file = RAW_DATA_DIR / f"prom_raw_{chunk_str}.parquet"
        else:
            # 아직 시간이 채워지고 있는 구간은 active 접미사를 붙여 캐싱 동결을 방지
            out_file = RAW_DATA_DIR / f"prom_raw_{chunk_str}_active.parquet"
        
        # 💡 마감 완료된 파일이 디스크에 존재할 때만 안전하게 SKIP 작동
        if not args.force and is_completed_chunk and out_file.exists() and out_file.stat().st_size > 0:
            print(f"⏭️  [SKIP] {current_start.strftime('%m-%d %H:%M')} 구간은 마감 완료되어 패스합니다.")
            current_start = current_end
            continue

        start_ts = int(current_start.timestamp())
        end_ts = int(current_end.timestamp())
        
        if is_completed_chunk:
            print(f"⏳ [청크 마감 수집] {current_start.strftime('%m-%d %H:%M')} ~ {current_end.strftime('%m-%d %H:%M')} UTC")
        else:
            print(f"🔄 [실시간 active 수집] {current_start.strftime('%m-%d %H:%M')} ~ {current_end.strftime('%m-%d %H:%M')} UTC (매번 덮어쓰기)")
            
        chunk_df_list = []
        for metric_name, query_expr in queries.items():
            raw_matrix = fetch_range_chunk(args.url, query_expr, start_ts, end_ts, args.step)
            df_metric = process_matrix_to_df(raw_matrix, metric_name)
            if not df_metric.empty:
                chunk_df_list.append(df_metric)
        
        if chunk_df_list:
            merged_chunk = chunk_df_list[0]
            for next_df in chunk_df_list[1:]:
                join_keys = ["timestamp", "cluster", "namespace", "pod", "container", "node"]
                merged_chunk = pd.merge(merged_chunk, next_df, on=join_keys, how="outer")
            
            merged_chunk.to_parquet(out_file, index=False)
            print(f"  -> 💾 저장 완료: {out_file.name} (행수: {len(merged_chunk):,}개)")
            
            # ── 🛡️ [패치 3] 중복 연산 방지 청소 레이어 ──
            # 실시간 active 상태였다가 시간이 흘러 마감 파일(.parquet)이 정상 생성되면, 예전의 임시 active 파일은 삭제
            if is_completed_chunk:
                active_file = RAW_DATA_DIR / f"prom_raw_{chunk_str}_active.parquet"
                if active_file.exists():
                    active_file.unlink()
                    print(f"  -> 🧹 이전 임시 active 파일을 안전하게 정리했습니다.")
        else:
            print("  -> ⚠️ 해당 구간 메트릭 없음.")

        current_start = current_end

    print("\n✅ 정각 정렬 및 데이터 정합성 보정이 완료되었습니다.")

if __name__ == "__main__":
    main()

0개의 댓글