from pydantic_settings import BaseSettings
class Settings(BaseSettings):
OLLAMA_BASE_URL: str = "http://ollama-service:11434"
OLLAMA_MODEL: str = "llama3.1:8b"
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
class Config:
env_file = ".env"
settings = Settings()
from typing import Annotated, TypedDict, Sequence
from langchain_core.messages import BaseMessage
import operator
class AgentState(TypedDict):
messages: Annotated[Sequence[BaseMessage], operator.add]
tool_results: Annotated[list[dict], operator.add]
final_report: str
time_range_minutes: int
import json
import requests
from datetime import datetime, timedelta
from typing import Optional
from langchain_core.tools import tool
from opensearchpy import OpenSearch
from config import settings
def get_opensearch_client() -> OpenSearch:
return OpenSearch(
hosts=[{"host": settings.OPENSEARCH_HOST, "port": settings.OPENSEARCH_PORT}],
http_auth=(settings.OPENSEARCH_USER, settings.OPENSEARCH_PASSWORD),
use_ssl=settings.OPENSEARCH_USE_SSL,
verify_certs=False,
ssl_show_warn=False,
)
@tool
def query_prometheus(promql: str, time_range_minutes: int = 30) -> str:
"""
Prometheus에서 PromQL 쿼리를 실행합니다.
에러율, CPU/메모리 급등, latency spike 등 메트릭 이상 탐지에 사용합니다.
예시 promql:
- 'rate(http_requests_total{status=~"5.."}[5m])' : HTTP 5xx 에러율
- 'node_memory_MemAvailable_bytes / node_memory_MemTotal_bytes' : 메모리 여유율
- 'histogram_quantile(0.99, rate(http_request_duration_seconds_bucket[5m]))' : P99 latency
"""
end_time = datetime.utcnow()
start_time = end_time - timedelta(minutes=time_range_minutes)
try:
response = requests.get(
f"{settings.PROMETHEUS_URL}/api/v1/query_range",
params={
"query": promql,
"start": start_time.isoformat() + "Z",
"end": end_time.isoformat() + "Z",
"step": f"{max(15, time_range_minutes * 2)}s",
},
timeout=15,
)
response.raise_for_status()
data = response.json()
if data["status"] != "success":
return f"Prometheus 쿼리 실패: {data.get('error', 'unknown')}"
results = data["data"]["result"]
if not results:
return f"쿼리 결과 없음 (PromQL: {promql})"
summary = []
for series in results[:5]:
labels = series["metric"]
values = series["values"]
if not values:
continue
vals = [float(v[1]) for v in values if v[1] != "NaN"]
if not vals:
continue
summary.append({
"labels": labels,
"min": round(min(vals), 4),
"max": round(max(vals), 4),
"avg": round(sum(vals) / len(vals), 4),
"latest": round(vals[-1], 4),
"spike_detected": max(vals) > (sum(vals) / len(vals)) * 2,
})
return json.dumps({"query": promql, "time_range_min": time_range_minutes, "series": summary}, ensure_ascii=False)
except Exception as e:
return f"Prometheus 조회 오류: {str(e)}"
@tool
def search_opensearch_logs(
keyword: str,
log_level: Optional[str] = None,
namespace: Optional[str] = None,
time_range_minutes: int = 30,
size: int = 20,
) -> str:
"""
OpenSearch에서 로그를 검색합니다.
에러 메시지, 스택트레이스, 특정 서비스 로그 검색에 사용합니다.
Args:
keyword: 검색할 키워드 (예: "OOMKilled", "connection refused", "timeout")
log_level: 로그 레벨 필터 (ERROR, WARN, INFO 등)
namespace: k8s 네임스페이스 필터
time_range_minutes: 조회 시간 범위 (분)
size: 반환할 최대 로그 수
"""
client = get_opensearch_client()
end_time = datetime.utcnow()
start_time = end_time - timedelta(minutes=time_range_minutes)
must_clauses = [
{"range": {"@timestamp": {"gte": start_time.isoformat(), "lte": end_time.isoformat()}}},
{
"multi_match": {
"query": keyword,
"fields": ["message", "log", "log.message", "msg"],
"type": "best_fields",
}
},
]
if log_level:
must_clauses.append({"match": {"level": log_level}})
if namespace:
must_clauses.append({"match": {"kubernetes.namespace_name": namespace}})
try:
response = client.search(
index=settings.OPENSEARCH_INDEX_PATTERN,
body={
"query": {"bool": {"must": must_clauses}},
"sort": [{"@timestamp": {"order": "desc"}}],
"size": size,
"_source": [
"@timestamp", "message", "log", "level",
"kubernetes.pod_name", "kubernetes.namespace_name",
"kubernetes.container_name",
],
"aggs": {
"error_over_time": {
"date_histogram": {
"field": "@timestamp",
"fixed_interval": f"{max(1, time_range_minutes // 10)}m",
}
},
"top_pods": {
"terms": {"field": "kubernetes.pod_name.keyword", "size": 5}
},
},
},
)
hits = response["hits"]["hits"]
total = response["hits"]["total"]["value"]
aggs = response.get("aggregations", {})
logs = []
for hit in hits:
src = hit["_source"]
logs.append({
"timestamp": src.get("@timestamp"),
"pod": src.get("kubernetes", {}).get("pod_name", "unknown"),
"namespace": src.get("kubernetes", {}).get("namespace_name", "unknown"),
"container": src.get("kubernetes", {}).get("container_name", "unknown"),
"message": (src.get("message") or src.get("log") or "")[:300],
})
top_pods = [b["key"] for b in aggs.get("top_pods", {}).get("buckets", [])]
return json.dumps({
"total_hits": total,
"keyword": keyword,
"top_pods_affected": top_pods,
"sample_logs": logs[:10],
}, ensure_ascii=False)
except Exception as e:
return f"OpenSearch 조회 오류: {str(e)}"
@tool
def get_error_pattern_summary(
namespace: Optional[str] = None,
time_range_minutes: int = 60,
) -> str:
"""
최근 에러 로그의 패턴을 집계하여 반환합니다.
어떤 에러가 얼마나 발생하는지 overview를 파악할 때 사용합니다.
"""
client = get_opensearch_client()
end_time = datetime.utcnow()
start_time = end_time - timedelta(minutes=time_range_minutes)
filter_clauses = [
{"range": {"@timestamp": {"gte": start_time.isoformat(), "lte": end_time.isoformat()}}},
{"terms": {"level.keyword": ["ERROR", "FATAL", "error", "fatal"]}},
]
if namespace:
filter_clauses.append({"match": {"kubernetes.namespace_name": namespace}})
try:
response = client.search(
index=settings.OPENSEARCH_INDEX_PATTERN,
body={
"query": {"bool": {"filter": filter_clauses}},
"size": 0,
"aggs": {
"by_container": {
"terms": {"field": "kubernetes.container_name.keyword", "size": 10},
"aggs": {
"error_count": {"value_count": {"field": "@timestamp"}},
"sample_messages": {
"top_hits": {
"size": 3,
"_source": ["message"],
"sort": [{"@timestamp": {"order": "desc"}}],
}
},
},
}
},
},
)
buckets = response["aggregations"]["by_container"]["buckets"]
result = []
for b in buckets:
msgs = [
h["_source"].get("message", "")[:200]
for h in b["sample_messages"]["hits"]["hits"]
]
result.append({
"container": b["key"],
"error_count": b["doc_count"],
"sample_messages": msgs,
})
return json.dumps({
"time_range_minutes": time_range_minutes,
"error_summary_by_container": result,
}, ensure_ascii=False)
except Exception as e:
return f"에러 패턴 집계 오류: {str(e)}"
@tool
def correlate_metric_and_logs(
promql: str,
log_keyword: str,
time_range_minutes: int = 30,
) -> str:
"""
메트릭 spike 시점과 로그 에러 발생 시점의 상관관계를 분석합니다.
예를 들어 CPU spike와 OOM 로그가 동시에 발생했는지 확인합니다.
"""
end_time = datetime.utcnow()
start_time = end_time - timedelta(minutes=time_range_minutes)
try:
prom_resp = requests.get(
f"{settings.PROMETHEUS_URL}/api/v1/query_range",
params={
"query": promql,
"start": start_time.isoformat() + "Z",
"end": end_time.isoformat() + "Z",
"step": "60s",
},
timeout=15,
)
prom_data = prom_resp.json()
results = prom_data.get("data", {}).get("result", [])
spike_windows = []
for series in results[:3]:
values = series["values"]
if not values:
continue
nums = [(v[0], float(v[1])) for v in values if v[1] != "NaN"]
if not nums:
continue
avg = sum(v for _, v in nums) / len(nums)
spikes = [(ts, val) for ts, val in nums if val > avg * 1.5]
for ts, val in spikes[:3]:
dt = datetime.utcfromtimestamp(ts)
spike_windows.append({
"timestamp": dt.isoformat(),
"value": round(val, 4),
"avg": round(avg, 4),
})
client = get_opensearch_client()
correlated_logs = []
for spike in spike_windows[:2]:
spike_time = datetime.fromisoformat(spike["timestamp"])
window_start = (spike_time - timedelta(minutes=2)).isoformat()
window_end = (spike_time + timedelta(minutes=2)).isoformat()
log_resp = client.search(
index=settings.OPENSEARCH_INDEX_PATTERN,
body={
"query": {
"bool": {
"must": [
{"range": {"@timestamp": {"gte": window_start, "lte": window_end}}},
{"multi_match": {"query": log_keyword, "fields": ["message", "log"]}},
]
}
},
"size": 5,
"_source": ["@timestamp", "message", "kubernetes.pod_name"],
},
)
correlated_logs.append({
"spike": spike,
"related_logs_count": log_resp["hits"]["total"]["value"],
"logs": [h["_source"] for h in log_resp["hits"]["hits"]],
})
return json.dumps({
"metric_query": promql,
"log_keyword": log_keyword,
"spike_count": len(spike_windows),
"correlations": correlated_logs,
}, ensure_ascii=False)
except Exception as e:
return f"상관관계 분석 오류: {str(e)}"
SYSTEM_PROMPT = """당신은 Kubernetes 기반 Data Lakehouse 환경의 로그 분석 전문가입니다.
## 사용 가능한 도구
- query_prometheus: PromQL로 메트릭 조회 (에러율, CPU/메모리, latency 등)
- search_opensearch_logs: 키워드로 로그 검색
- get_error_pattern_summary: 컨테이너별 에러 패턴 집계
- correlate_metric_and_logs: 메트릭 spike와 로그 에러 시점 상관관계 분석
## 분석 절차
1. 먼저 get_error_pattern_summary로 전체 에러 현황 파악
2. 이상 징후 발견 시 query_prometheus로 해당 시점 메트릭 확인
3. search_opensearch_logs로 에러 상세 내용 탐색
4. 필요시 correlate_metric_and_logs로 메트릭-로그 상관관계 분석
5. 원인과 해결 방향을 한국어로 명확하게 리포트
## 리포트 형식
- 🔴 Critical / 🟡 Warning / 🟢 Normal 로 심각도 표시
- 발생 시간, 영향 받은 서비스, 추정 원인, 권장 조치 포함
- 기술적 근거 (어떤 메트릭/로그에서 판단했는지) 명시
## 주의
- 확실하지 않은 내용은 추정임을 명시
- 데이터 없으면 데이터 없다고 솔직하게 답변
"""
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.prompts import SYSTEM_PROMPT
from config import settings
TOOLS = [
query_prometheus,
search_opensearch_logs,
get_error_pattern_summary,
correlate_metric_and_logs,
]
def build_graph() -> StateGraph:
llm = ChatOllama(
base_url=settings.OLLAMA_BASE_URL,
model=settings.OLLAMA_MODEL,
temperature=0,
num_ctx=8192,
).bind_tools(TOOLS)
tool_node = ToolNode(TOOLS)
def agent_node(state: AgentState):
messages = [SystemMessage(content=SYSTEM_PROMPT)] + list(state["messages"])
response = llm.invoke(messages)
return {"messages": [response]}
def should_continue(state: AgentState):
return tools_condition(state)
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", should_continue, {"tools": "tools", END: END})
graph.add_edge("tools", "agent")
return graph.compile()
_graph = None
def get_graph():
global _graph
if _graph is None:
_graph = build_graph()
return _graph
import asyncio
from langchain_core.messages import HumanMessage
from agent.graph import get_graph
from agent.state import AgentState
async def analyze(query: str, time_range_minutes: int = 30) -> str:
graph = get_graph()
initial_state: AgentState = {
"messages": [HumanMessage(content=query)],
"tool_results": [],
"final_report": "",
"time_range_minutes": time_range_minutes,
}
print(f"\n{'='*60}")
print(f"🔍 분석 시작: {query}")
print(f"{'='*60}\n")
final_state = None
async for event in graph.astream(initial_state, stream_mode="updates"):
for node_name, update in event.items():
if node_name == "agent":
msgs = update.get("messages", [])
for msg in msgs:
if hasattr(msg, "content") and msg.content:
print(f"[Agent] {msg.content[:200]}...")
if hasattr(msg, "tool_calls") and msg.tool_calls:
for tc in msg.tool_calls:
print(f" → 도구 호출: {tc['name']}({list(tc['args'].keys())})")
elif node_name == "tools":
msgs = update.get("messages", [])
for msg in msgs:
print(f" ✓ 도구 결과 수신 ({len(str(msg.content))} chars)")
final_state = event
if final_state:
last_node_output = list(final_state.values())[-1]
messages = last_node_output.get("messages", [])
if messages:
return messages[-1].content
return "분석 실패"
if __name__ == "__main__":
queries = [
"지난 1시간 동안 에러가 많이 발생한 서비스가 있는지 분석해줘",
"MinIO에서 connection timeout 에러가 발생하고 있어. 원인을 분석해줘",
"최근 30분간 CPU spike와 관련된 에러 로그를 상관분석 해줘",
"data-lakehouse 네임스페이스에서 OOM 관련 이슈가 있는지 확인해줘",
]
result = asyncio.run(analyze(queries[0], time_range_minutes=60))
print(f"\n{'='*60}")
print("📋 최종 분석 리포트")
print(f"{'='*60}")
print(result)
실행 예시 출력
============================================================
🔍 분석 시작: 지난 1시간 동안 에러가 많이 발생한 서비스가 있는지 분석해줘
============================================================
[Agent] 에러 패턴 집계를 먼저 확인하겠습니다...
→ 도구 호출: get_error_pattern_summary(['time_range_minutes'])
✓ 도구 결과 수신 (1842 chars)
[Agent] minio 컨테이너에서 에러가 집중되고 있습니다. 메트릭을 확인합니다...
→ 도구 호출: query_prometheus(['promql', 'time_range_minutes'])
✓ 도구 결과 수신 (923 chars)
→ 도구 호출: search_opensearch_logs(['keyword', 'namespace'])
✓ 도구 결과 수신 (2341 chars)
============================================================
📋 최종 분석 리포트
============================================================
🔴 Critical: minio 스토리지 연결 오류 집중 발생
**발생 현황**
- 컨테이너: minio (data-lakehouse 네임스페이스)
- 에러 수: 최근 1시간 347건
- 주요 메시지: "dial tcp: connection refused", "storage backend unavailable"
**메트릭 연계 분석**
- 14:23 ~ 14:31 구간에 CPU 사용률 92% spike 감지
- 동 시간대 에러 로그 급증과 일치
**추정 원인**
MinIO erasure set 중 일부 드라이브 응답 지연으로 인한 quorum 실패로 추정
**권장 조치**
1. `kubectl logs -n data-lakehouse minio-0` 로 상세 확인
2. `mc admin info` 로 드라이브 상태 점검
3. 디스크 I/O latency 확인: `node_disk_io_time_seconds_total`
langgraph>=0.2.0
langchain>=0.3.0
langchain-ollama>=0.2.0
langchain-core>=0.3.0
opensearch-py>=2.4.0
requests>=2.31.0
pydantic>=2.0.0
python-dotenv>=1.0.0