"""
[신규 1단계] Prometheus/Thanos API 직접 조회를 통한 FinOps 원천 데이터 분할 징수 엔진
실행: python step1_prom_fetch.py --days 3 --step 1m --namespace user-analytics
"""
import os
import argparse
import requests
import pandas as pd
from datetime import datetime, timedelta
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, 1h)")
parser.add_argument("--cluster", type=str, default=None, help="특정 클러스터 필터 키워드")
parser.add_argument("--namespace", type=str, default=None, help="특정 네임스페이스 필터 (Regex 지원, 예: 'analytics-.*')")
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):
"""Prometheus query_range API 호출 및 예외 처리"""
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:
print(f" ❌ API 에러 [{response.status_code}]: {response.text[:200]}")
return None
res_json = response.json()
if res_json.get("status") != "success":
return None
return res_json["data"]["result"]
except requests.exceptions.Timeout:
print(" 💥 [타임아웃 발령] 1분 주기 조회의 부하가 너무 큽니다. 해상도(step) 조절이 필요합니다.")
return None
except Exception as e:
print(f" ❌ 네트워크 연결 실패: {str(e)}")
return None
def process_matrix_to_df(matrix_result, metric_name):
"""대용량 매트릭스 JSON 결과를 고속 메모리 맵프 형태의 데이터프레임으로 플래시 플래트닝"""
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", []):
ts = int(val_pair[0])
val = float(val_pair[1])
records.append({
"timestamp": pd.to_datetime(ts, unit='s'),
"cluster": cluster,
"namespace": namespace,
"pod": pod,
"container": container,
"node": node,
metric_name: val
})
return pd.DataFrame(records)
def main():
args = parse_arguments()
if args.days:
end_dt = datetime.utcnow()
start_dt = end_dt - timedelta(days=args.days)
elif args.start_date and args.end_date:
start_dt = datetime.strptime(args.start_date, "%Y-%m-%d")
end_dt = datetime.strptime(args.end_date, "%Y-%m-%d") + timedelta(days=1) - timedelta(seconds=1)
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_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_usage": f"container_memory_working_set_bytes{selector}"
}
print(f"🚀 FinOps 수집 기동 -> URL: {args.url} | 해상도 주기: {args.step}")
print(f"📅 수집 타임라인 범위: {start_dt.strftime('%Y-%m-%d %H:%M:%S')} ~ {end_dt.strftime('%Y-%m-%d %H:%M:%S')}")
print(f"🛡️ 인프라 주입 필터: {selector}\n")
chunk_delta = timedelta(hours=6)
current_start = start_dt
while current_start < end_dt:
current_end = min(current_start + chunk_delta, end_dt)
start_ts = int(current_start.timestamp())
end_ts = int(current_end.timestamp())
chunk_str = current_start.strftime("%Y%m%d_%H")
print(f"⏳ [청크 수집 중] {current_start.strftime('%m-%d %H:%M')} -> {current_end.strftime('%m-%d %H:%M')}")
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")
out_file = RAW_DATA_DIR / f"prom_raw_{chunk_str}.parquet"
merged_chunk.to_parquet(out_file, index=False)
print(f" -> 저장 완료: {out_file} (행수: {len(merged_chunk):,}개)")
else:
print(" -> [경고] 해당 구간에 수집된 메트릭이 없습니다.")
current_start = current_end
print("\n✅ 모든 기간의 원천 데이터 징수 프로세스가 성공적으로 완료되었습니다.")
if __name__ == "__main__":
main()