26F26e

Young-Kyoo Kim·2026년 2월 25일
###config.py
from pydantic_settings import BaseSettings

class Settings(BaseSettings):
    # Ollama
    OLLAMA_BASE_URL: str = "http://ollama-service:11434"
    OLLAMA_MODEL: str = "llama3.1:8b"  # 또는 qwen2.5, mistral 등

    # Prometheus
    PROMETHEUS_URL: str = "http://prometheus:9090"

    # OpenSearch
    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()

###agent/state.py
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
    
    
###agent/tools.py
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

# ── OpenSearch 클라이언트 ──────────────────────────────
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 1: Prometheus 메트릭 조회 ────────────────────
@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:
        # range query
        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})"

        # 결과 요약 (LLM이 소화하기 쉽게)
        summary = []
        for series in results[:5]:  # 최대 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,  # avg의 2배 이상이면 spike
            })

        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 2: OpenSearch 로그 검색 ──────────────────────
@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],  # 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],  # LLM에 10개만 전달
        }, ensure_ascii=False)

    except Exception as e:
        return f"OpenSearch 조회 오류: {str(e)}"


# ── Tool 3: 에러 패턴 집계 ───────────────────────────
@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 4: 메트릭-로그 시간 상관관계 ────────────────
@tool
def correlate_metric_and_logs(
    promql: str,
    log_keyword: str,
    time_range_minutes: int = 30,
) -> str:
    """
    메트릭 spike 시점과 로그 에러 발생 시점의 상관관계를 분석합니다.
    예를 들어 CPU spike와 OOM 로그가 동시에 발생했는지 확인합니다.
    """
    # Prometheus에서 spike 시점 찾기
    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),
                })

        # spike 시점 전후 로그 검색
        client = get_opensearch_client()
        correlated_logs = []

        for spike in spike_windows[:2]:  # 최대 2개 spike 분석
            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)}"
        
        
###agent/prompts.py
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 로 심각도 표시
- 발생 시간, 영향 받은 서비스, 추정 원인, 권장 조치 포함
- 기술적 근거 (어떤 메트릭/로그에서 판단했는지) 명시

## 주의
- 확실하지 않은 내용은 추정임을 명시
- 데이터 없으면 데이터 없다고 솔직하게 답변
"""


###agent/graph.py
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:
    # Ollama LLM (airgap 로컬)
    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
    
    
###main.py
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. 전체 현황 파악
        "지난 1시간 동안 에러가 많이 발생한 서비스가 있는지 분석해줘",
        
        # 2. 특정 증상 조사
        "MinIO에서 connection timeout 에러가 발생하고 있어. 원인을 분석해줘",
        
        # 3. 메트릭-로그 연계
        "최근 30분간 CPU spike와 관련된 에러 로그를 상관분석 해줘",
        
        # 4. 특정 네임스페이스
        "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

0개의 댓글