"""
[신규 1단계 - 설정 이식 및 아규먼트 강화판] Prometheus/Thanos 마스터 증분 징수 엔진
실행 예시:
python step1_prom_fetch.py --days 1 --step 5m --namespace "user-analytics-.*"
"""
import os
import argparse
import requests
import pandas as pd
from datetime import datetime, timedelta, timezone
from pathlib import Path
from config import RAW_DIR, DEFAULT_THANOS_URL, get_finops_promql_queries
def parse_arguments():
"""인프라 리드가 배치를 자유자재로 튜닝할 수 있는 CLI 파라미터 제어 센터"""
parser = argparse.ArgumentParser(description="FinOps Thanos Raw Time-Series Data Fetcher")
parser.add_argument("--url", type=str, default=DEFAULT_THANOS_URL, help="Thanos/Prometheus Endpoint URL")
parser.add_argument("--days", type=int, default=None, help="최근 며칠간의 데이터를 당일 00:00 정각 정렬 기준으로 수집할지 정의 (정수)")
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, 1h)")
parser.add_argument("--cluster", type=str, default=None, help="특정 타겟 클러스터 라벨 필터 키워드")
parser.add_argument("--namespace", type=str, default=None, help="특정 네임스페이스 필터 (정규식/Regex 지원, 예: 'analytics-.*')")
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 any(char in namespace for char in ["|", ".*", "^", "$"]):
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)
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) 중 하나는 Argument 필수 인자입니다.")
return
selector = build_promql_selectors(args.cluster, args.namespace)
queries = get_finops_promql_queries(selector)
print(f"🚀 FinOps 데이터 징수 기동 -> 엔드포인트 URL: {args.url}")
print(f"📊 수집 해상도 주기: {args.step} | 인프라 식별 셀렉터: {selector}")
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")
is_completed_chunk = (current_start + chunk_delta) <= now_utc
if is_completed_chunk:
out_file = RAW_DIR / f"prom_raw_{chunk_str}.parquet"
else:
out_file = RAW_DIR / f"prom_raw_{chunk_str}_active.parquet"
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):,}개)")
if is_completed_chunk:
active_file = RAW_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()