from pydantic_settings import BaseSettings
class Settings(BaseSettings):
OLLAMA_BASE_URL: str = "http://ollama-service:11434"
OLLAMA_MODEL: str = "qwen2.5:14b"
OLLAMA_EMBED_MODEL: str = "nomic-embed-text"
OLLAMA_EMBED_DIM: int = 768
PROMETHEUS_URL: str = "http://prometheus:9090"
OPENSEARCH_HOST: str = "opensearch"
OPENSEARCH_PORT: int = 9200
OPENSEARCH_USER: str = "admin"
OPENSEARCH_PASSWORD: str = "admin"
OPENSEARCH_INDEX_PATTERN: str = "logs-*"
OPENSEARCH_USE_SSL: bool = False
MILVUS_HOST: str = "milvus"
MILVUS_PORT: int = 19530
MILVUS_USER: str = ""
MILVUS_PASSWORD: str = ""
MILVUS_INCIDENT_COLLECTION: str = "incident_history"
MILVUS_RUNBOOK_COLLECTION: str = "runbook_knowledge"
MILVUS_NORMAL_LOG_COLLECTION: str = "normal_log_baseline"
class Config:
env_file = ".env"
settings = Settings()
"""
Milvus 컬렉션 초기화 스크립트
실행: python scripts/init_milvus.py
"""
from pymilvus import (
connections, Collection, CollectionSchema,
FieldSchema, DataType, utility
)
from config import settings
def connect():
connections.connect(
alias="default",
host=settings.MILVUS_HOST,
port=settings.MILVUS_PORT,
user=settings.MILVUS_USER or None,
password=settings.MILVUS_PASSWORD or None,
)
print(f"✓ Milvus 연결: {settings.MILVUS_HOST}:{settings.MILVUS_PORT}")
def create_incident_collection():
"""과거 인시던트 히스토리 컬렉션"""
name = settings.MILVUS_INCIDENT_COLLECTION
if utility.has_collection(name):
print(f" 이미 존재: {name}")
return Collection(name)
schema = CollectionSchema(fields=[
FieldSchema("id", DataType.INT64, is_primary=True, auto_id=True),
FieldSchema("embedding", DataType.FLOAT_VECTOR, dim=settings.OLLAMA_EMBED_DIM),
FieldSchema("title", DataType.VARCHAR, max_length=500),
FieldSchema("error_summary", DataType.VARCHAR, max_length=2000),
FieldSchema("root_cause", DataType.VARCHAR, max_length=2000),
FieldSchema("resolution", DataType.VARCHAR, max_length=3000),
FieldSchema("severity", DataType.VARCHAR, max_length=20),
FieldSchema("service", DataType.VARCHAR, max_length=200),
FieldSchema("namespace", DataType.VARCHAR, max_length=200),
FieldSchema("occurred_at", DataType.VARCHAR, max_length=50),
FieldSchema("tags", DataType.VARCHAR, max_length=500),
], description="과거 인시던트 히스토리")
col = Collection(name=name, schema=schema)
col.create_index(
field_name="embedding",
index_params={
"metric_type": "COSINE",
"index_type": "HNSW",
"params": {"M": 16, "efConstruction": 200},
},
)
print(f" ✓ 생성 완료: {name}")
return col
def create_runbook_collection():
"""운영 Runbook 지식베이스 컬렉션"""
name = settings.MILVUS_RUNBOOK_COLLECTION
if utility.has_collection(name):
print(f" 이미 존재: {name}")
return Collection(name)
schema = CollectionSchema(fields=[
FieldSchema("id", DataType.INT64, is_primary=True, auto_id=True),
FieldSchema("embedding", DataType.FLOAT_VECTOR, dim=settings.OLLAMA_EMBED_DIM),
FieldSchema("title", DataType.VARCHAR, max_length=500),
FieldSchema("content", DataType.VARCHAR, max_length=5000),
FieldSchema("category", DataType.VARCHAR, max_length=100),
FieldSchema("source_file", DataType.VARCHAR, max_length=300),
FieldSchema("chunk_index", DataType.INT32),
], description="운영 Runbook 지식베이스")
col = Collection(name=name, schema=schema)
col.create_index(
field_name="embedding",
index_params={
"metric_type": "COSINE",
"index_type": "HNSW",
"params": {"M": 16, "efConstruction": 200},
},
)
print(f" ✓ 생성 완료: {name}")
return col
def create_normal_log_collection():
"""정상 로그 베이스라인 컬렉션 (이상 탐지용)"""
name = settings.MILVUS_NORMAL_LOG_COLLECTION
if utility.has_collection(name):
print(f" 이미 존재: {name}")
return Collection(name)
schema = CollectionSchema(fields=[
FieldSchema("id", DataType.INT64, is_primary=True, auto_id=True),
FieldSchema("embedding", DataType.FLOAT_VECTOR, dim=settings.OLLAMA_EMBED_DIM),
FieldSchema("log_message", DataType.VARCHAR, max_length=1000),
FieldSchema("container", DataType.VARCHAR, max_length=200),
FieldSchema("namespace", DataType.VARCHAR, max_length=200),
FieldSchema("sampled_at", DataType.VARCHAR, max_length=50),
], description="정상 로그 벡터 베이스라인")
col = Collection(name=name, schema=schema)
col.create_index(
field_name="embedding",
index_params={
"metric_type": "L2",
"index_type": "HNSW",
"params": {"M": 16, "efConstruction": 200},
},
)
print(f" ✓ 생성 완료: {name}")
return col
if __name__ == "__main__":
connect()
print("\n[컬렉션 초기화]")
create_incident_collection()
create_runbook_collection()
create_normal_log_collection()
print("\n✅ Milvus 초기화 완료")
"""
Runbook 문서 Milvus 적재 스크립트
실행: python scripts/ingest_runbook.py --dir data/runbooks/
"""
import argparse
import re
from pathlib import Path
from pymilvus import connections, Collection
from agent.tools_milvus import get_embedding
from config import settings
def chunk_markdown(text: str, chunk_size: int = 800, overlap: int = 100) -> list[str]:
"""마크다운을 섹션 단위로 청킹"""
sections = re.split(r'\n(?=#{1,3} )', text)
chunks = []
for section in sections:
if len(section) <= chunk_size:
if section.strip():
chunks.append(section.strip())
else:
words = section.split()
buf = []
buf_len = 0
for word in words:
buf.append(word)
buf_len += len(word) + 1
if buf_len >= chunk_size:
chunks.append(" ".join(buf))
buf = buf[-overlap:]
buf_len = sum(len(w) + 1 for w in buf)
if buf:
chunks.append(" ".join(buf))
return [c for c in chunks if len(c) > 50]
def ingest_directory(doc_dir: str):
connections.connect(host=settings.MILVUS_HOST, port=settings.MILVUS_PORT)
col = Collection(settings.MILVUS_RUNBOOK_COLLECTION)
doc_path = Path(doc_dir)
md_files = list(doc_path.glob("**/*.md"))
print(f"📄 발견된 문서: {len(md_files)}개")
total_chunks = 0
for md_file in md_files:
text = md_file.read_text(encoding="utf-8")
chunks = chunk_markdown(text)
category = md_file.parent.name if md_file.parent != doc_path else "general"
print(f" └ {md_file.name}: {len(chunks)} chunks (category: {category})")
embeddings = []
titles, contents, categories, sources, indices = [], [], [], [], []
for i, chunk in enumerate(chunks):
emb = get_embedding(chunk)
embeddings.append(emb)
titles.append(md_file.stem[:499])
contents.append(chunk[:4999])
categories.append(category[:99])
sources.append(str(md_file)[:299])
indices.append(i)
col.insert([embeddings, titles, contents, categories, sources, indices])
total_chunks += len(chunks)
col.flush()
print(f"\n✅ 총 {total_chunks}개 청크 적재 완료")
if __name__ == "__main__":
parser = argparse.ArgumentParser()
parser.add_argument("--dir", default="data/runbooks/")
args = parser.parse_args()
ingest_directory(args.dir)
import json
import requests
from typing import Optional
from datetime import datetime
from pymilvus import connections, Collection, utility
from langchain_core.tools import tool
from config import settings
_connected = False
def ensure_connected():
global _connected
if not _connected:
connections.connect(
alias="default",
host=settings.MILVUS_HOST,
port=settings.MILVUS_PORT,
user=settings.MILVUS_USER or None,
password=settings.MILVUS_PASSWORD or None,
)
_connected = True
def get_embedding(text: str) -> list[float]:
"""Ollama로 텍스트 임베딩 생성"""
response = requests.post(
f"{settings.OLLAMA_BASE_URL}/api/embeddings",
json={"model": settings.OLLAMA_EMBED_MODEL, "prompt": text},
timeout=30,
)
response.raise_for_status()
return response.json()["embedding"]
def get_embeddings_batch(texts: list[str]) -> list[list[float]]:
"""배치 임베딩 (여러 텍스트 한번에)"""
return [get_embedding(t) for t in texts]
@tool
def search_similar_incidents(
error_description: str,
top_k: int = 3,
service_filter: Optional[str] = None,
) -> str:
"""
현재 에러/증상과 유사한 과거 인시던트를 Milvus에서 검색합니다.
과거 해결책과 Root Cause를 참고해 현재 문제 해결에 활용합니다.
Args:
error_description: 현재 발생 중인 에러나 증상 설명
top_k: 반환할 유사 인시던트 수 (기본 3)
service_filter: 특정 서비스로 필터 (예: "minio", "spark")
"""
ensure_connected()
col = Collection(settings.MILVUS_INCIDENT_COLLECTION)
col.load()
try:
query_vector = get_embedding(error_description)
search_params = {"metric_type": "COSINE", "params": {"ef": 100}}
expr = f'service == "{service_filter}"' if service_filter else None
results = col.search(
data=[query_vector],
anns_field="embedding",
param=search_params,
limit=top_k,
expr=expr,
output_fields=[
"title", "error_summary", "root_cause",
"resolution", "severity", "service",
"namespace", "occurred_at", "tags",
],
)
incidents = []
for hit in results[0]:
entity = hit.entity
incidents.append({
"similarity_score": round(hit.score, 4),
"title": entity.get("title"),
"severity": entity.get("severity"),
"service": entity.get("service"),
"occurred_at": entity.get("occurred_at"),
"error_summary": entity.get("error_summary"),
"root_cause": entity.get("root_cause"),
"resolution": entity.get("resolution"),
"tags": entity.get("tags"),
})
if not incidents:
return json.dumps({"message": "유사한 과거 인시던트를 찾지 못했습니다.", "incidents": []})
return json.dumps({
"query": error_description,
"similar_incidents": incidents,
}, ensure_ascii=False)
except Exception as e:
return f"유사 인시던트 검색 오류: {str(e)}"
finally:
col.release()
@tool
def search_runbook(
situation_description: str,
category: Optional[str] = None,
top_k: int = 3,
) -> str:
"""
현재 상황에 맞는 운영 Runbook 내용을 검색합니다.
장애 대응 절차, 점검 체크리스트, 설정 가이드 등을 찾아 조치 방법을 제안합니다.
Args:
situation_description: 현재 상황 설명 (에러 증상, 하고 싶은 작업 등)
category: 카테고리 필터 (minio, k8s, cilium, spark, trino 등)
top_k: 반환할 문서 청크 수
"""
ensure_connected()
col = Collection(settings.MILVUS_RUNBOOK_COLLECTION)
col.load()
try:
query_vector = get_embedding(situation_description)
expr = f'category == "{category}"' if category else None
search_params = {"metric_type": "COSINE", "params": {"ef": 100}}
results = col.search(
data=[query_vector],
anns_field="embedding",
param=search_params,
limit=top_k,
expr=expr,
output_fields=["title", "content", "category", "source_file", "chunk_index"],
)
runbooks = []
for hit in results[0]:
entity = hit.entity
runbooks.append({
"relevance_score": round(hit.score, 4),
"title": entity.get("title"),
"category": entity.get("category"),
"source": entity.get("source_file"),
"content": entity.get("content"),
})
if not runbooks:
return json.dumps({"message": "관련 Runbook을 찾지 못했습니다.", "runbooks": []})
return json.dumps({
"situation": situation_description,
"relevant_runbooks": runbooks,
}, ensure_ascii=False)
except Exception as e:
return f"Runbook 검색 오류: {str(e)}"
finally:
col.release()
@tool
def detect_anomaly_logs(
log_messages: list[str],
container: Optional[str] = None,
anomaly_threshold: float = 0.8,
) -> str:
"""
로그 메시지들이 정상 베이스라인과 얼마나 다른지(거리) 측정하여 이상 로그를 탐지합니다.
threshold 이상의 거리를 가진 로그는 이상 로그로 판단합니다.
Args:
log_messages: 분석할 로그 메시지 리스트 (OpenSearch에서 가져온 원시 로그)
container: 베이스라인 필터용 컨테이너 이름
anomaly_threshold: 이상 판단 임계값 (L2 거리, 높을수록 이상)
"""
ensure_connected()
if not utility.has_collection(settings.MILVUS_NORMAL_LOG_COLLECTION):
return "정상 로그 베이스라인이 없습니다. 먼저 store_normal_log_baseline을 실행하세요."
col = Collection(settings.MILVUS_NORMAL_LOG_COLLECTION)
col.load()
try:
anomalies = []
normal_count = 0
for log_msg in log_messages[:30]:
if not log_msg or len(log_msg) < 10:
continue
query_vector = get_embedding(log_msg)
expr = f'container == "{container}"' if container else None
results = col.search(
data=[query_vector],
anns_field="embedding",
param={"metric_type": "L2", "params": {"ef": 50}},
limit=1,
expr=expr,
output_fields=["log_message"],
)
if results[0]:
distance = results[0][0].score
if distance > anomaly_threshold:
anomalies.append({
"log": log_msg[:300],
"distance_from_normal": round(distance, 4),
"most_similar_normal_log": results[0][0].entity.get("log_message", "")[:200],
})
else:
normal_count += 1
else:
anomalies.append({
"log": log_msg[:300],
"distance_from_normal": -1,
"most_similar_normal_log": "베이스라인 없음",
})
anomalies.sort(key=lambda x: x["distance_from_normal"], reverse=True)
return json.dumps({
"total_analyzed": len(log_messages),
"anomaly_count": len(anomalies),
"normal_count": normal_count,
"anomaly_rate": round(len(anomalies) / max(len(log_messages), 1), 3),
"top_anomalies": anomalies[:10],
}, ensure_ascii=False)
except Exception as e:
return f"이상 탐지 오류: {str(e)}"
finally:
col.release()
@tool
def store_resolved_incident(
title: str,
error_summary: str,
root_cause: str,
resolution: str,
severity: str,
service: str,
namespace: str = "",
tags: str = "",
) -> str:
"""
해결된 인시던트를 Milvus에 저장합니다.
향후 유사 인시던트 분석 시 참고 자료로 활용됩니다.
Args:
title: 인시던트 제목
error_summary: 에러 증상 요약 (임베딩 대상)
root_cause: 근본 원인
resolution: 해결 방법
severity: 심각도 (critical/warning/info)
service: 영향받은 서비스 (minio, spark, trino 등)
namespace: k8s 네임스페이스
tags: 태그 (쉼표 구분, 예: "oom,memory,jvm")
"""
ensure_connected()
col = Collection(settings.MILVUS_INCIDENT_COLLECTION)
try:
embed_text = f"{title}\n{error_summary}\n{root_cause}"
embedding = get_embedding(embed_text)
data = [
[embedding],
[title[:499]],
[error_summary[:1999]],
[root_cause[:1999]],
[resolution[:2999]],
[severity[:19]],
[service[:199]],
[namespace[:199]],
[datetime.utcnow().isoformat()],
[tags[:499]],
]
col.insert(data)
col.flush()
return json.dumps({
"status": "저장 완료",
"title": title,
"service": service,
"severity": severity,
}, ensure_ascii=False)
except Exception as e:
return f"인시던트 저장 오류: {str(e)}"
@tool
def store_normal_log_baseline(
log_messages: list[str],
container: str,
namespace: str = "",
) -> str:
"""
정상 상태의 로그들을 베이스라인으로 Milvus에 저장합니다.
이후 이상 탐지의 기준점이 됩니다.
평온한 시간대의 로그를 OpenSearch에서 가져와서 등록하세요.
Args:
log_messages: 정상 로그 메시지 리스트
container: 컨테이너 이름
namespace: k8s 네임스페이스
"""
ensure_connected()
col = Collection(settings.MILVUS_NORMAL_LOG_COLLECTION)
try:
valid_logs = [m for m in log_messages if m and len(m) > 10][:200]
embeddings = get_embeddings_batch(valid_logs)
now = datetime.utcnow().isoformat()
data = [
embeddings,
[m[:999] for m in valid_logs],
[container[:199]] * len(valid_logs),
[namespace[:199]] * len(valid_logs),
[now] * len(valid_logs),
]
col.insert(data)
col.flush()
return json.dumps({
"status": "베이스라인 저장 완료",
"container": container,
"stored_count": len(valid_logs),
}, ensure_ascii=False)
except Exception as e:
return f"베이스라인 저장 오류: {str(e)}"
SYSTEM_PROMPT = """당신은 Kubernetes 기반 Data Lakehouse 환경의 로그 분석 전문가입니다.
과거 인시던트 데이터베이스와 운영 Runbook을 참고해 정확한 원인 분석과 해결책을 제시합니다.
## 사용 가능한 도구
### 실시간 데이터 조회
- query_prometheus: PromQL로 메트릭 조회
- search_opensearch_logs: 키워드로 로그 검색
- get_error_pattern_summary: 컨테이너별 에러 패턴 집계
- correlate_metric_and_logs: 메트릭 spike와 로그 상관관계
### Milvus 벡터 검색 (의미 기반)
- search_similar_incidents: 현재 증상과 유사한 과거 인시던트 + 해결책 검색
- search_runbook: 상황에 맞는 Runbook 절차 검색
- detect_anomaly_logs: 정상 베이스라인 대비 이상 로그 탐지
- store_resolved_incident: 해결된 인시던트 저장 (학습)
## 권장 분석 순서
1. get_error_pattern_summary → 전체 에러 현황 파악
2. search_similar_incidents → 유사 과거 사례 확인 (빠른 해결 단서)
3. query_prometheus + search_opensearch_logs → 현재 상태 상세 확인
4. search_runbook → 대응 절차 확인
5. detect_anomaly_logs → 놓친 이상 로그 없는지 확인
6. 분석 완료 후 → store_resolved_incident로 저장 (사용자 확인 후)
## 리포트 형식
🔴 Critical / 🟡 Warning / 🟢 Normal
**[현재 증상]** 무슨 일이 일어나고 있는가
**[유사 사례]** 과거에 같은 패턴이 있었는가 (유사도 점수 포함)
**[Root Cause]** 기술적 근거와 함께 추정 원인
**[권장 조치]** Runbook 참조 절차 + 즉각 조치 커맨드
**[학습]** 이 인시던트를 저장할지 여부 제안
"""
from langchain_ollama import ChatOllama
from langchain_core.messages import SystemMessage
from langgraph.graph import StateGraph, END
from langgraph.prebuilt import ToolNode, tools_condition
from agent.state import AgentState
from agent.tools import (
query_prometheus, search_opensearch_logs,
get_error_pattern_summary, correlate_metric_and_logs,
)
from agent.tools_milvus import (
search_similar_incidents, search_runbook,
detect_anomaly_logs, store_resolved_incident,
store_normal_log_baseline,
)
from agent.prompts import SYSTEM_PROMPT
from config import settings
ALL_TOOLS = [
query_prometheus,
search_opensearch_logs,
get_error_pattern_summary,
correlate_metric_and_logs,
search_similar_incidents,
search_runbook,
detect_anomaly_logs,
store_resolved_incident,
store_normal_log_baseline,
]
def build_graph():
llm = ChatOllama(
base_url=settings.OLLAMA_BASE_URL,
model=settings.OLLAMA_MODEL,
temperature=0,
num_ctx=16384,
).bind_tools(ALL_TOOLS)
tool_node = ToolNode(ALL_TOOLS)
def agent_node(state: AgentState):
messages = [SystemMessage(content=SYSTEM_PROMPT)] + list(state["messages"])
response = llm.invoke(messages)
return {"messages": [response]}
graph = StateGraph(AgentState)
graph.add_node("agent", agent_node)
graph.add_node("tools", tool_node)
graph.set_entry_point("agent")
graph.add_conditional_edges("agent", tools_condition, {"tools": "tools", END: END})
graph.add_edge("tools", "agent")
return graph.compile()
await analyze("MinIO에서 storage backend unavailable 에러가 터지고 있어. 과거에 이런 적 있었어?")
await analyze("MinIO 드라이브 장애 대응 절차 알려줘. 지금 disk latency가 급등하고 있어.")
await analyze("data-lakehouse 네임스페이스에서 지난 30분간 이상한 로그가 있는지 탐지해줘")
await analyze("""
오늘 14:30부터 Trino 쿼리가 실패하고 있어.
원인 분석하고, 해결되면 인시던트로 기록해줘.
""")
Runbook 문서 예시 구조
data/runbooks/
├── minio/
│ ├── disk-failure-recovery.md
│ ├── erasure-set-troubleshooting.md
│ └── performance-tuning.md
├── k8s/
│ ├── oom-killed-investigation.md
│ └── pod-crashloopbackoff.md
├── cilium/
│ └── network-policy-debugging.md
└── spark/
└── executor-failure-analysis.md
전체 데이터 흐름 요약
[초기 설정]
python scripts/init_milvus.py # 컬렉션 생성
python scripts/ingest_runbook.py # Runbook 적재
→ (평온한 시간대에 Agent로 베이스라인 로그 수집)
[운영 중]
에러 발생 → Agent 분석 요청
↓
OpenSearch 실시간 로그 + Prometheus 메트릭 조회
↓ ↓
Milvus 유사 인시던트 검색 Milvus Runbook RAG
↓ ↓
통합 리포트 생성
↓
해결 후 store_resolved_incident → 다음 번엔 더 빠른 해결