기존에 빌드해 두신 2단계 배치 파이프라인 구조가 아주 견고하게 잡혀 있네요. 1단계(Step 1)는 Thanos의 타임아웃을 박멸하면서 1시간 단위로 원천 데이터를 안전하게 확보(Extraction)하는 데 집중하고, 2단계(Step 2)는 Polars의 스트리밍 엔진을 활용해 메모리(256MB) 한계 내에서 일별 정산 및 지표 고도화(Aggregation & Transformation)를 수행하는 구조입니다.
이 파이프라인에서 네트워크, 디스크 I/O, 인증 등의 추가 메트릭을 수집하고 융합할 때, 각 단계의 책임(Separation of Concerns)을 명확히 나누어 인젝션하는 것이 가장 체계적입니다.
기존 스크립트 기반으로 어느 단계에서 무엇을 처리해야 하는지 가이드와 수정 코드를 정리했습니다.
데이터의 수집과 고차원 분석을 가장 효율적으로 안착시키기 위한 두 단계의 역할 분담 규칙입니다.
[Thanos Matrix]
│
▼ (1분 해상도 수집)
┌────────────────────────────────────────────────────────┐
│ Step 1: prom_fetch.py (원천 데이터 자산화 레이어) │
│ - 11종 코어 + 7종 확장 메트릭 PromQL 동시 수집 │
│ - 다른 그라뉴얼리티(Pod vs Container) 그대로 unstack │
└────────────────────────────────────────────────────────┘
│
▼ (1시간 단위 Wide Parquet 로컬 캐시/AIStor 저장)
┌────────────────────────────────────────────────────────┐
│ Step 2: polars_reprocess.py (SRE 인텔리전스 레이어) │
│ - 1분 단위 로우 레벨에서 SRE 유도 특징량(Feature) 계산 │
│ - 일일 배치 스케일로 P95/Max 통계 집계 및 가버넌스 마킹 │
└────────────────────────────────────────────────────────┘
step1_prom_fetch.py)에서의 역할: 수집 및 1차 얼라인먼트config.py의 쿼리 매트릭스에 신규 확장 메트릭(네트워크, 디스크, 인증) PromQL을 추가하고, Thanos에서 그대로 긁어와 하나의 Wide 포맷 Parquet 파일로 unstack 링버퍼를 만드는 것까지 전담합니다.net_*, pv_*)는 수집 시 container 레이블이 없으므로, 기존 코드의 metric_labels.get("container", "all-volume") 로직에 의해 자연스럽게 "all-volume"이라는 가상 컨테이너 행으로 녹아들게 됩니다. 이 구조를 유지해야 데이터 결손 없이 깨끗하게 조인됩니다.💡 코드 가드레일 발견: 기존 Step 1 소스 코드 중
df_clean = df_chunk.unique(subset=["timestamp", "namespace", "pod", "container"])파트는 판다스 문법 오류(drop_duplicates가 맞음)일 뿐만 아니라, 피벗하기 전에 이 지정을 해버리면 metric_type 스코프가 무시되어 한 가지만 남고 데이터가 유실됩니다.groupby().max().unstack()구조가 중복을 알아서 제거하므로 이 라인은 과감히 생킵니다.
step2_polars_reprocess.py)에서의 역할: SRE 유도 특징량 생성 및 배치 집계collect(streaming=True) 직후)에 수학적 SRE 유도 특징량(, , 대역폭 비율 등) 컬럼을 대량 생성합니다..group_by().agg() 블록에서 롤업 집계하여 사내 LLM이 시계열의 복잡한 연관성을 역추적할 수 있도록 컨텍스트 자산화합니다.config.py 및 step1_prom_fetch.py 확장먼저 config.py 내부의 쿼리 맵에 확장 메트릭 7종을 추가 정의합니다.
# config.py 내부 혹은 스크립트 상단 쿼리 레포지토리 확장 예시
def get_finops_promql_queries(base_selector):
return {
# [기존 코어 11종]
"cpu_request": f"sum(kube_pod_container_resource_requests{{resource='cpu'}}) by (cluster, namespace, node, pod, container)",
"cpu_limit": f"sum(kube_pod_container_resource_limits{{resource='cpu'}}) by (cluster, namespace, node, pod, container)",
"cpu_usage": f"sum(rate(container_cpu_usage_seconds_total{base_selector}[2m])) by (cluster, namespace, node, pod, container) * 100",
"cpu_throttled": f"sum(rate(container_cpu_cfs_throttled_seconds_total{base_selector}[2m])) by (cluster, namespace, node, pod, container)",
"mem_request": f"sum(kube_pod_container_resource_requests{{resource='memory'}}) by (cluster, namespace, node, pod, container)",
"mem_limit": f"sum(kube_pod_container_resource_limits{{resource='memory'}}) by (cluster, namespace, node, pod, container)",
"mem_usage": f"sum(container_memory_working_set_bytes{base_selector}) by (cluster, namespace, node, pod, container)",
"mem_rss": f"sum(container_memory_rss{base_selector}) by (cluster, namespace, node, pod, container)",
"oom_event": f"sum(changes(kube_pod_container_status_terminated_reason{{reason='OOMKilled'}}[2m])) by (cluster, namespace, node, pod, container)",
"pv_capacity": f"sum(kubelet_volume_stats_capacity_bytes) by (cluster, namespace, node, pod)",
"pv_used": f"sum(kubelet_volume_stats_used_bytes) by (cluster, namespace, node, pod)",
# [신규 수집 확장 7종: CNI, 디스크 I/O, 인증 밀도]
"net_rx_bytes": f"sum(rate(container_network_receive_bytes_total{base_selector}[2m])) by (cluster, namespace, node, pod)",
"net_tx_bytes": f"sum(rate(container_network_transmit_bytes_total{base_selector}[2m])) by (cluster, namespace, node, pod)",
"net_drop_errors": f"sum(rate(container_network_receive_errors_total{base_selector}[2m]) + rate(container_network_transmit_errors_total{base_selector}[2m])) by (cluster, namespace, node, pod)",
"disk_read_bytes": f"sum(rate(container_fs_reads_bytes_total{base_selector}[2m])) by (cluster, namespace, node, pod, container)",
"disk_write_bytes": f"sum(rate(container_fs_writes_bytes_total{base_selector}[2m])) by (cluster, namespace, node, pod, container)",
"disk_iops": f"sum(rate(container_fs_reads_total{base_selector}[2m]) + rate(container_fs_writes_total{base_selector}[2m])) by (cluster, namespace, node, pod, container)",
"auth_metadata_tps": f"sum(rate(minio_api_requests_incoming_total[2m])) by (cluster, namespace, pod) or sum(rate(http_server_requests_seconds_count{{uri=~'.*token.*'}}[2m])) by (cluster, namespace, pod)"
}
이후 step1_prom_fetch.py 데이터 가공 파트를 아래와 같이 간결화 및 동기화합니다.
# step1_prom_fetch.py의 메인 루프 내부 가공 영역 교체
if records_master:
df_chunk = pd.DataFrame(records_master)
# 🧹 [가드레일 보정] 유실 유발 마스킹 라인 제거 후 안전한 피벗만 수행
df_pivot = df_chunk.groupby(
["timestamp", "cluster", "namespace", "pod", "container", "node", "metric_type"]
)["value"].max().unstack().reset_index()
# 존재하지 않는 컬럼 디폴트 0.0 바인딩 가드
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}")
step2_polars_reprocess.py SRE 연관 특징량 엔진 이식Polars의 지연 연산 트리에 신규 수집 필드 명세를 주입하고, 하루치 고해상도 타임라인 위에서 상호 인프라 간섭 현상을 찾아낼 특징량 산출 코드를 추가합니다.
# step2_polars_reprocess.py 내부의 컬럼 정의 영역 자산 확장
base_fields = ["timestamp", "namespace", "node", "pod", "container"]
metric_fields = [
"cpu_request", "cpu_limit", "cpu_usage", "cpu_throttled",
"mem_request", "mem_limit", "mem_usage", "mem_rss", "oom_event",
"pv_capacity", "pv_used",
# 확장 규격 컬럼 프로젝션 추가
"net_rx_bytes", "net_tx_bytes", "net_drop_errors",
"disk_read_bytes", "disk_write_bytes", "disk_iops", "auth_metadata_tps"
]
이후, for date_chunk in target_dates: 루프 내에서 데이터프레임을 수집(collect)하여 구체화한 바로 다음 영역에 SRE 벡터 특징량 계산 코드와 고급 집계 내역을 인젝션합니다.
# [Step 2 루프 내부 실행부]
df_day = lf_day.collect(streaming=True)
if df_day.is_empty():
del df_day; gc.collect()
continue
# 1. 인프라 정보 및 워크로드 타입 맵핑 선수행
df_day = df_day.with_columns([
pl.col("node").map_elements(
lambda x: {"cluster_type": classify_node_infrastructure(x)[0], "workload_domain": classify_node_infrastructure(x)[1]},
return_dtype=pl.Struct({"cluster_type": pl.String, "workload_domain": pl.String})
).alias("infra_struct"),
pl.col("pod").map_elements(get_workload_type, return_dtype=pl.String).alias("workload_type"),
pl.lit(date_chunk).alias("date")
]).with_columns([
pl.col("infra_struct").struct.field("cluster_type").alias("cluster_type"),
pl.col("infra_struct").struct.field("workload_domain").alias("workload_domain")
]).drop("infra_struct")
# 🚀 2. [신설] 1분 고해상도 로우 레벨 기반 SRE 복합 관계 지표 인라인 연산 가동
df_day = df_day.with_columns([
# Cilium 네트워크-컴퓨트 부하 점유 비율 (Network-to-Compute Ratio)
((pl.col("net_rx_bytes") + pl.col("net_tx_bytes")) / (pl.col("cpu_usage") + 1e-5)).alias("feat_net_to_cpu_ratio"),
# DirectPV NVMe 입출력 비대칭성 (Storage IO Asymmetry)
(pl.col("disk_write_bytes") / (pl.col("disk_read_bytes") + 1e-5)).alias("feat_io_asymmetry"),
# Keycloak 인증 / 메타데이터 트래픽 조밀도 (Auth Call Density)
(pl.col("auth_metadata_tps") / (pl.col("cpu_usage") + 1e-5)).alias("feat_auth_density")
])
# 🚀 3. [신설] 동일 윈도우 시간 분(timestamp) 내 분산 노드 간 부하 변동계수 계산 (Load Skew CV)
# 수식: 특정 분에 배포된 동일 워크로드들의 표준편차 / 평균
df_day = df_day.with_columns([
(pl.col("cpu_usage").std().over(["timestamp", "workload_type"]) /
(pl.col("cpu_usage").mean().over(["timestamp", "workload_type"]) + 1e-5)).fill_null(0.0).alias("feat_load_skew_cv")
])
# 분석 타겟 클러스터 분기 루프 진입
target_clusters = ["COMPUTE", "STORAGE"] if args.cluster.upper() == "ALL" else [args.cluster.upper()]
for cluster_chunk in target_clusters:
df_slice = df_day.filter(pl.col("cluster_type") == cluster_chunk)
df_slice = filter_target_nodes(df_slice, cluster_chunk)
if df_slice.is_empty(): continue
print(f" ⏳ [Polars 분산 집계] {date_chunk} ➡️ {cluster_chunk} 확장형 그룹바이 가동...")
# [테이블 1] 상세 메트릭 원부 집계 파트 확장 수용
df_pod = df_slice.group_by([
"date", "cluster_type", "workload_domain", "namespace", "workload_type", "node", "pod", "container"
]).agg([
pl.len().alias("minutes_running"),
pl.col("cpu_request").max().alias("cpu_request_max"),
pl.col("cpu_limit").max().alias("cpu_limit_max"),
pl.col("cpu_usage").quantile(0.95).alias("cpu_usage_p95"),
pl.col("cpu_usage").max().alias("cpu_usage_max"),
pl.col("cpu_throttled").max().alias("cpu_throttled_max"),
(pl.col("mem_request").max() / GB_DIV).alias("mem_request_max"),
(pl.col("mem_limit").max() / GB_DIV).alias("mem_limit_max"),
(pl.col("mem_usage").quantile(0.95) / GB_DIV).alias("mem_usage_p95"),
(pl.col("mem_rss").quantile(0.95) / GB_DIV).alias("mem_rss_p95"),
pl.col("oom_event").sum().alias("oom_strike_sum"),
(pl.col("pv_capacity").max() / GB_DIV).alias("pv_capacity_max"),
(pl.col("pv_used").quantile(0.95) / GB_DIV).alias("pv_usage_p95"),
# ─── 📊 [신규 확장 집계 필드 연동] ───
pl.col("net_rx_bytes").quantile(0.95).alias("net_rx_p95"),
pl.col("net_tx_bytes").quantile(0.95).alias("net_tx_p95"),
pl.col("net_drop_errors").sum().alias("net_drop_errors_total"),
pl.col("disk_iops").quantile(0.95).alias("disk_iops_p95"),
pl.col("auth_metadata_tps").quantile(0.95).alias("auth_tps_p95"),
# SRE 유도 행위 지표 대표값 도출
pl.col("feat_net_to_cpu_ratio").quantile(0.95).alias("feat_net_to_cpu_p95"),
pl.col("feat_io_asymmetry").quantile(0.95).alias("feat_io_asymmetry_p95"),
pl.col("feat_auth_density").quantile(0.95).alias("feat_auth_density_p95"),
pl.col("feat_load_skew_cv").max().alias("feat_load_skew_cv_max")
]).fill_null(0.0)
# (이하 기존 FinOps 사칙연산 및 거버넌스 상태 진단 로직 유지...)
이렇게 Step 1과 Step 2의 데이터 인터페이스를 완전히 정렬해 두면, 사람이 대시보드 눈금만 보고는 절대 인지할 수 없는 하이브리드 병목 패턴을 일별 배치 보고서 생성 시 사내 LLM이 완벽하게 발라낼 수 있게 됩니다.
feat_net_to_cpu_p95 지표가 비정상적으로 높고 cpu_throttled_max가 동반 상승한 패턴을 LLM이 확인하면, "단순한 계산 부하가 아니라 연산 노드 간 대량 데이터 셔플로 인한 Cilium CNI 레이어 컨텍스트 스위칭 포화"로 원인을 단번에 요약 판독합니다.disk_iops_p95는 엄청나게 격하게 튀는데 실제 디스크 전송 대역폭과 메모리 효율 지표는 유휴 상태인 좀비 패턴을 포착하여, LLM이 "Storage Engine 컴팩션(Compaction) 미비로 인한 소형 파일(Small-file) IO Read Backpressure가 의심되니 물리 파티션 버킷 재조정을 검토하라"는 고차원 솔루션을 사출해 주는 아키텍처 토대가 완성됩니다.