기존 step2 고도화 파이프라인에서 메모리 안전성과 속도를 극대화하기 위해 폴라(Polars)로 완벽하게 전환한 것처럼, 2차 심층 분석(Deep-Dive) 엔진의 데이터 적재 및 연산 레이어도 판다스 잔재를 완전히 걷어내고 Pure Polars 기반으로 마이그레이션을 완료했습니다.
특히 step2 마감 정산 단계에서 메모리 단위 보정(Bytes → GB) 및 4대 SRE 핵심 특징량 지표 집계가 이미 완료되었으므로, 2차 분석 스크립트에서는 중복 계산을 배제하고 폴라의 고속 벡터 집계 엔진을 활용해 워크로드 스택별 거동 프로필을 최적화된 JSON 구조로 사내 LLM에 인젝션하도록 재설계했습니다.
폴라의 고속 스캔 메커니즘을 가동하여 전일 정산 완료된 아키텍처 자산으로부터 컴포넌트 고유 징후를 격리 추출하는 흐름입니다.
./data/merged/ 내부에 정산된 daily_enriched_*.parquet 파일들을 폴라의 pl.read_parquet 명세로 고속 취합합니다.step2에서 정규화 완료된 workload_type 컬럼을 기준으로 [starrocks, spark, airflow, postgresql, aistor] 등 코어 스택 그룹의 장기 배치 이력 특징량을 롤업 집계합니다..to_dicts()로 고속 익스포트하여, 단순 사용량 통계와 4대 CNI/스토리지 상관관계 지표를 사내 LLM 지 지시문 포맷에 일대일 가속 결합합니다.이 스크립트는 내부에서 판다스(pd)를 전혀 사용하지 않으며, 메모리 복사 없이 폴라 고유의 병렬 처리 스펙으로 대량의 일별 정산 Parquet를 통합 처리합니다.
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-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 원부가 존재하지 않습니다.")
print(f" -> 총 {len(parquet_files)}개의 일별 고고도 파티션 감지. Polars 데이터 로딩...")
# 폴라 멀티 파일 고속 통합 로드
df_all = pl.read_parquet(parquet_files)
# 💡 [SRE 분석 정밀화]: 데이터 레이크하우스 주요 엔진 스택별 롤업 집계
# step2에서 이미 GB 및 기본 통계량 전처리가 끝났으므로, 여기서는 종합 대표 거동 특징량을 산출합니다.
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"),
# step2에서 이미 1024**3 스케일링이 끝난 자산이므로 추가 나누기 없이 다이렉트 바인딩
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"),
# ─── 🌐 [차세대 SRE 융합 핵심 유도 지표 대표값 도출] ───
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")
])
# 사내 LLM 조인 피딩을 위해 Polars 프레임을 딕셔너리 구조로 초고속 변환
profiles = {}
for row in df_profile.to_dicts():
w_type = row.pop("workload_type")
profiles[w_type] = row
return profiles
# ==========================================
# 3. 클라우드 네이티브 아키텍처 특화형 심층 프롬프트 빌더
# ==========================================
def build_cloud_native_deep_prompt(group_name, profile_dict):
profile_json_str = json.dumps(profile_dict, indent=2)
prompt = f"""
[Role]: 당신은 대규모 분산 고동시성 Cloud Native Data Lakehouse 인프라를 총괄 최적화하는 글로벌 SRE 인프라 수석 설계자입니다.
당신은 Kubernetes 쿼팅 메커니즘, Cilium의 eBPF 기반 네트워크 바이패스, MinIO AIStor의 DirectPV 로컬 NVMe 다이렉트 아이오, Keycloak의 OpenID Connect(OIDC) 인증 오버헤드를 완벽히 이해하고 있습니다.
[Target Workload Group]: {group_name}
[Rich Infrastructure Context]:
다음 데이터는 전일 [{group_name}] 엔진 가동 시 아키텍처 전체 레이어에서 실시간으로 역산 정제된 다차원 거동 지표 세트입니다.
{profile_json_str}
[SRE Expert Diagnosis 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 코어 데이터 노드)인지 인프라 레이어 간 연계 분석을 제공하세요.
2. **Keycloak IAM 및 메타데이터 인증 구조 진단**:
- `auth_density_avg` 수치를 판독하여 이 워크로드가 Keycloak 인증 레이어에 과도한 토큰 검증 부하를 가하고 있는지 격리 분석하세요.
3. **K8s 할당량 마모도 분석 및 하이브리드 리사이징 제안 (Quota Right-sizing)**:
- Request/Limit 가드레일 대비 P95 사용 곡선을 진단하고, 자원 절감(Slack 해소) 조치안 또는 OOM/Throttling 위험 차단을 위한 CPU/Memory 커스텀 튜닝 인사이트를 도출하세요.
[Output Format Rules]:
- 미사여구나 서론 없이, SRE 분석관의 엄격하고 건조한 기술적 어조로 마크다운 형태로만 출력하세요.
"""
return prompt
# ==========================================
# 4. 실행 파이프라인 메인 제어 루틴
# ==========================================
if __name__ == "__main__":
try:
print("1. Polars 커널 가동: 일별 Enriched Parquet 자산 통합 및 엔진 스택 추출...")
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 기반으로 사내 LLM을 경유해 종합 추론 분석한 결과서입니다.\n")
print(f"2. 감지된 {len(workload_profiles)}개 클라우드 네이티브 워크로드 그룹별 순회 분석 시작...")
for group_name, profile in workload_profiles.items():
if group_name == "unknown":
continue
print(f" -> [{group_name}] Polars 집계 맵 기반 LLM 쿼리 구성 중...")
cn_prompt = build_cloud_native_deep_prompt(group_name, profile)
payload = {
"model": "internal-cloudnative-sre-brain",
"messages": [{"role": "user", "content": cn_prompt}],
"temperature": 0.1,
"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}] 클라우드 네이티브 통합 분석 캔버스 빌드 완료.")
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"🚀 [최종 완료] Polars 파싱 기반 2차 연계 분석서가 '{output_path}'로 사출되었습니다.")
print(f"==========================================================")
except Exception as e:
print(f"💥 파이프라인 가동 에러 발생: {str(e)}")
classify_workload_group 정규식 파이썬 함수 연산을 배제했습니다. 이미 step2에서 정교하게 처리되어 파티셔닝된 workload_type 컬럼을 폴라가 그대로 스캔하기 때문에 성능 오버헤드가 완전히 소멸되었습니다.step2_pipeline.py 마감부에서 메모리 및 용량 지표(mem_request_max, mem_usage_p95, pv_usage_p95 등)에 대해 이미 1024**3 분할 보정이 적용되어 저장되므로, 해당 폴라 스크립트에서는 추가 변동 연산 없이 평균 지표(mean())로 깨끗하게 바인딩하여 데이터 무결성을 유지했습니다.step2/step3/step4로 이어지는 최신 명세 컨벤션인 feat_net_to_cpu_p95, feat_io_asymmetry_p95, feat_auth_density_p95, feat_load_skew_cv_max 네이밍을 정확하게 파싱하여 사내 LLM 브레인 인젝션 포맷으로 재정렬했습니다.