From 661fa2f02b27c571fd7f27ee49f891f9f40e8975 Mon Sep 17 00:00:00 2001 From: Pipeline Agent Date: Sat, 15 Aug 2026 11:40:10 +0800 Subject: [PATCH] =?UTF-8?q?develop:=20=E9=87=8D=E5=81=9A=20fault-log-analy?= =?UTF-8?q?zer=20=E5=88=86=E6=9E=90=E6=A8=A1=E5=9D=97=EF=BC=88=E9=81=B5?= =?UTF-8?q?=E5=AE=88=E6=A8=A1=E5=9D=97=E5=BC=80=E5=8F=91=E8=A7=84=E8=8C=83?= =?UTF-8?q?=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/fault_log_analyzer/__init__.py | 9 +- src/fault_log_analyzer/classifier.py | 94 +++++++++++++- src/fault_log_analyzer/cluster.py | 179 +++++++++++++++++++++++++- src/fault_log_analyzer/config.py | 111 +++++++++++++++- src/fault_log_analyzer/filters.py | 55 +++++++- src/fault_log_analyzer/fingerprint.py | 90 ++++++++++++- src/fault_log_analyzer/models.py | 170 +++++++++++++++++++++++- src/fault_log_analyzer/parser.py | 165 +++++++++++++++++++++++- src/fault_log_analyzer/root_cause.py | 125 +++++++++++++++++- 9 files changed, 989 insertions(+), 9 deletions(-) diff --git a/src/fault_log_analyzer/__init__.py b/src/fault_log_analyzer/__init__.py index 2e7c81a..94a177d 100644 --- a/src/fault_log_analyzer/__init__.py +++ b/src/fault_log_analyzer/__init__.py @@ -1 +1,8 @@ -(已提交至远程 main 分支) \ No newline at end of file +"""fault-log-analyzer:HMS 故障日志捕获与分析模块。 + +提供故障日志捕获管道、归类聚类、根因分析以及 REST API。 +核心逻辑零强制第三方依赖(纯 Python 标准库实现),生产环境可 +按需接入 Kafka / Elasticsearch / MySQL / Redis 等真实后端。 +""" + +__version__ = "0.1.0" diff --git a/src/fault_log_analyzer/classifier.py b/src/fault_log_analyzer/classifier.py index 2e7c81a..8469fb1 100644 --- a/src/fault_log_analyzer/classifier.py +++ b/src/fault_log_analyzer/classifier.py @@ -1 +1,93 @@ -(已提交至远程 main 分支) \ No newline at end of file +"""簇 -> fault_type 归类。 + +依据 architecture.md 5.3.3 第 4 步: +1. 已有人工标注 fault_type.pattern 正则匹配优先。 +2. 簇标注(cluster_id -> fault_type)次之。 +3. 启发式候选(关键字 -> fault_type)兜底。 +""" + +from __future__ import annotations + +import re +from typing import Dict, List, Optional, Sequence + +from .models import FaultType + +# 启发式关键字 -> 候选 fault_type(与 root_cause 规则库一致) +_KEYWORD_HINTS: List[tuple] = [ + ("disk_full", re.compile(r"no space left|disk full|enospc|device full", re.IGNORECASE)), + ("memory_exhausted", re.compile(r"out of memory|oom|oomkilled|memory exhausted", re.IGNORECASE)), + ("cpu_saturation", re.compile(r"cpu (?:throttl|saturat|overload)|high cpu|load average", re.IGNORECASE)), + ("network_unreachable", re.compile(r"network unreachable|connection refused|connect timeout|dns (?:failure|resolution)", re.IGNORECASE)), + ("timeout", re.compile(r"\btimeout\b|timed out|deadline exceeded", re.IGNORECASE)), + ("permission_denied", re.compile(r"permission denied|access denied|forbidden", re.IGNORECASE)), + ("file_corruption", re.compile(r"corrupt|checksum mismatch|integrity check failed", re.IGNORECASE)), +] + +DEFAULT_FAULT_TYPES: List[FaultType] = [ + FaultType( + fault_type="disk_full", + name="磁盘空间不足", + description="磁盘使用率过高或无剩余空间", + pattern=r"no space left|disk full|enospc", + severity="critical", + ), + FaultType( + fault_type="memory_exhausted", + name="内存耗尽", + description="内存不足或 OOM 被杀", + pattern=r"out of memory|oom|oomkilled", + severity="critical", + ), + FaultType( + fault_type="network_unreachable", + name="网络不可达", + description="网络连接失败", + pattern=r"network unreachable|connection refused|connect timeout", + severity="error", + ), +] + + +class Classifier: + """将日志消息归类到 fault_type。""" + + def __init__(self, fault_types: Optional[Sequence[FaultType]] = None): + self.fault_types = list(fault_types) if fault_types else list(DEFAULT_FAULT_TYPES) + self._cluster_map: Dict[int, str] = {} + + def register_cluster(self, cluster_id: int, fault_type: str) -> None: + self._cluster_map[cluster_id] = fault_type + + def classify(self, message: str, cluster_id: Optional[int] = None) -> Optional[str]: + # 1) 簇标注 + if cluster_id is not None and cluster_id in self._cluster_map: + return self._cluster_map[cluster_id] + + # 2) fault_type.pattern 匹配 + for ft in self.fault_types: + if not ft.enabled or not ft.pattern: + continue + try: + if re.search(ft.pattern, message, re.IGNORECASE): + return ft.fault_type + except re.error: + continue + + # 3) 启发式候选 + return self.hint(message) + + @staticmethod + def hint(message: str) -> Optional[str]: + for fault_type, regex in _KEYWORD_HINTS: + if regex.search(message): + return fault_type + return None + + def classify_cluster(self, cluster_id: int, messages: Sequence[str]) -> Optional[str]: + """对一个簇的代表消息集合进行归类。""" + joined = " ".join(messages) + fault_type = self.classify(joined, cluster_id) + if fault_type: + self._cluster_map[cluster_id] = fault_type + return fault_type diff --git a/src/fault_log_analyzer/cluster.py b/src/fault_log_analyzer/cluster.py index 2e7c81a..843d03e 100644 --- a/src/fault_log_analyzer/cluster.py +++ b/src/fault_log_analyzer/cluster.py @@ -1 +1,178 @@ -(已提交至远程 main 分支) \ No newline at end of file +"""TF-IDF 特征 + DBSCAN 聚类(纯 Python 后端)。 + +依据 architecture.md 5.3.3: +1. 特征提取:message 分词 + 模板化。 +2. 向量化:TF-IDF。 +3. 相似度:余弦相似度。 +4. 聚类:DBSCAN(eps=0.75,min_samples=5,可配置;余弦距离 = 1 - 余弦相似度)。 + +生产环境可通过 scikit-learn 后端加速(可选),本模块保证零依赖可用。 +""" + +from __future__ import annotations + +import math +from collections import Counter +from typing import Dict, List, Optional, Sequence, Tuple + +from .parser import LogParser + +Vector = Dict[str, float] + + +def tokenize_all(documents: Sequence[str]) -> List[List[str]]: + return [LogParser.tokenize(d) for d in documents] + + +def tfidf_vectors(documents: Sequence[str]) -> List[Vector]: + """为文档集合计算 TF-IDF 向量。""" + tokenized = tokenize_all(documents) + n = len(tokenized) + if n == 0: + return [] + + df: Counter = Counter() + for tokens in tokenized: + for term in set(tokens): + df[term] += 1 + + vectors: List[Vector] = [] + for tokens in tokenized: + if not tokens: + vectors.append({}) + continue + tf = Counter(tokens) + length = len(tokens) + vec: Vector = {} + for term, count in tf.items(): + idf = math.log((1.0 + n) / (1.0 + df[term])) + 1.0 + vec[term] = (count / length) * idf + vectors.append(vec) + return vectors + + +def _dot(a: Vector, b: Vector) -> float: + if len(a) > len(b): + a, b = b, a + return sum(v * b.get(k, 0.0) for k, v in a.items()) + + +def _norm(v: Vector) -> float: + return math.sqrt(sum(x * x for x in v.values())) + + +def cosine_similarity(a: Vector, b: Vector) -> float: + na, nb = _norm(a), _norm(b) + if na == 0.0 or nb == 0.0: + return 0.0 + return _dot(a, b) / (na * nb) + + +def cosine_distance(a: Vector, b: Vector) -> float: + return 1.0 - cosine_similarity(a, b) + + +def _region_query(i: int, vectors: Sequence[Vector], eps: float) -> List[int]: + return [ + j + for j in range(len(vectors)) + if cosine_distance(vectors[i], vectors[j]) <= eps + ] + + +def dbscan(vectors: Sequence[Vector], eps: float, min_samples: int) -> List[int]: + """DBSCAN 聚类(余弦距离)。返回每个样本的簇标签,-1 表示噪声。 + + 注意:当样本数少于 min_samples 时,所有样本都将是噪声(-1)。 + """ + n = len(vectors) + labels: List[Optional[int]] = [None] * n + cluster_id = 0 + + for p in range(n): + if labels[p] is not None: + continue + neighbors = _region_query(p, vectors, eps) + if len(neighbors) < min_samples: + labels[p] = -1 + continue + + labels[p] = cluster_id + seeds = neighbors[:] + # 用索引扫描 seeds,避免在迭代中修改列表导致重复处理过多 + i = 0 + while i < len(seeds): + q = seeds[i] + i += 1 + if labels[q] == -1: + labels[q] = cluster_id + if labels[q] is not None: + continue + labels[q] = cluster_id + q_neighbors = _region_query(q, vectors, eps) + if len(q_neighbors) >= min_samples: + for r in q_neighbors: + if labels[r] is None or labels[r] == -1: + seeds.append(r) + cluster_id += 1 + + return [-1 if lbl is None else lbl for lbl in labels] + + +class Clusterer: + """面向日志消息的聚类器。""" + + def __init__(self, eps: float = 0.75, min_samples: int = 5): + self.eps = eps + self.min_samples = min_samples + self.labels: List[int] = [] + self.centroids: Dict[int, Vector] = {} + self.n_clusters = 0 + + def fit_predict(self, messages: Sequence[str]) -> List[int]: + vectors = tfidf_vectors(list(messages)) + self.labels = dbscan(vectors, self.eps, self.min_samples) + self.centroids = self._compute_centroids(vectors, self.labels) + self.n_clusters = len([c for c in self.centroids if c >= 0]) + return self.labels + + @staticmethod + def _compute_centroids(vectors: Sequence[Vector], labels: Sequence[int]) -> Dict[int, Vector]: + sums: Dict[int, List[float]] = {} + counts: Dict[int, int] = {} + terms: Dict[int, set] = {} + for vec, label in zip(vectors, labels): + if label < 0: + continue + if label not in sums: + sums[label] = [] + terms[label] = set() + counts[label] = 0 + counts[label] += 1 + terms[label].update(vec.keys()) + centroids: Dict[int, Vector] = {} + for label, term_set in terms.items(): + centroid: Vector = {} + for term in term_set: + total = sum( + v.get(term, 0.0) + for i, v in enumerate(vectors) + if labels[i] == label + ) + centroid[term] = total / counts[label] + centroids[label] = centroid + return centroids + + def assign(self, message: str) -> int: + """增量归类:新消息分配到余弦距离最近的簇,否则为噪声(-1)。""" + if not self.centroids: + return -1 + vec = tfidf_vectors([message])[0] + best_label = -1 + best_dist = self.eps + for label, centroid in self.centroids.items(): + dist = cosine_distance(vec, centroid) + if dist <= best_dist: + best_dist = dist + best_label = label + return best_label diff --git a/src/fault_log_analyzer/config.py b/src/fault_log_analyzer/config.py index 2e7c81a..1227687 100644 --- a/src/fault_log_analyzer/config.py +++ b/src/fault_log_analyzer/config.py @@ -1 +1,110 @@ -(已提交至远程 main 分支) \ No newline at end of file +"""配置模型与环境变量加载。 + +零第三方依赖,通过环境变量(带合理默认值)与可选显式参数配置。 +""" + +from __future__ import annotations + +import os +from dataclasses import dataclass, field, asdict +from typing import Any, Dict + + +def _env(name: str, default: str) -> str: + return os.environ.get(name, default) + + +def _env_float(name: str, default: float) -> float: + raw = os.environ.get(name) + if raw is None or raw == "": + return default + try: + return float(raw) + except ValueError: + return default + + +def _env_int(name: str, default: int) -> int: + raw = os.environ.get(name) + if raw is None or raw == "": + return default + try: + return int(raw) + except ValueError: + return default + + +def _env_bool(name: str, default: bool) -> bool: + raw = os.environ.get(name) + if raw is None or raw == "": + return default + return raw.strip().lower() in ("1", "true", "yes", "on") + + +@dataclass +class Config: + """fault-log-analyzer 配置。""" + + # Kafka + kafka_bootstrap_servers: str = "localhost:9092" + kafka_group_id: str = "fault-log-analyzer" + kafka_logs_raw_topic: str = "logs.raw" + kafka_logs_fault_topic: str = "logs.fault" + + # Elasticsearch + es_hosts: str = "http://localhost:9200" + es_index_fault_prefix: str = "hms-fault-log" + + # MySQL + mysql_host: str = "localhost" + mysql_port: int = 3306 + mysql_user: str = "hms" + mysql_password: str = "" + mysql_db: str = "hms" + + # Redis + redis_url: str = "redis://localhost:6379/0" + + # 聚类 + cluster_eps: float = 0.75 + cluster_min_samples: int = 5 + cluster_num_perm: int = 64 + + # 去重窗口(秒) + dedup_window_seconds: int = 300 + + # 根因分析时间窗口(分钟) + root_cause_window_minutes: int = 5 + + # REST API + api_host: str = "127.0.0.1" + api_port: int = 8080 + storage: str = "memory" + + @classmethod + def from_env(cls) -> "Config": + return cls( + kafka_bootstrap_servers=_env("KAFKA_BOOTSTRAP_SERVERS", "localhost:9092"), + kafka_group_id=_env("KAFKA_GROUP_ID", "fault-log-analyzer"), + kafka_logs_raw_topic=_env("KAFKA_LOGS_RAW_TOPIC", "logs.raw"), + kafka_logs_fault_topic=_env("KAFKA_LOGS_FAULT_TOPIC", "logs.fault"), + es_hosts=_env("ES_HOSTS", "http://localhost:9200"), + es_index_fault_prefix=_env("ES_INDEX_FAULT_PREFIX", "hms-fault-log"), + mysql_host=_env("MYSQL_HOST", "localhost"), + mysql_port=_env_int("MYSQL_PORT", 3306), + mysql_user=_env("MYSQL_USER", "hms"), + mysql_password=_env("MYSQL_PASSWORD", ""), + mysql_db=_env("MYSQL_DB", "hms"), + redis_url=_env("REDIS_URL", "redis://localhost:6379/0"), + cluster_eps=_env_float("CLUSTER_EPS", 0.75), + cluster_min_samples=_env_int("CLUSTER_MIN_SAMPLES", 5), + cluster_num_perm=_env_int("CLUSTER_NUM_PERM", 64), + dedup_window_seconds=_env_int("DEDUP_WINDOW_SECONDS", 300), + root_cause_window_minutes=_env_int("ROOT_CAUSE_WINDOW_MINUTES", 5), + api_host=_env("API_HOST", "127.0.0.1"), + api_port=_env_int("API_PORT", 8080), + storage=_env("STORAGE", "memory"), + ) + + def to_dict(self) -> Dict[str, Any]: + return asdict(self) diff --git a/src/fault_log_analyzer/filters.py b/src/fault_log_analyzer/filters.py index 2e7c81a..5402985 100644 --- a/src/fault_log_analyzer/filters.py +++ b/src/fault_log_analyzer/filters.py @@ -1 +1,54 @@ -(已提交至远程 main 分支) \ No newline at end of file +"""故障日志过滤规则。 + +依据 architecture.md 5.3.2:过滤规则配置化(fault_filter_rule 表, +支持 level / pattern / exclude_pattern)。未配置规则时默认捕获 +ERROR 与 FATAL 级别。 +""" + +from __future__ import annotations + +import re +from typing import Iterable, List, Optional + +from .models import FaultFilterRule, ParsedLog + +DEFAULT_LEVELS = {"ERROR", "FATAL"} + + +class FaultFilter: + """按规则判断一条日志是否为故障日志。""" + + def __init__(self, rules: Optional[Iterable[FaultFilterRule]] = None): + self.rules: List[FaultFilterRule] = list(rules) if rules else [] + + def add_rule(self, rule: FaultFilterRule) -> None: + self.rules.append(rule) + + def should_capture(self, parsed: ParsedLog) -> bool: + level = (parsed.level or "INFO").upper() + message = parsed.message or "" + + if not self.rules: + return level in DEFAULT_LEVELS + + # 命中任意一条 enabled 规则即捕获;被明确排除则跳过 + for rule in self.rules: + if not rule.enabled: + continue + if self._match_rule(rule, level, message): + return True + return False + + @staticmethod + def _match_rule(rule: FaultFilterRule, level: str, message: str) -> bool: + if rule.level and level != rule.level.upper(): + return False + if rule.pattern and not re.search(rule.pattern, message, re.IGNORECASE): + return False + if rule.exclude_pattern and re.search(rule.exclude_pattern, message, re.IGNORECASE): + return False + return True + + @staticmethod + def compile_rules(rows: Iterable[dict]) -> List[FaultFilterRule]: + return [FaultFilterRule.from_dict(r) for r in rows] diff --git a/src/fault_log_analyzer/fingerprint.py b/src/fault_log_analyzer/fingerprint.py index 2e7c81a..5d6c5c4 100644 --- a/src/fault_log_analyzer/fingerprint.py +++ b/src/fault_log_analyzer/fingerprint.py @@ -1 +1,89 @@ -(已提交至远程 main 分支) \ No newline at end of file +"""MinHash 指纹与去重。 + +纯 Python 实现 MinHash(字符 3-gram shingles + 稳定哈希 + 随机置换 +模拟),生成稳定十六进制指纹。用于相同 (host, service, message_fingerprint) +短窗口去重(architecture.md 5.3.2,Redis `hms:log:fp:{fingerprint}`)。 +""" + +from __future__ import annotations + +import hashlib +import random +import re +import struct +import zlib +from typing import Dict, Iterable, Optional, Set + +_MAX_HASH = (1 << 32) - 1 + + +def _h(token: str) -> int: + return zlib.crc32(token.encode("utf-8")) & _MAX_HASH + + +def shingles(text: str, k: int = 3) -> Set[str]: + """将文本切分为 k-gram 集合(字符级,跨空白压缩)。""" + cleaned = re.sub(r"\s+", " ", text or "").strip() + if len(cleaned) < k: + # 短文本退化为整串 + return {cleaned} if cleaned else set() + return {cleaned[i : i + k] for i in range(len(cleaned) - k + 1)} + + +class MinHash: + """MinHash 指纹器,相同(或高度相似)文本产生相同/相近指纹。""" + + def __init__(self, num_perm: int = 64, seed: int = 42): + self.num_perm = num_perm + rng = random.Random(seed) + self.a = [rng.randint(1, _MAX_HASH) for _ in range(num_perm)] + self.b = [rng.randint(0, _MAX_HASH) for _ in range(num_perm)] + + def signature(self, tokens: Iterable[str]) -> list: + hashes = [_h(t) for t in tokens] + if not hashes: + return [0] * self.num_perm + sig = [] + for ai, bi in zip(self.a, self.b): + sig.append(min(((ai * x + bi) & _MAX_HASH) for x in hashes)) + return sig + + def fingerprint(self, text: str) -> str: + """返回文本的 16 字符十六进制指纹。""" + tokens = shingles(text) + if not tokens: + return "0" * 16 + sig = self.signature(tokens) + digest = hashlib.sha256() + for v in sig: + digest.update(struct.pack(">I", v)) + return digest.hexdigest()[:16] + + +class Deduplicator: + """短窗口去重器。生产可替换为 Redis SETNX,测试用进程内实现。""" + + def __init__(self, window_seconds: int = 300, clock=None): + self.window_seconds = window_seconds + self._clock = clock or _time_now + self._seen: Dict[str, float] = {} + + def is_duplicate(self, fingerprint: str) -> bool: + now = self._clock() + last = self._seen.get(fingerprint) + if last is not None and (now - last) <= self.window_seconds: + return True + self._seen[fingerprint] = now + return False + + def mark(self, fingerprint: str) -> None: + self._seen[fingerprint] = self._clock() + + def clear(self) -> None: + self._seen.clear() + + +def _time_now() -> float: + import time + + return time.time() diff --git a/src/fault_log_analyzer/models.py b/src/fault_log_analyzer/models.py index 2e7c81a..974077f 100644 --- a/src/fault_log_analyzer/models.py +++ b/src/fault_log_analyzer/models.py @@ -1 +1,169 @@ -(已提交至远程 main 分支) \ No newline at end of file +"""数据模型定义。 + +依据 docs/01-design/database-design.md 中 fault_log / root_cause / +fault_type / fault_filter_rule / event 表结构,提供进程内数据模型。 +核心实现零第三方依赖,全部使用标准库 dataclass。 +""" + +from __future__ import annotations + +from dataclasses import dataclass, field, asdict +from datetime import datetime, timezone +from typing import Any, Dict, List, Optional + + +def _now() -> datetime: + return datetime.now(timezone.utc) + + +def _to_dt(value: Any) -> datetime: + """将字符串/时间戳转换为带时区的 datetime。""" + if value is None: + return _now() + if isinstance(value, datetime): + if value.tzinfo is None: + return value.replace(tzinfo=timezone.utc) + return value + if isinstance(value, (int, float)): + return datetime.fromtimestamp(value, tz=timezone.utc) + text = str(value).strip() + if text.endswith("Z"): + text = text[:-1] + "+00:00" + try: + dt = datetime.fromisoformat(text) + except ValueError: + dt = datetime.fromisoformat(text.replace(" ", "T")) + if dt.tzinfo is None: + dt = dt.replace(tzinfo=timezone.utc) + return dt + + +@dataclass +class FaultLog: + """故障日志(归类后),对应 fault_log 表。""" + + fault_log_id: str + host_id: str + message: str + level: str = "ERROR" + fingerprint: str = "" + fault_type: Optional[str] = None + cluster_id: Optional[str] = None + service: Optional[str] = None + trace_id: Optional[str] = None + occurred_at: datetime = field(default_factory=_now) + count: int = 1 + + def to_dict(self) -> Dict[str, Any]: + data = asdict(self) + data["occurred_at"] = self.occurred_at.isoformat() + return data + + @classmethod + def from_dict(cls, data: Dict[str, Any]) -> "FaultLog": + payload = dict(data) + payload["occurred_at"] = _to_dt(payload.get("occurred_at")) + return cls(**{k: v for k, v in payload.items() if k in cls.__dataclass_fields__}) + + +@dataclass +class RootCause: + """根因结论,对应 root_cause 表。""" + + fault_log_id: str + cause_type: str + evidence: List[Dict[str, Any]] = field(default_factory=list) + confidence: float = 0.0 + analysis_at: datetime = field(default_factory=_now) + + def to_dict(self) -> Dict[str, Any]: + data = asdict(self) + data["analysis_at"] = self.analysis_at.isoformat() + return data + + @classmethod + def from_dict(cls, data: Dict[str, Any]) -> "RootCause": + payload = dict(data) + payload["analysis_at"] = _to_dt(payload.get("analysis_at")) + return cls(**{k: v for k, v in payload.items() if k in cls.__dataclass_fields__}) + + +@dataclass +class FaultType: + """故障类型,对应 fault_type 表。""" + + fault_type: str + name: str + description: str = "" + pattern: str = "" + severity: str = "warning" + enabled: bool = True + + def to_dict(self) -> Dict[str, Any]: + return asdict(self) + + @classmethod + def from_dict(cls, data: Dict[str, Any]) -> "FaultType": + return cls(**{k: v for k, v in data.items() if k in cls.__dataclass_fields__}) + + +@dataclass +class FaultFilterRule: + """故障日志过滤规则,对应 fault_filter_rule 表。""" + + name: str + level: str = "ERROR" + pattern: str = "" + exclude_pattern: str = "" + enabled: bool = True + + def to_dict(self) -> Dict[str, Any]: + return asdict(self) + + @classmethod + def from_dict(cls, data: Dict[str, Any]) -> "FaultFilterRule": + return cls(**{k: v for k, v in data.items() if k in cls.__dataclass_fields__}) + + +@dataclass +class Event: + """检测事件(用于根因分析的关联输入),对应 event 表。""" + + event_id: str + host_id: str + metric: str + value: float + threshold: float + operator: str = "gt" + status: str = "firing" + severity: str = "warning" + fired_at: datetime = field(default_factory=_now) + + def to_dict(self) -> Dict[str, Any]: + data = asdict(self) + data["fired_at"] = self.fired_at.isoformat() + return data + + @classmethod + def from_dict(cls, data: Dict[str, Any]) -> "Event": + payload = dict(data) + payload["fired_at"] = _to_dt(payload.get("fired_at")) + return cls(**{k: v for k, v in payload.items() if k in cls.__dataclass_fields__}) + + +@dataclass +class ParsedLog: + """解析后的结构化日志条目(捕获管道内部表示)。""" + + timestamp: datetime = field(default_factory=_now) + host_id: str = "" + service: str = "" + level: str = "INFO" + message: str = "" + trace_id: Optional[str] = None + raw: str = "" + + def to_dict(self) -> Dict[str, Any]: + data = asdict(self) + data["timestamp"] = self.timestamp.isoformat() + return data diff --git a/src/fault_log_analyzer/parser.py b/src/fault_log_analyzer/parser.py index 2e7c81a..59f6d7a 100644 --- a/src/fault_log_analyzer/parser.py +++ b/src/fault_log_analyzer/parser.py @@ -1 +1,164 @@ -(已提交至远程 main 分支) \ No newline at end of file +"""日志结构化解析与模板化。 + +支持常见日志形态: +- JSON 对象(含 timestamp/level/message/host_id/service/trace_id) +- log4j 风格:`2025-01-01 10:18:00 ERROR [app] message` +- syslog 风格:`Jan 1 10:18:00 host app[pid]: message` + +模板化用于后续聚类特征提取:将数字、IP、UUID、十六进制串等替换为 +稳定占位符,从而将“同一类”日志归并。 +""" + +from __future__ import annotations + +import json +import re +from typing import Any, Dict, Optional + +from .models import ParsedLog, _to_dt + +# 常见级别关键词 +_LEVELS = ("FATAL", "ERROR", "WARN", "WARNING", "INFO", "DEBUG", "TRACE") + +# log4j / 通用时间前缀:2025-01-01 10:18:00,123 或 2025-01-01T10:18:00Z +_TS_RE = re.compile( + r"^(?P\d{4}-\d{2}-\d{2}[T ]\d{2}:\d{2}:\d{2}(?:[.,]\d{3})?(?:Z|[+-]\d{2}:?\d{2})?)" +) + +_LOG4J_RE = re.compile( + r"^\s*(\d{4}-\d{2}-\d{2}[T ]\d{2}:\d{2}:\d{2}(?:[.,]\d{3})?(?:Z|[+-]\d{2}:?\d{2})?)\s+" + r"([A-Za-z]+)\s+\[?([^\]]*?)\]?\s*(?:-)?\s*(.*)$" +) + +_SYSLOG_RE = re.compile( + r"^<(?P\d{1,3})>\d?\s*(?P\w{3}\s+\d{1,2}\s+\d{2}:\d{2}:\d{2})\s+" + r"(?P\S+)\s+(?P\S+?)(?:\[(?P\d+)\])?:\s*(?P.*)$" +) + + +class LogParser: + """将原始日志(字符串或映射)解析为结构化 ParsedLog。""" + + def parse(self, raw: Any) -> ParsedLog: + if isinstance(raw, ParsedLog): + return raw + if isinstance(raw, dict): + return self._parse_dict(raw) + if isinstance(raw, (bytes, bytearray)): + raw = bytes(raw).decode("utf-8", errors="replace") + return self._parse_line(str(raw)) + + def _parse_dict(self, data: Dict[str, Any]) -> ParsedLog: + message = str(data.get("message", data.get("msg", ""))) + level = str(data.get("level", data.get("severity", "INFO"))).upper() + host_id = str(data.get("host_id", data.get("host", ""))) + service = str(data.get("service", data.get("app", ""))) + trace_id = data.get("trace_id", data.get("traceId")) or None + timestamp = data.get("timestamp", data.get("time")) + ts = _to_dt(timestamp) if timestamp is not None else ParsedLog.__dataclass_fields__["timestamp"].default_factory() # type: ignore[attr-defined] + return ParsedLog( + timestamp=ts, + host_id=host_id, + service=service, + level=level, + message=message, + trace_id=str(trace_id) if trace_id is not None else None, + raw=json.dumps(data, ensure_ascii=False, default=str), + ) + + def _parse_line(self, line: str) -> ParsedLog: + stripped = line.strip() + # 1) JSON 整行 + if stripped.startswith("{") and stripped.endswith("}"): + try: + return self._parse_dict(json.loads(stripped)) + except (json.JSONDecodeError, ValueError): + pass + + # 2) log4j 风格 + m = _LOG4J_RE.match(stripped) + if m: + ts_text, level, thread, message = m.groups() + level = level.upper() + thread = (thread or "").strip() + return ParsedLog( + timestamp=_to_dt(ts_text), + host_id=self._extract_host(message) or thread, + service=thread, + level=level, + message=message.strip(), + trace_id=self._extract_trace(message), + raw=line, + ) + + # 3) syslog 风格 + m = _SYSLOG_RE.match(stripped) + if m: + host = m.group("host") + app = m.group("app") + message = m.group("message") or "" + return ParsedLog( + timestamp=_to_dt(m.group("ts")), + host_id=self._extract_host(message) or host, + service=app, + level=self._guess_level(message), + message=message.strip(), + trace_id=self._extract_trace(message), + raw=line, + ) + + # 4) 兜底:整行作为 message,尝试在行内识别级别 + return ParsedLog( + level=self._guess_level(line), + message=stripped, + host_id=self._extract_host(line), + trace_id=self._extract_trace(line), + raw=line, + ) + + @staticmethod + def _extract_host(text: str) -> str: + m = re.search(r"\bhost[=:]\s*([A-Za-z0-9_.\-]+)", text, re.IGNORECASE) + if m: + return m.group(1) + m = re.search(r"\bhost_id[=:]\s*([A-Za-z0-9_.\-]+)", text, re.IGNORECASE) + return m.group(1) if m else "" + + @staticmethod + def _extract_trace(text: str) -> Optional[str]: + m = re.search( + r"\b(?:trace_id|traceId|trace)[=:]\s*([A-Za-z0-9\-]{6,64})", text + ) + return m.group(1) if m else None + + @staticmethod + def _guess_level(text: str) -> str: + upper = text.upper() + for level in _LEVELS: + # 级别通常出现在行首或紧跟时间戳后 + if re.search(rf"\b{level}\b", upper): + return level + return "INFO" + + # ------------------------------------------------------------------ + # 模板化(供聚类特征提取) + # ------------------------------------------------------------------ + @staticmethod + def template(message: str) -> str: + """将 message 中的变量替换为稳定占位符,用于聚类归并。""" + text = str(message) + text = re.sub(r"\b\d{1,3}(?:\.\d{1,3}){3}\b", "", text) + text = re.sub(r"\b[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-" + r"[0-9a-fA-F]{4}-[0-9a-fA-F]{12}\b", "", text) + text = re.sub(r"\b0x[0-9a-fA-F]+\b", "", text) + text = re.sub(r"\b\d{10,}\b", "", text) + text = re.sub(r"\b\d+(?:\.\d+)?\b", "", text) + text = re.sub(r"\s+", " ", text) + return text.strip() + + @staticmethod + def tokenize(message: str) -> list: + """分词:模板化后按非字母数字切分并过滤过短词。""" + templated = LogParser.template(message) + tokens = re.findall(r"[A-Za-z_][A-Za-z0-9_\-]*", templated) + return [t for t in tokens if len(t) > 1] diff --git a/src/fault_log_analyzer/root_cause.py b/src/fault_log_analyzer/root_cause.py index 2e7c81a..77d27af 100644 --- a/src/fault_log_analyzer/root_cause.py +++ b/src/fault_log_analyzer/root_cause.py @@ -1 +1,124 @@ -(已提交至远程 main 分支) \ No newline at end of file +"""根因分析(辅助性)。 + +依据 architecture.md 5.3.4: +1. 时间/主机关联:fault_log.host_id 与 event.host_id 匹配,时间差 ≤ 窗口。 +2. 指标-故障规则库:指标阈值 -> 归因;日志关键字 -> 归因。 +3. trace 关联:trace_id 相同的日志与事件串链。 +4. 产出 root_cause(evidence + confidence)。 + +规则优先,统计辅助,置信度评分。 +""" + +from __future__ import annotations + +import math +import re +from datetime import timedelta +from typing import Dict, List, Optional, Sequence + +from .models import Event, FaultLog, RootCause + +_METRIC_RULES: List[dict] = [ + {"metric": "disk_used_percent", "operator": "gte", "threshold": 90.0, "cause": "disk_full"}, + {"metric": "mem_used_percent", "operator": "gte", "threshold": 90.0, "cause": "memory_exhausted"}, + {"metric": "cpu_usage", "operator": "gte", "threshold": 95.0, "cause": "cpu_saturation"}, +] + +_LOG_CAUSE_HINTS: List[tuple] = [ + ("disk_full", re.compile(r"no space left|disk full|enospc|device full", re.IGNORECASE)), + ("memory_exhausted", re.compile(r"out of memory|oom|oomkilled|memory exhausted", re.IGNORECASE)), + ("cpu_saturation", re.compile(r"cpu throttl|high cpu|load average|saturat", re.IGNORECASE)), + ("network_unreachable", re.compile(r"network unreachable|connection refused|connect timeout", re.IGNORECASE)), +] + + +def _op_ok(operator: str, value: float, threshold: float) -> bool: + if operator in ("gt", ">"): + return value > threshold + if operator in ("gte", ">="): + return value >= threshold + if operator in ("lt", "<"): + return value < threshold + if operator in ("lte", "<="): + return value <= threshold + return False + + +class RootCauseAnalyzer: + """根据故障日志与近期事件进行根因分析。""" + + def __init__(self, window_minutes: int = 5): + self.window_minutes = window_minutes + + def analyze( + self, + fault_log: FaultLog, + events: Sequence[Event], + related_logs: Optional[Sequence[FaultLog]] = None, + ) -> RootCause: + related_logs = list(related_logs or []) + now = fault_log.occurred_at + window = timedelta(minutes=self.window_minutes) + + host_events = [ + e + for e in events + if e.host_id == fault_log.host_id + and abs((e.fired_at - now).total_seconds()) <= window.total_seconds() + ] + + evidence: List[Dict] = [] + cause_scores: Dict[str, float] = {} + + for message in [fault_log.message] + [l.message for l in related_logs]: + for cause, regex in _LOG_CAUSE_HINTS: + if regex.search(message or ""): + cause_scores[cause] = cause_scores.get(cause, 0.0) + 0.6 + evidence.append({"type": "log", "message": message, "cause": cause}) + + for event in host_events: + for rule in _METRIC_RULES: + if event.metric != rule["metric"]: + continue + if _op_ok(event.operator, event.value, rule["threshold"]): + cause = rule["cause"] + cause_scores[cause] = cause_scores.get(cause, 0.0) + 0.8 + evidence.append( + { + "type": "event", + "event_id": event.event_id, + "metric": event.metric, + "value": event.value, + "threshold": rule["threshold"], + "cause": cause, + } + ) + + if fault_log.trace_id and related_logs: + trace_logs = [l for l in related_logs if l.trace_id == fault_log.trace_id] + if trace_logs: + evidence.append({"type": "trace", "trace_id": fault_log.trace_id, "count": len(trace_logs)}) + for l in trace_logs: + for cause, regex in _LOG_CAUSE_HINTS: + if regex.search(l.message or ""): + cause_scores[cause] = cause_scores.get(cause, 0.0) + 0.3 + + if not cause_scores: + return RootCause( + fault_log_id=fault_log.fault_log_id, + cause_type="unknown", + evidence=evidence, + confidence=0.0, + analysis_at=now, + ) + + best_cause = max(cause_scores, key=cause_scores.get) + raw_score = cause_scores[best_cause] + confidence = round(1.0 - math.exp(-raw_score), 4) + return RootCause( + fault_log_id=fault_log.fault_log_id, + cause_type=best_cause, + evidence=evidence, + confidence=confidence, + analysis_at=now, + )