기존에 Kubernetes 오퍼레이션 컨트롤러 종류(StatefulSet, Deployment)로 추측했던 워크로드 개념을 데이터 레이크하우스 생태계의 실제 비즈니스 및 엔진 스택인 StarRocks(OLAP 분산 DB), Spark(분산 대용량 연산 엔진), Airflow(워크플로우 오케스트레이터), PostgreSQL(OLTP RDBMS) 등으로 전면 명확화하여 리빌딩했습니다.
merged 소스 파일 내부의 실제 워크로드 명칭 명세를 그대로 파싱하되, 사내 LLM이 각각의 데이터 스택 특성(예: 대용량 분산 쿼리를 수행하는 StarRocks와 배치 스파이크 성격이 강한 Spark의 리소스 소비 양상 차이)을 정확히 인지하여 맞춤형 아키텍처 가이드를 도출하도록 프롬프트 구조를 전면 선개조했습니다.
데이터 아키텍처 스택 단위로 데이터를 집계하고 복구 포인트를 역추적하는 실시간 데이터 파이프라인 흐름입니다.
workload_name 컬럼에 이미 기록된 starrocks, spark, airflow, postgresql 등의 실데이터 구조를 있는 그대로 메모리에 파싱합니다.import os
import glob
import json
import pandas as pd
import requests
# ==========================================
# 1. 인프라 환경 설정 및 데이터 경로 지정
# ==========================================
MERGED_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. 실제 애플리케이션 스택 기반 자원 편중(Skew) 분석 함수
# ==========================================
def analyze_engine_stack_metrics(directory):
csv_files = glob.glob(os.path.join(directory, "*.csv"))
if not csv_files:
raise FileNotFoundError(f"'{directory}' 내에 분석 가능한 CSV 파일이 존재하지 않습니다.")
full_df = pd.concat([pd.read_csv(f) for f in csv_files], ignore_index=True)
# [인프라 실제 데이터 스키마 기반 매핑]
# expected: cluster, namespace, node, workload_name (starrocks, spark, airflow 등), container,
# cpu_usage, cpu_req, cpu_limit, mem_usage, mem_req, mem_limit
# 1단계: 노드 단위로 분산 배치된 각 엔진 스택 컨테이너의 핵심 통계량 선산출
node_group = full_df.groupby([
'cluster', 'namespace', 'workload_name', 'container', 'node'
])
node_summary = node_group.agg(
cpu_p95=('cpu_usage', lambda x: x.quantile(0.95)),
cpu_max=('cpu_usage', 'max'),
cpu_req=('cpu_req', 'last'),
cpu_limit=('cpu_limit', 'last'),
mem_p95=('mem_usage', lambda x: x.quantile(0.95)),
mem_max=('mem_usage', 'max'),
mem_req=('mem_req', 'last'),
mem_limit=('mem_limit', 'last')
).reset_index()
# 2단계: 플랫폼 엔진 스택별 글로벌 Aggregation 및 노드 간 Skew(불균형 격차) 수식화
global_engine_group = node_summary.groupby([
'cluster', 'namespace', 'workload_name', 'container'
])
reduced_records = []
for (cluster, namespace, engine_name, container), group in global_engine_group:
# 분산 노드 간의 불균형 편차(Skew Delta) 계산
cpu_skew = group['cpu_p95'].max() - group['cpu_p95'].min()
mem_skew = group['mem_p95'].max() - group['mem_p95'].min()
# Quota 기준 정보 추적
cpu_req_avg = group['cpu_req'].mean()
cpu_lim_avg = group['cpu_limit'].mean()
mem_req_avg = group['mem_req'].mean()
mem_lim_avg = group['mem_limit'].mean()
# 실제 탑재 평균 실효성 지표 계산
cpu_p95_avg = group['cpu_p95'].mean()
mem_p95_avg = group['mem_p95'].mean()
cpu_efficiency = (cpu_p95_avg / cpu_req_avg * 100) if cpu_req_avg > 0 else 0
mem_efficiency = (mem_p95_avg / mem_req_avg * 100) if mem_req_avg > 0 else 0
reduced_records.append({
"Cluster": cluster,
"Namespace": namespace,
"EngineStack": engine_name, # starrocks, spark, airflow, postgresql 등이 동적 바인딩됨
"Container": container,
"CPU_Request": round(cpu_req_avg, 2),
"CPU_Limit": round(cpu_lim_avg, 2),
"CPU_P95_Avg": round(cpu_p95_avg, 2),
"CPU_Skew_Delta": round(cpu_skew, 2),
"CPU_Eff_%": round(cpu_efficiency, 1),
"MEM_Req_GB": round(mem_req_avg / 1024/1024/1024, 1) if mem_req_avg > 100000 else round(mem_req_avg, 1),
"MEM_Lim_GB": round(mem_lim_val / 1024/1024/1024, 1) if mem_lim_avg > 100000 else round(mem_lim_avg, 1),
"MEM_P95_GB": round(mem_p95_avg / 1024/1024/1024, 1) if mem_p95_avg > 100000 else round(mem_p95_avg, 1),
"MEM_Skew_Delta": round(mem_skew / 1024/1024/1024, 1) if mem_skew > 100000 else round(mem_skew, 1),
"MEM_Eff_%": round(mem_efficiency, 1)
})
return pd.DataFrame(reduced_records)
# ==========================================
# 3. 플랫폼 엔진 스택 분류 특화형 SRE 프롬프트
# ==========================================
def build_engine_specialized_prompt(summary_df):
data_table_str = summary_df.to_markdown(index=False)
prompt = f"""
[Role]: 당신은 대규모 데이터 레이크하우스 아키텍처 인프라를 전담 운영 및 최적화하는 수석 SRE 및 플랫폼 엔지니어입니다.
[Context]: 다음 테이블은 사내 데이터 인프라의 핵심 엔진 스택별(`EngineStack`) 자원 Quota 대비 실사용 효율성 및 분산 노드 간 자원 편중도(`Skew_Delta`) 요약 정보입니다.
### 📊 전일 인프라 코어 엔진 스택별 Quota & Skew 분석 테이블
{data_table_str}
[Instruction]:
제시된 통계 데이터를 바탕으로 각 엔진 스택의 '소프트웨어 고유 아키텍처 특성'을 전적으로 투영하여 다음 3가지 섹션으로 정밀 진단 리포트를 작성해 주세요. 단순 수치 요약이 아닌 각 시스템의 작동 메커니즘에 기반해야 합니다.
1. **전반적인 자원 분포 및 스택별 Skew(편중) 상황 Overview**:
- `EngineStack` 컬럼에 명시된 워크로드(starrocks, spark, airflow, postgresql 등)별로 전체적인 리소스 불균형 상태를 해석하세요.
- **StarRocks**: OLAP 분산 질의 엔진 특성상 분산 계산 노드 간에 `Skew_Delta`가 크다면, 이는 특정 버킷/파티션 셔플 키 불균형(Data Skewness)으로 인해 특정 컨테이너 코어만 독식하여 타오르는 현상인지 커널 관점에서 판독하세요.
- **Spark**: 대규모 임시 분산 연산 시 발생하는 스파이크 특성을 고려하여 자원 분포를 평가하세요.
2. **할당량(Request/Limit) 대비 실효성 진단 (Over/Under Provisioning)**:
- **StarRocks / PostgreSQL**: 상시 가동형 대형 분산 DB 및 관계형 DB 컴포넌트입니다. 이들의 메모리 효율(`MEM_Eff_%`)이 너무 낮다면 shared_buffers나 캐시 영역이 낭비(Slack)되는 상태인지, 반대로 `Limit`선에 딱 붙어있어 OOM(Out-of-Memory) 킬러의 타겟이 되기 일보 직전인지 물리적 자원 한계점을 진단하세요.
- **Spark**: Ephemeral(일시적) 배치 워커 엔진입니다. `CPU_Eff_%`나 Limit 마진을 대조하여 과도한 버스팅(Bursting)으로 인해 CPU 스로틀링(Throttling)이 발생해 대형 모델 학습 및 연산 파이프라인의 전체 타임라인을 지연시키고 있지는 않은지 격리 진단하세요.
- **Airflow**: 스케줄링 관리자로 평소 베이스라인 자원은 낮으나 특정 시간대 DAG 기동 시 튀는 특성을 반영하세요.
3. **엔진 스택별 인프라 조치 우선순위 및 인사이트**:
- 운영 긴급도 및 성능 아키텍처 충격도를 고려하여 엔지니어가 오늘 당장 튜닝해야 할 리사이징(Resizing) 또는 구성 보정 조치 명령과 플레이북 우선순위 목록을 도출하세요.
- 단순 리소스 확대/축소가 아닌 '파티션 키 재조정', 'Spark Executor Memory 할당량 보정', 'DB Connection Pool 조정' 등 스택에 부합하는 솔루션을 제시해야 합니다.
[Output Format Rules]:
- 서론과 맺음말 없이 핵심 기술 통찰 보고서 형태로 마크다운으로 깔끔하게 사출해 주세요.
"""
return prompt
# ==========================================
# 4. 실행 파이프라인 메인 루틴
# ==========================================
if __name__ == "__main__":
try:
print("1. merged 데이터 파일에서 starrocks, spark, airflow 등 실제 엔진 스택 요약 파싱 가동...")
summary_data = analyze_engine_stack_metrics(MERGED_DATA_DIR)
print(f" -> 압축 성공! 총 {len(summary_data)}개의 독립 플랫폼 엔진 컴포넌트 추출.")
print("2. 스택 고유 특성을 반영한 하이브리드 SRE 프롬프트 동적 생성 중...")
final_prompt = build_engine_specialized_prompt(summary_data)
payload = {
"model": "internal-sre-engine-brain",
"messages": [{"role": "user", "content": final_prompt}],
"temperature": 0.15,
"max_tokens": 3072
}
headers = {
"Authorization": f"Bearer {INTERNAL_LLM_API_KEY}",
"Content-Type": "application/json"
}
print(f"3. 사내 LLM 엔드포인트로 스택 융합 인프라 증적 송신...")
response = requests.post(INTERNAL_LLM_URL, headers=headers, data=json.dumps(payload), timeout=90)
if response.status_code == 200:
result_json = response.json()
llm_analysis = result_json['choices'][0]['message']['content']
print("\n==========================================================")
print("🚀 사내 LLM 플랫폼 코어 스택별 자원 분석 및 아키텍처 리포트")
print("==========================================================\n")
print(llm_analysis)
with open("daily_engine_stack_insight.md", "w", encoding="utf-8") as f:
f.write(llm_analysis)
print("\n[완료] daily_engine_stack_insight.md 문서 동기화 완료.")
else:
print(f"❌ 사내 LLM 호출 실패: {response.status_code}")
print(response.text)
except Exception as e:
print(f"💥 파이프라인 구동 에러: {str(e)}")
사내 LLM에 이 정제본을 밀어 넣은 뒤 나오는 답변 형태를 모니터링하면서 프롬프트 템플릿의 지시어 파트를 튜닝(Fine-tuning)해 나갈 핵심 운영 가이드라인입니다.
[Spark 튜닝 규칙: Spark 엔진은 P95 대비 Max 수치의 격차가 4배 이상 벌어지는 것이 정상적인 배치 양상이므로, 단순히 P95가 낮다는 이유로 Request를 무리하게 감축하지 말고 Limit 버스팅 폭을 넓히도록 제안할 것]이라는 가드레일 조건을 보강합니다.Skew_Delta가 튈 때 LLM의 해석 일관성을 부여합니다.[Database Skew 판정: StarRocks 컴포넌트의 MEM_Skew_Delta가 32GB 이상으로 기록될 경우 최우선 순위 조치 대상으로 지정하고 파티션 분산 해시 키 검증을 솔루션으로 강제 매핑할 것]을 명시해 둠으로써 내부 의사결정 브레인을 지속 튜닝해 나갈 수 있습니다.실제 데이터 스택 명세에 정확하게 밀착하여 분산 DB와 대용량 연산 파이프라인의 효율을 정밀 판독하는 생산용 분석 프레임워크가 가동 준비 완료되었습니다.