文章

LLM 辅助 RCA 与日志分析

LLM 辅助 RCA 与日志分析

核心目标:利用大语言模型(LLM)的自然语言理解、模式识别和推理能力,将根因分析和日志解读从”人海战术”升级为”AI辅助+人机协同”,大幅缩短 MTTI(平均识别时间)。

概述

传统 RCA 和日志分析面临三大痛点:日志量爆炸(日均TB级)、上下文碎片化(分散在多系统)、经验难以传承(资深工程师的”第六感”难以文档化)。LLM 的引入为这三个痛点提供了新解法:自然语言交互降低使用门槛、多模态上下文整合打破信息孤岛、推理链路可追溯让经验沉淀可复用。

flowchart TB
    subgraph 数据输入层
        A[结构化日志<br/>ELK/Loki] 
        B[指标数据<br/>Prometheus]
        C[事件流<br/>K8s Events/告警]
        D[链路追踪<br/>Jaeger/Tempo]
        E[变更记录<br/>GitOps/CI]
    end

    subgraph LLM处理层
        F[日志聚类<br/>+模式提取]
        G[上下文组装<br/>+Prompt构建]
        H[根因推理<br/>+假设验证]
        I[自然语言交互<br/>Q&A]
    end

    subgraph 输出层
        J[根因报告]
        K[修复建议]
        L[Action Items]
        M[知识沉淀<br/>更新RAG库]
    end

    A --> F
    B --> G
    C --> G
    D --> G
    E --> G
    F --> G
    G --> H
    H --> I
    H --> J
    H --> K
    K --> L
    L --> M
    I --> M

一、LLM 辅助日志分析

1. 日志聚类与异常模式提取

传统方法(如 Drain 算法)能做日志模板提取,但缺乏语义理解。LLM 能在聚类基础上做语义级别的模式识别和异常解释。

"""LLM 辅助日志聚类与异常模式提取"""

from dataclasses import dataclass, field
from enum import Enum
from typing import Optional
import re
import hashlib


class LogLevel(Enum):
    DEBUG = "DEBUG"
    INFO = "INFO"
    WARN = "WARN"
    ERROR = "ERROR"
    FATAL = "FATAL"


@dataclass
class LogEntry:
    timestamp: str
    level: LogLevel
    service: str
    message: str
    raw: str = ""
    cluster_id: str = ""
    is_anomaly: bool = False


@dataclass
class LogCluster:
    cluster_id: str
    template: str            # 日志模板(变量替换后)
    sample: str              # 原始示例
    count: int = 0           # 出现次数
    level: LogLevel = LogLevel.INFO
    first_seen: str = ""
    last_seen: str = ""
    entries: list[LogEntry] = field(default_factory=list)
    anomaly_score: float = 0.0


class LogClusterer:
    """基于模板提取的日志聚类器"""

    # 变量替换正则
    VARIABLE_PATTERNS = [
        (re.compile(r'\b\d{1,3}\.\d{1,3}\.\d{1,3}\.\d{1,3}\b'), '<IP>'),
        (re.compile(r'\b[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}\b',
                    re.I), '<UUID>'),
        (re.compile(r'\b\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}'), '<TIMESTAMP>'),
        (re.compile(r'\b\d+\.?\d*(ms|s|us)?\b'), '<NUM>'),
        (re.compile(r'\b0x[0-9a-f]+\b', re.I), '<HEX>'),
        (re.compile(r'/[\w/.-]+'), '<PATH>'),
    ]

    def cluster(self, logs: list[LogEntry]) -> list[LogCluster]:
        """对日志进行聚类"""
        clusters: dict[str, LogCluster] = {}

        for log in logs:
            template = self._extract_template(log.message)
            cluster_id = hashlib.md5(template.encode()).hexdigest()[:8]

            if cluster_id not in clusters:
                clusters[cluster_id] = LogCluster(
                    cluster_id=cluster_id,
                    template=template,
                    sample=log.message,
                    level=log.level,
                    first_seen=log.timestamp,
                    last_seen=log.timestamp,
                )

            cluster = clusters[cluster_id]
            cluster.count += 1
            cluster.last_seen = log.timestamp
            log.cluster_id = cluster_id
            cluster.entries.append(log)

        return list(clusters.values())

    def _extract_template(self, message: str) -> str:
        """提取日志模板(将变量替换为占位符)"""
        template = message
        for pattern, replacement in self.VARIABLE_PATTERNS:
            template = pattern.sub(replacement, template)
        return template


class LLMLogAnalyzer:
    """LLM 日志异常分析器"""

    def __init__(self, clusterer: LogClusterer):
        self.clusterer = clusterer

    def analyze(self, logs: list[LogEntry],
                historical_baseline: Optional[dict] = None) -> dict:
        """分析日志,检测异常并生成LLM可读的上下文"""
        clusters = self.clusterer.cluster(logs)

        # 统计异常模式
        anomalies = []
        for cluster in clusters:
            # 频率异常:该模板出现次数显著偏离基线
            baseline_count = (historical_baseline or {}).get(
                cluster.template, 0
            )
            if baseline_count > 0:
                ratio = cluster.count / baseline_count
                if ratio > 3 or ratio < 0.3:
                    cluster.anomaly_score = min(ratio, 1.0)
                    cluster.is_anomaly = True
                    anomalies.append(cluster)

            # 新模板异常:从未见过的模板
            if baseline_count == 0 and cluster.level in (
                LogLevel.ERROR, LogLevel.FATAL
            ):
                cluster.anomaly_score = 0.8
                cluster.is_anomaly = True
                anomalies.append(cluster)

        # 构建 LLM 分析 prompt
        prompt = self._build_analysis_prompt(anomalies, clusters)
        return {
            "total_logs": len(logs),
            "total_clusters": len(clusters),
            "anomaly_clusters": len(anomalies),
            "prompt": prompt,
            "anomalies": [
                {
                    "template": a.template,
                    "sample": a.sample,
                    "count": a.count,
                    "level": a.level.value,
                    "score": a.anomaly_score,
                }
                for a in anomalies
            ],
        }

    def _build_analysis_prompt(self, anomalies: list[LogCluster],
                                all_clusters: list[LogCluster]) -> str:
        """构建 LLM 分析 prompt"""
        prompt_parts = [
            "你是一位资深 SRE 工程师,请分析以下日志聚类结果,"
            "识别异常模式并推断可能的根因。\n\n",
            f"## 日志概览\n",
            f"- 总日志数: {sum(c.count for c in all_clusters)}\n",
            f"- 聚类数: {len(all_clusters)}\n",
            f"- 异常聚类数: {len(anomalies)}\n\n",
            "## 异常日志聚类\n",
        ]

        for i, anomaly in enumerate(anomalies[:20], 1):
            prompt_parts.extend([
                f"### 异常 {i}\n",
                f"- 模板: {anomaly.template}\n",
                f"- 示例: {anomaly.sample}\n",
                f"- 出现次数: {anomaly.count}\n",
                f"- 级别: {anomaly.level.value}\n",
                f"- 异常评分: {anomaly.anomaly_score:.2f}\n\n",
            ])

        prompt_parts.extend([
            "## 请输出\n",
            "1. 异常模式总结(2-3句话)\n",
            "2. 可能的根因推断(列出2-3个假设)\n",
            "3. 建议的排查步骤\n",
            "4. 需要进一步查看的指标/日志\n",
        ])

        return "".join(prompt_parts)


# 使用示例
if __name__ == "__main__":
    logs = [
        LogEntry("10:30:01", LogLevel.ERROR, "api", 
                 "Connection refused to mysql-primary:3306 (10.0.1.5)"),
        LogEntry("10:30:02", LogLevel.ERROR, "api",
                 "Connection refused to mysql-primary:3306 (10.0.1.5)"),
        LogEntry("10:30:03", LogLevel.WARN, "api",
                 "Retrying connection attempt 1/3"),
        LogEntry("10:30:05", LogLevel.FATAL, "worker",
                 "Failed to acquire database connection after 5000ms timeout"),
        LogEntry("10:30:06", LogLevel.ERROR, "gateway",
                 "Upstream timeout: 504 Gateway Timeout from 10.0.1.10:8080"),
    ]

    analyzer = LLMLogAnalyzer(LogClusterer())
    result = analyzer.analyze(logs)
    print(f"总日志: {result['total_logs']}")
    print(f"异常聚类: {result['anomaly_clusters']}")
    print(f"\n--- LLM Prompt ---\n{result['prompt'][:500]}...")

2. 多源日志关联分析

"""多源日志关联分析 —— 构建故障时间线"""

from dataclasses import dataclass, field
from typing import Optional
from datetime import datetime
import re


@dataclass
class CorrelatedEvent:
    timestamp: datetime
    source: str          # log / metric / event / trace / change
    service: str
    message: str
    severity: str       # info / warn / error / critical
    correlation_id: str = ""


class LogCorrelator:
    """跨源日志关联器"""

    def __init__(self):
        self.events: list[CorrelatedEvent] = []

    def add_logs(self, logs: list[dict], source: str = "log"):
        for log in logs:
            self.events.append(CorrelatedEvent(
                timestamp=self._parse_time(log.get("timestamp")),
                source=source,
                service=log.get("service", "unknown"),
                message=log.get("message", ""),
                severity=log.get("level", "info").lower(),
            ))

    def add_k8s_events(self, events: list[dict]):
        for ev in events:
            self.events.append(CorrelatedEvent(
                timestamp=self._parse_time(ev.get("lastTimestamp")),
                source="event",
                service=ev.get("involvedObject", {}).get("name", ""),
                message=ev.get("message", ""),
                severity=ev.get("type", "Normal").lower(),
            ))

    def add_changes(self, changes: list[dict]):
        for ch in changes:
            self.events.append(CorrelatedEvent(
                timestamp=self._parse_time(ch.get("timestamp")),
                source="change",
                service=ch.get("service", ""),
                message=f"变更: {ch.get('action')} by {ch.get('author')}",
                severity="info",
            ))

    def build_timeline(self) -> list[CorrelatedEvent]:
        """按时间排序构建时间线"""
        return sorted(self.events, key=lambda e: e.timestamp)

    def build_llm_context(self, window_minutes: int = 30) -> str:
        """构建 LLM 可读的上下文"""
        timeline = self.build_timeline()

        prompt_parts = [
            "以下是故障时间窗口内(最近{window}分钟)的多源事件时间线,\n"
            "请分析事件之间的因果关系并推断根因。\n\n".format(
                window=window_minutes
            ),
            "## 事件时间线\n\n",
            "| 时间 | 来源 | 服务 | 级别 | 事件 |\n",
            "|------|------|------|------|------|\n",
        ]

        for event in timeline[-50:]:  # 最近50条
            prompt_parts.append(
                f"| {event.timestamp:%H:%M:%S} | {event.source} | "
                f"{event.service} | {event.severity} | "
                f"{event.message[:80]} |\n"
            )

        prompt_parts.extend([
            "\n## 请分析\n",
            "1. 事件因果关系链(哪个事件最可能是根因)\n",
            "2. 故障传播路径\n",
            "3. 与变更的关联性\n",
            "4. 建议的验证步骤\n",
        ])

        return "".join(prompt_parts)

    @staticmethod
    def _parse_time(ts) -> datetime:
        if isinstance(ts, datetime):
            return ts
        if isinstance(ts, str):
            try:
                return datetime.fromisoformat(ts.replace("Z", "+00:00"))
            except ValueError:
                pass
        return datetime.now()

二、LLM 辅助根因分析

1. 5-Whys 自动推理

"""LLM 驱动的 5-Whys 自动推理链"""

from dataclasses import dataclass, field
from typing import Optional
import json


@dataclass
class WhyStep:
    step: int
    question: str
    answer: str
    evidence: str           # 支撑证据
    confidence: float       # 置信度 0-1
    needs_verification: bool = False


@dataclass
class RCAResult:
    incident: str
    why_chain: list[WhyStep] = field(default_factory=list)
    root_cause: str = ""
    contributing_factors: list[str] = field(default_factory=list)
    action_items: list[str] = field(default_factory=list)
    confidence: float = 0.0


class LLMRCAEngine:
    """LLM 驱动的根因分析引擎"""

    SYSTEM_PROMPT = """你是一位资深 SRE 根因分析专家。你的任务是通过 5-Whys 方法
    逐步深入分析故障根因。每一层 Why 必须基于提供的证据,不能臆测。
    如果证据不足,明确标注"需要验证"并列出验证方法。

    输出 JSON 格式:
    {
      "why_chain": [
        {
          "step": 1,
          "question": "为什么...?",
          "answer": "因为...",
          "evidence": "具体证据(日志/指标/事件)",
          "confidence": 0.0-1.0,
          "needs_verification": true/false
        }
      ],
      "root_cause": "根本原因总结",
      "contributing_factors": ["因素1", "因素2"],
      "action_items": ["改进项1", "改进项2"],
      "confidence": 0.0-1.0
    }
    """

    def analyze(self, incident_context: dict) -> str:
        """构建 RCA 分析 prompt"""
        prompt = self._build_rca_prompt(incident_context)
        return prompt  # 传给 LLM 执行

    def _build_rca_prompt(self, ctx: dict) -> str:
        """构建 RCA 分析的完整 prompt"""
        prompt_parts = [self.SYSTEM_PROMPT, "\n---\n\n"]

        # 1. 故障概述
        prompt_parts.append("## 故障概述\n")
        prompt_parts.append(f"- 故障描述: {ctx.get('description', '')}\n")
        prompt_parts.append(f"- 影响范围: {ctx.get('impact', '')}\n")
        prompt_parts.append(f"- 开始时间: {ctx.get('start_time', '')}\n")
        prompt_parts.append(f"- 持续时间: {ctx.get('duration', '')}\n")
        prompt_parts.append(
            f"- 严重级别: {ctx.get('severity', '')}\n\n"
        )

        # 2. 关键指标变化
        if ctx.get("metrics"):
            prompt_parts.append("## 关键指标变化\n\n")
            for metric in ctx["metrics"]:
                prompt_parts.append(
                    f"- {metric['name']}: "
                    f"正常值 {metric['baseline']} → "
                    f"异常值 {metric['current']} "
                    f"({metric['change']})\n"
                )
            prompt_parts.append("\n")

        # 3. 事件时间线
        if ctx.get("timeline"):
            prompt_parts.append("## 事件时间线\n\n")
            for event in ctx["timeline"]:
                prompt_parts.append(
                    f"- [{event['time']}] [{event['source']}] "
                    f"{event['message']}\n"
                )
            prompt_parts.append("\n")

        # 4. 变更记录
        if ctx.get("changes"):
            prompt_parts.append("## 近期变更\n\n")
            for change in ctx["changes"]:
                prompt_parts.append(
                    f"- [{change['time']}] {change['action']} "
                    f"(by {change['author']}): {change['detail']}\n"
                )
            prompt_parts.append("\n")

        # 5. 异常日志
        if ctx.get("anomaly_logs"):
            prompt_parts.append("## 异常日志聚类\n\n")
            for log in ctx["anomaly_logs"][:10]:
                prompt_parts.append(
                    f"- [{log['level']}] ({log['count']}次) "
                    f"{log['template']}\n"
                )
            prompt_parts.append("\n")

        # 6. 拓扑信息
        if ctx.get("topology"):
            prompt_parts.append("## 服务拓扑\n\n")
            prompt_parts.append(f"```\n{ctx['topology']}\n```\n\n")

        prompt_parts.append("## 请执行 5-Whys 根因分析\n")
        return "".join(prompt_parts)


# ====== 使用示例 ======
if __name__ == "__main__":
    engine = LLMRCAEngine()

    incident_context = {
        "description": "支付服务 502 错误率从 0.1% 飙升至 45%",
        "impact": "支付功能不可用,影响约 2 万用户",
        "start_time": "2026-08-03 10:30:00",
        "duration": "8分钟",
        "severity": "SEV1",
        "metrics": [
            {"name": "HTTP 502率", "baseline": "0.1%",
             "current": "45%", "change": "+450x"},
            {"name": "P99延迟", "baseline": "200ms",
             "current": "5000ms", "change": "+25x"},
            {"name": "MySQL连接数", "baseline": "50",
             "current": "980", "change": "+19.6x"},
        ],
        "timeline": [
            {"time": "10:28:00", "source": "change",
             "message": "部署 payment-service v3.2.0"},
            {"time": "10:29:30", "source": "metric",
             "message": "MySQL连接数开始上升"},
            {"time": "10:30:00", "source": "alert",
             "message": "502错误率告警触发"},
            {"time": "10:30:05", "source": "log",
             "message": "Connection pool exhausted: max 1000"},
        ],
        "changes": [
            {"time": "10:28:00", "action": "deploy",
             "author": "zhangsan",
             "detail": "v3.2.0: 优化支付回调逻辑"},
        ],
        "anomaly_logs": [
            {"level": "ERROR", "count": 342,
             "template": "Connection pool exhausted: max <NUM>"},
            {"level": "ERROR", "count": 128,
             "template": "Failed to get connection from pool"},
            {"level": "WARN", "count": 56,
             "template": "Request timeout after <NUM>ms"},
        ],
    }

    prompt = engine.analyze(incident_context)
    print(prompt[:1000])
    # 实际调用 LLM 后解析 JSON 结果

2. 假设生成与验证

"""LLM 驱动的故障假设生成与验证"""

from dataclasses import dataclass, field
from enum import Enum
from typing import Optional


class HypothesisStatus(Enum):
    PROPOSED = "proposed"       # 已提出
    VERIFYING = "verifying"      # 验证中
    CONFIRMED = "confirmed"      # 已确认
    REFUTED = "refuted"          # 已排除
    INCONCLUSIVE = "inconclusive"  # 无法确认


@dataclass
class VerificationStep:
    description: str
    query: str              # 验证查询(PromQL/SQL/命令)
    expected_result: str    # 预期结果(确认/排除)
    actual_result: str = ""
    status: HypothesisStatus = HypothesisStatus.PROPOSED


@dataclass
class Hypothesis:
    id: str
    statement: str             # 假设陈述
    probability: float         # 初始概率
    evidence_for: list[str] = field(default_factory=list)
    evidence_against: list[str] = field(default_factory=list)
    verification_steps: list[VerificationStep] = field(default_factory=list)
    status: HypothesisStatus = HypothesisStatus.PROPOSED
    final_probability: float = 0.0


class HypothesisEngine:
    """假设生成与验证引擎"""

    HYPOTHESIS_PROMPT = """基于以下故障上下文,生成 3-5 个可能的根因假设。
    每个假设包含:
    1. 假设陈述
    2. 初始概率(0-1)
    3. 支持证据
    4. 反对证据
    5. 验证步骤(具体查询命令)

    故障上下文:
    {context}

    输出 JSON 格式:
    {{
      "hypotheses": [
        {{
          "id": "H1",
          "statement": "...",
          "probability": 0.0-1.0,
          "evidence_for": ["..."],
          "evidence_against": ["..."],
          "verification_steps": [
            {{
              "description": "...",
              "query": "具体查询命令",
              "expected_result": "如果假设成立应该看到的结果"
            }}
          ]
        }}
      ]
    }}
    """

    def generate_hypotheses(self, context: dict) -> str:
        """生成假设 prompt"""
        import json
        return self.HYPOTHESIS_PROMPT.format(
            context=json.dumps(context, ensure_ascii=False, indent=2)
        )

    def evaluate_hypothesis(self, hypothesis: Hypothesis,
                             verification_results: list[dict]) -> Hypothesis:
        """根据验证结果更新假设状态"""
        for i, result in enumerate(verification_results):
            if i < len(hypothesis.verification_steps):
                step = hypothesis.verification_steps[i]
                step.actual_result = result.get("output", "")
                if result.get("matches_expected"):
                    step.status = HypothesisStatus.CONFIRMED
                    hypothesis.evidence_for.append(
                        f"验证{i+1}: {step.description} → 确认"
                    )
                else:
                    step.status = HypothesisStatus.REFUTED
                    hypothesis.evidence_against.append(
                        f"验证{i+1}: {step.description} → 排除"
                    )

        # 更新假设状态
        confirmed = sum(
            1 for s in hypothesis.verification_steps
            if s.status == HypothesisStatus.CONFIRMED
        )
        refuted = sum(
            1 for s in hypothesis.verification_steps
            if s.status == HypothesisStatus.REFUTED
        )

        if refuted > 0 and confirmed == 0:
            hypothesis.status = HypothesisStatus.REFUTED
            hypothesis.final_probability = 0.1
        elif confirmed > 0 and refuted == 0:
            hypothesis.status = HypothesisStatus.CONFIRMED
            hypothesis.final_probability = 0.9
        else:
            hypothesis.status = HypothesisStatus.INCONCLUSIVE
            hypothesis.final_probability = 0.5

        return hypothesis

三、自然语言运维查询

1. NL-to-PromQL

将自然语言查询转化为 PromQL,降低可观测性数据的使用门槛。

"""自然语言 → PromQL 转换器"""

from dataclasses import dataclass
from typing import Optional


@dataclass
class PromQLQuery:
    natural_language: str
    promql: str
    explanation: str
    confidence: float


class NLToPromQLConverter:
    """自然语言到 PromQL 的转换器"""

    SYSTEM_PROMPT = """你是一位 PromQL 专家。将用户的自然语言查询
    转换为有效的 PromQL 表达式。

    可用指标:
    - http_requests_total{code, method, service}  — HTTP请求计数
    - http_request_duration_seconds{quantile, service}  — HTTP延迟
    - container_cpu_usage_seconds_total{container, pod}  — CPU使用
    - container_memory_usage_bytes{container, pod}  — 内存使用
    - kube_pod_status_phase{phase, namespace}  — Pod状态
    - mysql_connections{state}  — MySQL连接数
    - kafka_consumergroup_lag{topic, consumergroup}  — Kafka消费延迟
    - node_cpu_seconds_total{mode}  — 节点CPU
    - node_memory_MemAvailable_bytes  — 节点可用内存

    规则:
    1. 只输出一个 PromQL 表达式
    2. 使用 rate() 处理 counter 类型
    3. 使用 histogram_quantile() 处理延迟分位数
    4. 标签过滤用 {} 语法

    输出JSON:
    {
      "promql": "表达式",
      "explanation": "解释",
      "confidence": 0.0-1.0
    }
    """

    EXAMPLES = [
        {
            "nl": "支付服务最近5分钟的P99延迟",
            "promql": "histogram_quantile(0.99, "
                      "rate(http_request_duration_seconds_bucket"
                      "{service=\"payment\"}[5m]))",
            "explanation": "使用histogram_quantile计算payment服务的P99延迟"
        },
        {
            "nl": "所有服务的错误率排名",
            "promql": "sum(rate(http_requests_total"
                      "{code=~\"5..\"}[5m])) by (service) / "
                      "sum(rate(http_requests_total[5m])) by (service) * 100",
            "explanation": "计算每个服务的5xx错误率百分比"
        },
        {
            "nl": "Kafka消费积压最严重的topic",
            "promql": "topk(10, sum(kafka_consumergroup_lag) by (topic))",
            "explanation": "按消费延迟排名取前10的topic"
        },
    ]

    def convert(self, natural_language: str) -> str:
        """构建转换 prompt"""
        examples_text = "\n".join(
            f"问: {ex['nl']}\n"
            f"答: {json.dumps({k: v for k, v in ex.items() if k != 'nl'})}\n"
            for ex in self.EXAMPLES
        )
        return f"""{self.SYSTEM_PROMPT}

## 示例
{examples_text}

## 用户查询
{natural_language}
"""

2. NL-to-LogQL (Loki)

"""自然语言 → LogQL (Loki 查询语言) 转换器"""

class NLToLogQLConverter:
    """自然语言到 LogQL 的转换器"""

    SYSTEM_PROMPT = """你是一位 Loki/LogQL 专家。将自然语言查询
    转换为 LogQL 表达式。

    可用标签:
    - service: 服务名
    - level: 日志级别 (DEBUG/INFO/WARN/ERROR/FATAL)
    - namespace: K8s命名空间
    - pod: Pod名

    规则:
    1. 日志查询格式: {label="value"} |= "filter"
    2. 指标查询格式: rate({label="value"}[5m])
    3. 聚合查询格式: sum by (service) (rate({}[5m]))

    输出JSON: {"logql": "...", "explanation": "...", "confidence": 0.0-1.0}
    """

    EXAMPLES = [
        {
            "nl": "支付服务的所有错误日志",
            "logql": '{service="payment"} |= "ERROR"',
        },
        {
            "nl": "最近10分钟连接超时的日志",
            "logql": '{namespace="production"} |= "timeout" '
                      'or |= "connection refused" |= "Connection pool"',
        },
        {
            "nl": "各服务的错误日志数量统计",
            "logql": 'sum by (service) (count_overhead('
                     '{level="ERROR"}[10m]))',
        },
    ]

3. 对话式故障诊断

"""对话式故障诊断 Agent"""

from dataclasses import dataclass, field
from typing import Optional


@dataclass
class DiagnosticContext:
    """诊断会话上下文"""
    session_id: str
    incident_id: str
    user_question: str
    collected_data: dict = field(default_factory=dict)
    diagnostic_steps: list[dict] = field(default_factory=list)
    current_hypothesis: str = ""
    confidence: float = 0.0


class DiagnosticChatAgent:
    """对话式故障诊断 Agent"""

    def __init__(self):
        self.context = None

    def start_session(self, incident_id: str) -> str:
        """启动诊断会话"""
        self.context = DiagnosticContext(
            session_id=f"diag-{incident_id}",
            incident_id=incident_id,
            user_question="",
        )
        return (
            f"故障诊断会话已启动 (incident: {incident_id})。\n"
            "你可以用自然语言提问,例如:\n"
            "- '当前有哪些服务在报错?'\n"
            "- '支付服务最近有什么变更?'\n"
            "- 'MySQL连接池的使用情况?'\n"
            "- '帮我看一下这个时间段的日志'\n"
            "- '你觉得根因可能是什么?'\n"
        )

    def process_query(self, question: str) -> dict:
        """处理用户查询"""
        self.context.user_question = question

        # 1. 意图识别
        intent = self._classify_intent(question)

        # 2. 根据意图选择处理路径
        if intent == "query_metrics":
            return self._handle_metrics_query(question)
        elif intent == "query_logs":
            return self._handle_log_query(question)
        elif intent == "query_changes":
            return self._handle_change_query(question)
        elif intent == "hypothesis":
            return self._handle_hypothesis_query(question)
        elif intent == "timeline":
            return self._handle_timeline_query(question)
        else:
            return {"response": "请更具体地描述你想查看的内容。"}

    def _classify_intent(self, question: str) -> str:
        """简单意图分类"""
        question_lower = question.lower()
        if any(kw in question_lower for kw in
               ["指标", "监控", "metric", "qps", "延迟", "错误率"]):
            return "query_metrics"
        if any(kw in question_lower for kw in
               ["日志", "log", "报错", "错误信息"]):
            return "query_logs"
        if any(kw in question_lower for kw in
               ["变更", "发布", "部署", "change", "deploy"]):
            return "query_changes"
        if any(kw in question_lower for kw in
               ["根因", "原因", "为什么", "hypothesis", "假设"]):
            return "hypothesis"
        if any(kw in question_lower for kw in
               ["时间线", "timeline", "发生了什么"]):
            return "timeline"
        return "unknown"

    def _handle_metrics_query(self, question: str) -> dict:
        """处理指标查询"""
        # 1. NL → PromQL
        # 2. 执行查询
        # 3. LLM 解读结果
        return {
            "intent": "query_metrics",
            "action": "nl_to_promql",
            "prompt": f"将以下查询转为PromQL并执行: {question}",
        }

    def _handle_log_query(self, question: str) -> dict:
        """处理日志查询"""
        return {
            "intent": "query_logs",
            "action": "nl_to_logql",
            "prompt": f"将以下查询转为LogQL并执行: {question}",
        }

    def _handle_change_query(self, question: str) -> dict:
        """处理变更查询"""
        return {
            "intent": "query_changes",
            "action": "query_change_db",
            "prompt": f"查询变更记录: {question}",
        }

    def _handle_hypothesis_query(self, question: str) -> dict:
        """处理根因假设查询"""
        return {
            "intent": "hypothesis",
            "action": "generate_hypotheses",
            "prompt": "基于已收集的上下文生成根因假设",
        }

    def _handle_timeline_query(self, question: str) -> dict:
        """处理时间线查询"""
        return {
            "intent": "timeline",
            "action": "build_timeline",
            "prompt": "聚合多源数据构建故障时间线",
        }

四、RAG 知识检索增强

"""SRE 知识 RAG —— 检索增强的故障诊断"""

from dataclasses import dataclass, field
from typing import Optional


@dataclass
class KnowledgeChunk:
    source: str           # 文件名/文档名
    content: str         # 文本内容
    metadata: dict        # 标签、创建时间等
    similarity_score: float = 0.0


class SREKnowledgeRAG:
    """SRE 知识库 RAG 检索器"""

    def __init__(self):
        self.knowledge_base: list[KnowledgeChunk] = []
        self.embedding_model = None  # 实际使用 sentence-transformers / OpenAI

    def add_document(self, doc_path: str, content: str, tags: list[str]):
        """添加文档到知识库"""
        # 分块
        chunks = self._chunk_text(content, chunk_size=500, overlap=50)
        for i, chunk in enumerate(chunks):
            self.knowledge_base.append(KnowledgeChunk(
                source=f"{doc_path}#chunk-{i}",
                content=chunk,
                metadata={"tags": tags, "doc": doc_path},
            ))

    def retrieve(self, query: str, top_k: int = 5) -> list[KnowledgeChunk]:
        """检索最相关的知识块"""
        # 实际使用向量相似度检索
        # 这里简化为关键词匹配
        results = []
        query_lower = query.lower()
        for chunk in self.knowledge_base:
            score = self._keyword_similarity(query_lower, chunk.content.lower())
            chunk.similarity_score = score
            results.append(chunk)

        results.sort(key=lambda x: x.similarity_score, reverse=True)
        return results[:top_k]

    def build_rag_prompt(self, query: str, retrieved: list[KnowledgeChunk]) -> str:
        """构建 RAG 增强的 LLM prompt"""
        context_parts = []
        for i, chunk in enumerate(retrieved, 1):
            context_parts.append(
                f"[{i}] 来源: {chunk.source}\n"
                f"内容: {chunk.content}\n"
                f"相关度: {chunk.similarity_score:.2f}\n"
            )

        return f"""你是一位 SRE 故障诊断专家。以下是从知识库中检索到的相关信息:

## 检索到的知识
{chr(10).join(context_parts)}

## 当前故障
{query}

## 请基于检索到的知识回答
1. 知识库中是否有类似故障的处理经验?
2. 推荐的排查步骤是什么?
3. 需要关注哪些关键指标?
"""

    @staticmethod
    def _chunk_text(text: str, chunk_size: int = 500,
                     overlap: int = 50) -> list[str]:
        """文本分块"""
        chunks = []
        start = 0
        while start < len(text):
            end = start + chunk_size
            chunks.append(text[start:end])
            start = end - overlap
        return chunks

    @staticmethod
    def _keyword_similarity(query: str, text: str) -> float:
        """简单关键词相似度"""
        query_words = set(query.split())
        text_words = set(text.split())
        if not query_words:
            return 0.0
        overlap = query_words & text_words
        return len(overlap) / len(query_words)

五、LLM 辅助 RCA 工作流

sequenceDiagram
    participant Alert as 告警系统
    participant Agent as RCA Agent
    participant RAG as 知识库RAG
    participant LLM as LLM引擎
    participant Human as SRE工程师

    Alert->>Agent: 触发告警 + 上下文
    Agent->>Agent: 1. 日志聚类 + 异常检测
    Agent->>Agent: 2. 多源数据关联时间线
    Agent->>RAG: 3. 检索历史类似故障
    RAG-->>Agent: 返回相关知识块
    Agent->>LLM: 4. 生成根因假设(3-5个)
    LLM-->>Agent: 返回假设列表
    Agent->>Agent: 5. 自动验证假设(执行查询)
    Agent->>Human: 6. 推送分析报告
    Human->>Agent: 7. 追问/补充信息
    Agent->>LLM: 8. 对话式深入分析
    LLM-->>Agent: 更新推理
    Agent->>RAG: 9. 沉淀新经验
    Agent->>Human: 10. 最终RCA报告 + Action Items

安全与隐私考量

维度风险缓解措施
敏感数据日志中含密码/Token发送前自动脱敏(正则匹配+替换)
数据泄露LLM API 数据外泄本地部署 / 私有化模型
幻觉LLM 编造不存在的根因强制引用证据 + 置信度标注 + 人工复核
权限LLM 执行危险操作只读分析 + 人工确认执行
成本Token 消耗过大日志预聚合 + 上下文截断 + 缓存
"""日志脱敏处理器"""

import re


class LogSanitizer:
    """日志脱敏 —— 发送给 LLM 前去除敏感信息"""

    PATTERNS = {
        "password": re.compile(
            r'(password|passwd|pwd)\s*[=:]\s*\S+', re.I
        ),
        "token": re.compile(
            r'(token|auth|bearer)\s*[=:]\s*\S+', re.I
        ),
        "api_key": re.compile(
            r'(api[_-]?key|secret[_-]?key)\s*[=:]\s*\S+', re.I
        ),
        "phone": re.compile(r'\b1[3-9]\d{9}\b'),
        "id_card": re.compile(r'\b\d{17}[\dXx]\b'),
        "email": re.compile(r'\b[\w.-]+@[\w.-]+\.\w+\b'),
        "credit_card": re.compile(r'\b\d{16}\b'),
    }

    @classmethod
    def sanitize(cls, text: str) -> str:
        for name, pattern in cls.PATTERNS.items():
            if name in ("email",):
                text = pattern.sub('<EMAIL>', text)
            elif name in ("phone",):
                text = pattern.sub('<PHONE>', text)
            elif name in ("id_card",):
                text = pattern.sub('<ID_CARD>', text)
            elif name in ("credit_card",):
                text = pattern.sub('<CARD>', text)
            else:
                text = pattern.sub(r'\1=<REDACTED>', text)
        return text

效果度量

指标传统方式LLM辅助提升
日志分析时间20-30min3-5min5-10x
MTTI(识别时间)10-15min2-5min3-5x
根因命中率60%80%++20%
初级工程师独立处理率30%60%+2x
知识复用率10%50%+5x

常见坑点

坑点现象解决方案
上下文超长日志太多超出 Token 限制预聚类 + Top-K异常 + 截断
LLM 幻觉编造不存在的根因强制证据引用 + 置信度 + 人工复核
延迟太高LLM 推理慢影响实时性流式输出 + 轻量模型 + 缓存
知识过时RAG 库中的方案已失效定期更新 + 有效性验证
过度依赖工程师不再独立思考LLM为辅,决策在人

关联知识

参考资源

  • OpenAI Cookbook — RAG (Retrieval Augmented Generation)
  • LangChain 文档 — Agent + Tools 架构
  • 《Observability Engineering》— LLM 与可观测性的结合
  • 论文: “LogGPT: Exploring Log Analysis with LLMs”
  • 论文: “Root Cause Analysis with Large Language Models”

学习时间

约 6-8 小时(含搭建一个最小可用的 LLM 日志分析原型)

状态

  • 理解 LLM 在 RCA 中的定位和价值
  • 掌握日志聚类与异常模式提取
  • 能构建 5-Whys 自动推理 prompt
  • 能实现假设生成与验证流程
  • 掌握 NL-to-PromQL/LogQL 转换
  • 理解 RAG 知识检索增强
  • 了解安全隐私考量与脱敏
  • 实际搭建一个对话式故障诊断原型
  • 验证 LLM RCA 在真实故障中的效果