26Y19x2

요청하신 대로 시분초 단위를 배제하고 시간(Hour) 단위 범위를 직관적으로 지정할 수 있도록 파싱 및 인자 처리 로직을 수정했습니다.

이제 2026-07-20 11처럼 날짜와 시간(Hour)만 적거나, --date--start-hour, --end-hour 옵션을 조합하여 11:00:00 \le tt < 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()

💻 실행 명령어 예시

1) 11시 딱 1시간만 수집할 때 (11:00:00 \le tt < 12:00:00)

python step1_prom_fetch.py --date "2026-07-20" --start-hour 11

(※ --end-hour를 생략하면 자동으로 시작시간 + 1인 12시로 계산됩니다.)

2) 특정 시간 범위 지정 수집 (예: 11시부터 15시 전까지)

python step1_prom_fetch.py --date "2026-07-20" --start-hour 11 --end-hour 15

3) 문자열로 직접 시간 지정 수집

python step1_prom_fetch.py --start-time "2026-07-20 11" --end-time "2026-07-20 12"

0개의 댓글