아, 새로 정립한 클라우드 네이티브 3대 섹션 구조(Cilium/DirectPV, Keycloak, Quota)를 대장 뼈대로 삼고, 그 하위에 엔진별(StarRocks, Spark, Airflow, PostgreSQL, AIStor) 아키텍처 메커니즘 지침을 조화롭게 배치했어야 했는데 구형 템플릿의 순서에 너무 얽매였습니다. 제 불찰입니다.
원하셨던 [Instruction]: 1, 2, 3 대분류 위계 구조와 [Output Format Rules] 스펙을 정확하게 복원하고, 그 내부에 플랫폼 엔진 스택별 고유의 분산/캐시 작동 방식을 빈틈없이 짜 넣은 최종 완성형 step5_llm_diagnoser.py 전체 스크립트입니다.
step5_llm_diagnoser.py (완전 통합형 프롬프트 엔진)이 스크립트는 폴라(Polars)를 이용해 데이터를 효율적으로 핸들링하며, 사내 LLM에 쏘는 프롬프트의 지시문([Instruction])을 정확하게 약속된 3대 클라우드 네이티브 도메인 분류로 격리 사출합니다.
import os
import glob
import json
import polars as pl
import requests
from pathlib import Path
# ==========================================
# 1. 인프라 환경 설정 및 데이터 경로 지정
# ==========================================
BASE_DATA_DIR = Path("./data")
MERGED_DIR = BASE_DATA_DIR / "merged"
INTERNAL_LLM_URL = "https://llm.internal.company.net/v1/chat/completions"
INTERNAL_LLM_API_KEY = os.getenv("INTERNAL_LLM_API_KEY", "your-sre-token-here")
# ==========================================
# 2. Polars 기반 엔진 스택 요약 집계 추출 커널
# ==========================================
def extract_rich_features_with_polars():
# COMPUTE 및 STORAGE 모든 격리 하위 폴더의 enriched 자산 동적 스캔
parquet_pattern = str(MERGED_DIR / "**" / "daily_enriched_*.parquet")
parquet_files = glob.glob(parquet_pattern, recursive=True)
if not parquet_files:
raise FileNotFoundError(f"'{MERGED_DIR}' 하방에 정산된 daily_enriched_*.parquet 원부가 존재하지 않습니다.")
# 폴라 멀티 파일 고속 통합 로드
df_all = pl.read_parquet(parquet_files)
# 1분 해상도 기반으로 정산 완료된 아키텍처 자산의 메타 그룹 롤업 집계[cite: 2]
df_profile = df_all.group_by("workload_type").agg([
pl.col("minutes_running").sum().alias("total_minutes_running"),
pl.col("cpu_request_max").mean().round(1).alias("cpu_req_avg"),
pl.col("cpu_limit_max").mean().round(1).alias("cpu_lim_avg"),
pl.col("cpu_usage_p95").mean().round(2).alias("cpu_usage_p95_avg"),
pl.col("cpu_usage_max").max().round(2).alias("cpu_usage_max_highest"),
pl.col("cpu_throttled_max").max().round(2).alias("cpu_throttled_peak"),
pl.col("mem_request_max").mean().round(1).alias("mem_req_gb_avg"),
pl.col("mem_limit_max").mean().round(1).alias("mem_lim_gb_avg"),
pl.col("mem_usage_p95").mean().round(1).alias("mem_usage_p95_gb_avg"),
pl.col("oom_strike_sum").sum().alias("oom_strike_total"),
# 🌐 [클라우드 네이티브 확장 지표 대표값 도출][cite: 2]
pl.col("feat_net_to_cpu_p95").mean().round(2).alias("net_to_cpu_ratio_avg"),
pl.col("feat_io_asymmetry_p95").mean().round(2).alias("io_asymmetry_avg"),
pl.col("feat_auth_density_p95").mean().round(3).alias("auth_density_avg"),
pl.col("feat_load_skew_cv_max").max().round(3).alias("load_skew_cv_highest"),
pl.col("net_drop_errors_total").sum().alias("net_drop_errors_sum")
])
profiles = {}
for row in df_profile.to_dicts():
w_type = row.pop("workload_type")
profiles[w_type] = row
return profiles
# ==========================================
# 3. 오리지널 스택 아키텍처 + 네이티브 통합 프롬프트 빌더
# ==========================================
def build_combined_sre_prompt(group_name, profile_dict):
profile_json_str = json.dumps(profile_dict, indent=2)
prompt = f"""
[Role]: 당신은 멀티 클러스터 환경에서 Kubernetes 기반 대규모 데이터 레이크하우스(K8s, Cilium, MinIO AIStor, Keycloak, StarRocks, Spark) 인프라를 총괄 최적화하는 수석 SRE 및 플랫폼 아키텍처 엔지니어입니다.
[Target Workload Group]: {group_name}
[Rich Infrastructure Context]:
다음 데이터는 전일 [{group_name}] 엔진 가동 시 아키텍처 전체 레이어에서 1분 해상도로 역산 정제된 다차원 거동 지표 세트입니다.
{profile_json_str}
[SRE Expert Cross-Layer Rules]:
1. **Cilium Network Pipeline 포화 규칙**: `net_to_cpu_ratio_avg` 수치가 극단적으로 크다면, 순수 컴포넌트 CPU 연산이 아니라 노드 간 대량 패킷 전송(Cilium 인프라 전송망)이 병목의 주원인망임을 뜻함.
2. **AIStor / DirectPV 스토리지 스트레스 규칙**: `io_asymmetry_avg` 가 극단적으로 높다면(>10.0) 데이터 백엔드 플러시 및 복제(Replication) 오퍼레이션이 밀리고 있음을 의미.
3. **Keycloak 인증 폭풍(Token Storm) 규칙**: `auth_density_avg` 가 높고 CPU 피크가 동반된다면, 대량의 단기 커넥션 세션 요청이나 토큰 검증(OIDC Validation Hook) 부하가 엔진의 발목을 잡고 있음을 뜻함.
4. **Data Skewness 판단 규칙**: `load_skew_cv_highest` 수치가 0.25를 상회하면 분산 노드 간 쿼리 샤딩 밸런싱이 무너진 것으로 판단함.
[Instruction]:
이 규칙 가이드라인을 데이터와 1:1 매핑 대조하여, 전일 하루 동안 [{group_name}] 엔진이 우리 클라우드 네이티브 레이크하우스 전체 인프라에 가한 충격과 거동 특성을 정밀 분석 리포트로 작성해 주세요. 단순 수치 요약이 아닌 각 시스템의 작동 메커니즘에 기반해야 합니다.
1. **인프라 파이프라인 연계 거동 분석 (Cilium & DirectPV Storage Impact)**:
- 계산된 네트워크-컴퓨트 비율 및 스토리지 비대칭 지수를 바탕으로, 이 엔진이 Cilium 패킷 전송 부하를 주로 일으키는 녀석(e.g., Spark Shuffle, StarRocks MPP Node Exchange)인지, 아니면 DirectPV를 통한 로컬 NVMe 스토리지 I/O 단독 병목을 일으키는 녀석(e.g., AIStor 코어 데이터 노드)인지 인프라 레이어 간 연계 분석을 제공하세요.
- **StarRocks**: OLAP 분산 질의 엔진 특성상 분산 계산 노드 간에 부하 Skew(`load_skew_cv_highest`)가 크다면, 이는 특정 버킷/파티션 셔플 키 불균형(Data Skewness)으로 인해 특정 계산 컨테이너 코어만 독식하여 타오르는 현상인지 커널 및 MPP 데이터 익스체인지 관점에서 판독하세요.
- **Spark**: 대규모 임시 분산 연산 시 발생하는 스파이크 특성을 고려하여 자원 분포를 평가하고, Cilium의 대역폭 셔플 지수(`net_to_cpu_ratio_avg`) 및 패킷 드롭 오염 상태와 연계하여 노드 밸런싱을 분석하세요.
- **AIStor (objectstore / directpv)**: 로컬 NVMe 스토리지 I/O 단독 병목이나 특정 볼륨으로의 분산 토폴로지 균열 징후가 감지되는지 추론하세요.
2. **Keycloak IAM 및 메타데이터 인증 구조 진단**:
- `auth_density_avg` 수치를 판독하여 이 워크로드가 Keycloak 인증 레이어에 과도한 토큰 검증 부하를 가하고 있는지 격리 분석하세요.
- **Airflow**: 스케줄링 관리자로 평소 베이스라인 자원은 낮으나 특정 시간대 DAG 기동 시 튀는 특성을 반영하세요. 특히, 대량 태스크 생성 시 Keycloak 인증 부하 밀도(`auth_density_avg`)를 자극해 토큰 검증 폭풍(Token Storm) 오버헤드를 인프라에 전이시키고 있는지 추론하세요.
- **Spark / StarRocks**: Executor 또는 MPP 노드의 단기 커넥션 세션 요청이나 무차별 API 호출이 Keycloak 세션 고갈 및 토큰 인증 오버헤드를 동반하고 있는지 판정하세요.
3. **K8s 할당량 마모도 분석 및 하이브리드 리사이징 제안 (Quota Right-sizing)**:
- Request/Limit 가드레일 대비 P95 사용 곡선을 진단하고, 자원 절감(Slack 해소) 조치안 또는 OOM/Throttling 위험 차단을 위한 CPU/Memory 커스텀 튜닝 인사이트를 도출하세요.
- **StarRocks / PostgreSQL**: 상시 가동형 대형 분산 DB 및 관계형 DB 컴포넌트입니다. 이들의 메모리 효율이 너무 낮다면 shared_buffers나 캐시 영역이 낭비(Slack)되는 상태인지, 반대로 `Limit`선에 딱 붙어있거나 `oom_strike_total`이 발생하여 OOM(Out-of-Memory) 킬러의 타겟이 되기 일보 직전인지 물리적 자원 한계점을 진단하세요.
- **Spark**: Ephemeral(일시적) 배치 워커 엔진입니다. 실사용 P95/Max 사용률과 Limit 마진을 대조하여 과도한 버스팅(Bursting)으로 인해 `cpu_throttled_peak`이 발생해 연산 파이프라인의 전체 타임라인을 지연시키고 있지는 않은지 격리 진단하세요.
- **엔진별 조치 권고**: 단순 리소스 확대/축소가 아닌 '파티션 키 재조정(StarRocks Bucket Tuning)', 'Spark Executor Memory 및 Buffer 할당량 보정', 'DB Connection Pool 및 shared_buffers 조정', 'Keycloak OIDC 토큰 캐시 라이프타임 보정' 등 스택에 부합하는 솔루션을 명확한 수치 근거와 함께 제시해야 합니다.
[Output Format Rules]:
- 서론과 맺음말 없이 핵심 기술 통찰 보고서 형태로 마크다운으로 깔끔하게 사출해 주세요.
"""
return prompt
# ==========================================
# 4. 실행 파이프라인 메인 제어 루틴
# ==========================================
if __name__ == "__main__":
try:
print("1. [step5] Polars 기반 하이브리드 인프라 자산 프로파일러 가동...")
workload_profiles = extract_rich_features_with_polars()
combined_report = []
combined_report.append("# 🌐 2nd-Stage Cloud-Native Data Lakehouse Workload Behavioral Deep-Dive\n")
combined_report.append("본 보고서는 Polars 코어로 병합된 v4 확장 지표 Parquet를 기반으로, 소프트웨어 커널 메커니즘과 네이티브 CNI/스토리지 역학관계를 연계 분석한 SRE 리포트입니다.\n")
print(f"2. 감지된 {len(workload_profiles)}개 코어 엔진 그룹별 심층 추론 시작...")
for group_name, profile in workload_profiles.items():
if group_name == "unknown":
continue
print(f" -> [{group_name}] 스택 아키텍처 특화형 프롬프트 주입 및 송신 중...")
cn_prompt = build_combined_sre_prompt(group_name, profile)
payload = {
"model": "internal-cloudnative-sre-brain",
"messages": [{"role": "user", "content": cn_prompt}],
"temperature": 0.15,
"max_tokens": 3000
}
headers = {
"Authorization": f"Bearer {INTERNAL_LLM_API_KEY}",
"Content-Type": "application/json"
}
response = requests.post(INTERNAL_LLM_URL, headers=headers, data=json.dumps(payload), timeout=90)
if response.status_code == 200:
result_json = response.json()
analysis_text = result_json['choices'][0]['message']['content']
combined_report.append(f"## 🛠️ Engine Topology Analysis: {group_name.upper()}")
combined_report.append(analysis_text)
combined_report.append("\n---\n")
print(f" ✓ [{group_name}] SRE 결합 분석 완료.")
else:
print(f" ❌ [{group_name}] LLM 추론 실패 (코드: {response.status_code})")
# 종합 파일 저장
output_path = "daily_cloud_native_behavior_report.md"
with open(output_path, "w", encoding="utf-8") as f:
f.write("\n".join(combined_report))
print(f"\n==========================================================")
print(f"🚀 [최종 성료] 아키텍처 메커니즘이 완전히 복원된 리포트가 '{output_path}'로 동기화되었습니다.")
print(f"==========================================================")
except Exception as e:
print(f"💥 step5 파이프라인 가동 에러 발생: {str(e)}")