요청하신 대로 시분초 단위를 배제하고 시간(Hour) 단위 범위를 직관적으로 지정할 수 있도록 파싱 및 인자 처리 로직을 수정했습니다.
이제 2026-07-20 11처럼 날짜와 시간(Hour)만 적거나, --date와 --start-hour, --end-hour 옵션을 조합하여 11:00:00 < 12:00:00 범위를 손쉽게 제어할 수 있습니다.
step1_prom_fetch.py 소스 코드"""
step1_prom_fetch.py — 시간(Hour) 단위 분할 징수 및 AIStor 업로드 후 로컬 캐시 보존 엔진
"""
import os
import sys
import time
import argparse
import requests
import boto3
import pandas as pd
from datetime import datetime, timedelta, timezone
import config
from config import RAW_DIR, get_finops_promql_queries
KST = timezone(timedelta(hours=9))
def parse_arguments():
parser = argparse.ArgumentParser(description="Thanos Resource Fetcher for Data Lakehouse")
# 1. 일자 지정 옵션
parser.add_argument("--date", type=str, default=None, help="대상 날짜 (YYYY-MM-DD)")
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("--days", type=int, default=None, help="수집할 과거 일수")
# 2. 시간(Hour) 단위 지정 옵션 (0~23)
parser.add_argument("--start-hour", type=int, default=None, help="시작 시간 (0~23 시)")
parser.add_argument("--end-hour", type=int, default=None, help="종료 시간 (1~24 시)")
# 3. 통합 일시(Hour) 지정 옵션 ("YYYY-MM-DD HH")
parser.add_argument("--start-time", type=str, default=None, help="시작 일시 (예: '2026-07-20 11')")
parser.add_argument("--end-time", type=str, default=None, help="종료 일시 (예: '2026-07-20 12')")
parser.add_argument("--step", type=str, default="1m", help="PromQL 조회 간격")
return parser.parse_args()
def parse_hour_datetime(dt_str, is_end=False):
"""'YYYY-MM-DD HH' 또는 'YYYY-MM-DD' 포맷을 KST 시간 단위 datetime으로 정제"""
if not dt_str:
return None
dt_str = dt_str.strip()
# 1. 'YYYY-MM-DD HH' 형태인 경우
try:
dt = datetime.strptime(dt_str, "%Y-%m-%d %H")
return dt.replace(minute=0, second=0, microsecond=0, tzinfo=KST)
except ValueError:
pass
# 2. 'YYYY-MM-DD' 형태인 경우
try:
dt = datetime.strptime(dt_str, "%Y-%m-%d")
if is_end:
dt += timedelta(days=1) # 종료일의 경우 다음날 00:00:00까지 포함
return dt.replace(tzinfo=KST)
except ValueError:
pass
raise ValueError(f"❌ 날짜/시간 포맷 오류: '{dt_str}'. 'YYYY-MM-DD 11' 또는 'YYYY-MM-DD' 형식으로 입력해주세요.")
def query_thanos_matrix(thanos_url, query, start_time, end_time, step="1m", metric_name="metric"):
"""Thanos 타임아웃 방지를 위한 15분 분할 징수 엔진"""
is_epoch = isinstance(start_time, (int, float))
if is_epoch:
chunk_delta = 15 * 60
current_start = int(start_time)
end_time = int(end_time)
else:
chunk_delta = timedelta(minutes=15)
current_start = start_time
headers = {
"Accept": "application/json",
"Thanos-Allow-Partial-Response": "false"
}
all_chunks_data = []
log_start = datetime.fromtimestamp(start_time).strftime('%H:%M') if is_epoch else current_start.strftime('%H:%M')
log_end = datetime.fromtimestamp(end_time).strftime('%H:%M') if is_epoch else end_time.strftime('%H:%M')
print(f"🕵️♂️ [Thanos 분할 수집] 🎯 대상: {metric_name} | {log_start} ~ {log_end} (15분 청크)")
while current_start < end_time:
current_end = min(current_start + chunk_delta, end_time)
start_param = current_start if is_epoch else current_start.strftime("%Y-%m-%dT%H:%M:%SZ")
end_param = current_end if is_epoch else current_end.strftime("%Y-%m-%dT%H:%M:%SZ")
params = {
"query": query,
"start": start_param,
"end": end_param,
"step": step
}
for retry in range(1, 4):
try:
response = requests.get(
f"{thanos_url}/api/v1/query_range",
params=params,
headers=headers,
timeout=300
)
if response.status_code == 200:
res_json = response.json()
chunk_result = res_json.get("data", {}).get("result", [])
all_chunks_data.extend(chunk_result)
break
else:
time.sleep(2)
except requests.exceptions.ReadTimeout:
print(f" 🚨 [Thanos 읽기 타임아웃] {retry}회차 실패. 5초 후 리트라이...")
time.sleep(5)
except Exception:
time.sleep(3)
else:
print(f" 💥 [최종 수집 실패] {current_start} ~ {current_end} 구간 유실.")
current_start = current_end
return all_chunks_data
def main():
args = parse_arguments()
thanos_url = os.getenv("THANOS_QUERY_URL", "http://thanos-query.internal.zone:9090")
bucket_name = os.getenv("MINIO_RAW_BUCKET", "devops-test")
endpoint = os.getenv("MINIO_ENDPOINT", "http://minio-service.internal.zone:9000")
access_key = os.getenv("MINIO_ACCESS_KEY", "admin")
secret_key = os.getenv("MINIO_SECRET_KEY", "password")
now_kst = datetime.now(KST)
# 💡 [시간 범위 확정 로직]
# 방식 1) --start-time "YYYY-MM-DD 11" & --end-time "YYYY-MM-DD 12"
if args.start_time and args.end_time:
start_dt = parse_hour_datetime(args.start_time)
end_dt = parse_hour_datetime(args.end_time, is_end=True)
# 방식 2) --date "YYYY-MM-DD" + --start-hour 11 + --end-hour 12
elif (args.date or args.start_date) and (args.start_hour is not None):
target_date_str = args.date or args.start_date
base_dt = datetime.strptime(target_date_str, "%Y-%m-%d").replace(tzinfo=KST)
start_dt = base_dt.replace(hour=args.start_hour, minute=0, second=0)
end_h = args.end_hour if args.end_hour is not None else args.start_hour + 1
if end_h == 24:
end_dt = base_dt.replace(hour=0, minute=0, second=0) + timedelta(days=1)
else:
end_dt = base_dt.replace(hour=end_h, minute=0, second=0)
# 방식 3) --start-date & --end-date (일 단위 전체)
elif args.start_date and args.end_date:
start_dt = parse_hour_datetime(args.start_date)
end_dt = parse_hour_datetime(args.end_date, is_end=True)
# 방식 4) --days (상대적 과거 일수)
elif args.days:
start_dt = (now_kst - timedelta(days=args.days)).replace(hour=0, minute=0, second=0, microsecond=0)
end_dt = now_kst
# 기본값: 어제 하루 전체 (00:00:00 ~ 24:00:00)
else:
yesterday = (now_kst - timedelta(days=1)).date()
start_dt = datetime.combine(yesterday, datetime.min.time()).replace(tzinfo=KST)
end_dt = datetime.combine(now_kst.date(), datetime.min.time()).replace(tzinfo=KST)
print(f"🚀 [Step1 개시] Thanos 수집 구간: {start_dt.strftime('%Y-%m-%d %H:%M:%S')} <= t < {end_dt.strftime('%Y-%m-%d %H:%M:%S')}")
base_selector = '{container!="", namespace!~"kube-system|istio-system|monitoring"}'
queries_matrix = get_finops_promql_queries(base_selector)
s3_client = boto3.client(
"s3", endpoint_url=endpoint, aws_access_key_id=access_key, aws_secret_access_key=secret_key,
region_name="us-east-1", config=boto3.session.Config(signature_version="s3v4")
)
chunk_delta = timedelta(hours=1)
current_chunk_start = start_dt
while current_chunk_start < end_dt:
current_chunk_end = min(current_chunk_start + chunk_delta, end_dt)
chunk_str = current_chunk_start.strftime("%Y%m%d_%H")
start_ts = int(current_chunk_start.timestamp())
end_ts = int(current_chunk_end.timestamp())
print(f"⏳ 파티션 징수 및 AIStor 업로드 중... [{chunk_str}시]")
records_master = []
for metric_name, query_str in queries_matrix.items():
result_set = query_thanos_matrix(thanos_url, query_str, start_ts, end_ts, args.step, metric_name=metric_name)
if not result_set: continue
for item in result_set:
metric_labels = item.get("metric", {})
values_list = item.get("values", [])
cluster_label = metric_labels.get("cluster", "prod-cluster")
namespace = metric_labels.get("namespace", "unknown")
pod = metric_labels.get("pod", "unknown")
container = metric_labels.get("container", "all-volume")
node = metric_labels.get("node", metric_labels.get("instance", "unknown")).split(":")[0]
for val_pair in values_list:
records_master.append({
"timestamp": datetime.fromtimestamp(float(val_pair[0]), tz=timezone.utc).astimezone(KST).replace(tzinfo=None),
"cluster": cluster_label, "namespace": namespace, "pod": pod, "container": container, "node": node,
"metric_type": metric_name, "value": float(val_pair[1])
})
if records_master:
df_chunk = pd.DataFrame(records_master)
df_clean = df_chunk.drop_duplicates(subset=["timestamp", "namespace", "pod", "container", "metric_type"])
df_pivot = df_clean.groupby(["timestamp", "cluster", "namespace", "pod", "container", "node", "metric_type"])["value"].max().unstack().reset_index()
for expected_col in queries_matrix.keys():
if expected_col not in df_pivot.columns: df_pivot[expected_col] = 0.0
out_file = RAW_DIR / f"prom_raw_{chunk_str}.parquet"
df_pivot.to_parquet(str(out_file), index=False)
object_key = f"raw/prom_raw_{chunk_str}.parquet"
s3_client.upload_file(str(out_file), bucket_name, object_key)
print(f" ✅ AIStor 백업 완료 ➡️ 로컬 캐시 지속 보존: {out_file.name}")
else:
print(f" ⚠️ 청크 [{chunk_str}시] 데이터 공백 스킵.")
current_chunk_start = current_chunk_end
print("\n🏁 === [Step1 전송 성료] 로컬 작업공간에 원천 Parquet 풀이 보존되었습니다. ===")
if __name__ == "__main__":
main()
11:00:00 < 12:00:00)python step1_prom_fetch.py --date "2026-07-20" --start-hour 11
(※ --end-hour를 생략하면 자동으로 시작시간 + 1인 12시로 계산됩니다.)
python step1_prom_fetch.py --date "2026-07-20" --start-hour 11 --end-hour 15
python step1_prom_fetch.py --start-time "2026-07-20 11" --end-time "2026-07-20 12"