step1 수집 스크립트에서 하루 전체(24시간) 수집 대신 특정 시간대(Range)만 타겟팅하도록 파라미터 제어 로직을 추가한 소스 코드입니다.
프로메테우스/타노스 API 호출 시 필요한 Unix Epoch Timestamp 변환과 수집된 데이터를 Polars로 가공할 때의 시간 필터링까지 모두 처리할 수 있도록 argparse 인터페이스를 유연하게 구성했습니다.
step1_collector.py 소스 코드import argparse
from datetime import datetime, timezone
import polars as pl
import requests # 프로메테우스/타노스 API 호출 시 사용
def parse_args():
parser = argparse.ArgumentParser(description="Step1: Targeted Time-Range Metrics Collector")
# 1. 대상 클러스터 지정
parser.add_argument("--cluster", type=str, required=True, choices=["COMPUTE", "STORAGE"], help="Target Cluster")
# 2. 유연한 시간 범위 인자 (ISO 8601 포맷: YYYY-MM-DD HH:MM:SS)
parser.add_argument("--start-time", type=str, required=True, help="Start time (e.g. '2026-07-20 09:00:00')")
parser.add_argument("--end-time", type=str, required=True, help="End time (e.g. '2026-07-20 18:00:00')")
# 3. (선택) 출력 파일 경로 지정
parser.add_argument("--output", type=str, default="step1_raw_metrics.parquet", help="Output Parquet path")
return parser.parse_args()
def convert_to_epoch(dt_str: str) -> float:
"""문자열 형태의 날짜시간을 Unix Timestamp(Epoch초)로 변환"""
dt = datetime.strptime(dt_str, "%Y-%m-%d %H:%M:%S")
# 타임존 설정이 필요할 경우 local 또는 UTC 적용
return dt.timestamp()
def fetch_prometheus_metrics(thanos_url: str, query: str, start_epoch: float, end_epoch: float, step: str = "1m"):
"""Thanos/Prometheus Range Query API 호출"""
params = {
"query": query,
"start": start_epoch,
"end": end_epoch,
"step": step
}
print(f"📡 Requesting Thanos Range Query: {query}")
print(f" ⏱️ Window: {datetime.fromtimestamp(start_epoch)} ~ {datetime.fromtimestamp(end_epoch)}")
# response = requests.get(f"{thanos_url}/api/v1/query_range", params=params)
# response.raise_for_status()
# return response.json()
def filter_polars_by_time_range(df: pl.DataFrame, start_dt: str, end_dt: str) -> pl.DataFrame:
"""수집 완료된 Polars DataFrame에서 지정 시간대로 슬라이싱 (파일/DB 로딩 시)"""
start_ts = datetime.strptime(start_dt, "%Y-%m-%d %H:%M:%S")
end_ts = datetime.strptime(end_dt, "%Y-%m-%d %H:%M:%S")
return df.filter(
pl.col("timestamp").is_between(start_ts, end_ts)
)
def main():
args = parse_args()
# Epoch Timestamp 변환
start_epoch = convert_to_epoch(args.start_time)
end_epoch = convert_to_epoch(args.end_time)
if start_epoch >= end_epoch:
raise ValueError("❌ Error: --start-time 은 --end-time 보다 이전이어야 합니다.")
print(f"🚀 [Step1] 데이터 수집 시작 | Cluster: {args.cluster}")
print(f"📅 Target Range: {args.start_time} ~ {args.end_time}")
# -------------------------------------------------------------
# 예시 A: Thanos/Prometheus Range Query API로 직접 특정 시간대만 조회
# -------------------------------------------------------------
# thanos_endpoint = "http://thanos-querier.monitoring:9090"
# metric_data = fetch_prometheus_metrics(
# thanos_url=thanos_endpoint,
# query='node_cpu_seconds_total{mode="idle"}',
# start_epoch=start_epoch,
# end_epoch=end_epoch,
# step="1m"
# )
# -------------------------------------------------------------
# 예시 B: Polars 파이프라인 연동 예시
# -------------------------------------------------------------
# raw_df = pl.read_parquet("raw_telemetry_dump.parquet")
# filtered_df = filter_polars_by_time_range(raw_df, args.start_time, args.end_time)
# filtered_df.write_parquet(args.output)
print(f"✅ [Step1] 수집 및 슬라이싱 완료! 저장 경로: {args.output}")
if __name__ == "__main__":
main()
python step1_collector.py \
--cluster COMPUTE \
--start-time "2026-07-20 09:00:00" \
--end-time "2026-07-20 18:00:00" \
--output "step1_compute_business_hours.parquet"
python step1_collector.py \
--cluster STORAGE \
--start-time "2026-07-20 01:30:00" \
--end-time "2026-07-20 03:30:00" \
--output "step1_storage_incident_window.parquet"
--start-time / --end-time 도입: 단순 날짜 단위(YYYY-MM-DD)가 아니라 초 단위(HH:MM:SS) 범위 지정 가능.convert_to_epoch() 변환기: 프로메테우스 API 규격인 Unix Timestamp를 자동으로 파싱하여 Thanos/Prometheus Range Query에 바로 대입 가능.is_between() 슬라이싱 지원: 파일 통으로 읽을 때도 메모리 낭비 없이 해당 시간대 행(Row)만 빠르게 필터링.