From 1095788bb4a8578d17910595282494dfb28bbf72 Mon Sep 17 00:00:00 2001 From: Pipeline Agent Date: Sat, 15 Aug 2026 13:01:43 +0800 Subject: [PATCH] =?UTF-8?q?test:=20=E9=87=8D=E5=81=9A=20fault-log-analyzer?= =?UTF-8?q?=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=EF=BC=88test=E9=98=B6=E6=AE=B5=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pyproject.toml | 32 +++- requirements.txt | 9 +- src/fault_log_analyzer/__init__.py | 9 +- src/fault_log_analyzer/__main__.py | 52 ++++++- src/fault_log_analyzer/api.py | 192 +++++++++++++++++++++++- 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/integrations.py | 200 ++++++++++++++++++++++++- src/fault_log_analyzer/models.py | 184 ++++++++++++++++++++++- src/fault_log_analyzer/parser.py | 166 +++++++++++++++++++- src/fault_log_analyzer/pipeline.py | 98 +++++++++++- src/fault_log_analyzer/root_cause.py | 151 ++++++++++++++++++- src/fault_log_analyzer/storage.py | 181 +++++++++++++++++++++- src/fault_log_analyzer/workers.py | 109 +++++++++++++- 17 files changed, 1895 insertions(+), 17 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 0a75ce1..f36dfb8 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1 +1,31 @@ -项目构建与依赖声明(已提交至 git) \ No newline at end of file +[build-system] +requires = ["setuptools>=61.0"] +build-backend = "setuptools.build_meta" + +[project] +name = "fault-log-analyzer" +version = "0.1.0" +description = "HMS fault log capture, classification and root cause analysis module" +readme = "README.md" +requires-python = ">=3.10" +license = { text = "Apache-2.0" } +authors = [{ name = "HMS Team" }] +dependencies = [] + +[project.optional-dependencies] +kafka = ["kafka-python>=2.0.2"] +elasticsearch = ["elasticsearch>=8.0.0"] +mysql = ["pymysql>=1.1.0"] +redis = ["redis>=4.5.0"] +ml = ["numpy>=1.24.0", "scikit-learn>=1.2.0"] +dev = ["pytest>=7.0.0"] + +[project.scripts] +fault-log-analyzer = "fault_log_analyzer.__main__:main" + +[tool.setuptools.packages.find] +where = ["src"] + +[tool.pytest.ini_options] +testpaths = ["tests"] +pythonpath = ["src"] diff --git a/requirements.txt b/requirements.txt index a023a9b..92ad8ae 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1 +1,8 @@ -可选依赖(已提交至 git) \ No newline at end of file +# fault-log-analyzer 核心逻辑零第三方依赖(纯 Python 标准库),可直接运行。 +# 生产环境对接真实存储/消息总线时,按需安装以下可选依赖: +kafka-python>=2.0.2 # Kafka 消费/生产 +elasticsearch>=8.0.0 # Elasticsearch 日志存储 +pymysql>=1.1.0 # MySQL 元数据/结果存储 +redis>=4.5.0 # 去重/缓存 +numpy>=1.24.0 # 可选:向量化加速(聚类) +scikit-learn>=1.2.0 # 可选:生产级 DBSCAN/TF-IDF 后端 diff --git a/src/fault_log_analyzer/__init__.py b/src/fault_log_analyzer/__init__.py index 618a5f8..94a177d 100644 --- a/src/fault_log_analyzer/__init__.py +++ b/src/fault_log_analyzer/__init__.py @@ -1 +1,8 @@ -包元信息与版本 __version__ = 0.1.0(已提交至 git) \ 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/__main__.py b/src/fault_log_analyzer/__main__.py index 8b1628b..bd81913 100644 --- a/src/fault_log_analyzer/__main__.py +++ b/src/fault_log_analyzer/__main__.py @@ -1 +1,51 @@ -CLI 入口:构建存储、启动 REST API(已提交至 git) \ No newline at end of file +"""CLI 入口。 + +用法: + python -m fault_log_analyzer --storage memory --api-host 127.0.0.1 --api-port 8080 +""" + +from __future__ import annotations + +import argparse +import sys +from typing import List, Optional + +from .api import run_server +from .config import Config +from .storage import MemoryStorage, Storage + + +def build_storage(storage_name: str) -> Storage: + """根据配置构建存储。生产后端(mysql 等)在 integrations.py 中。""" + if storage_name in ("memory", "mem"): + return MemoryStorage() + if storage_name == "mysql": + from .integrations import MySQLFaultRepository + + cfg = Config.from_env() + return MySQLFaultRepository( + host=cfg.mysql_host, + port=cfg.mysql_port, + user=cfg.mysql_user, + password=cfg.mysql_password, + db=cfg.mysql_db, + ) + raise ValueError(f"unknown storage: {storage_name}") + + +def main(argv: Optional[List[str]] = None) -> int: + cfg = Config.from_env() + parser = argparse.ArgumentParser(prog="fault-log-analyzer", description="HMS 故障日志捕获与分析模块") + parser.add_argument("--storage", default=cfg.storage, help="存储后端(memory/mysql)") + parser.add_argument("--api-host", default=cfg.api_host, help="REST API 监听地址") + parser.add_argument("--api-port", type=int, default=cfg.api_port, help="REST API 监听端口") + args = parser.parse_args(argv) + + storage = build_storage(args.storage) + print(f"fault-log-analyzer serving on http://{args.api_host}:{args.api_port} (storage={args.storage})") + run_server(args.api_host, args.api_port, storage) + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/src/fault_log_analyzer/api.py b/src/fault_log_analyzer/api.py index 65ec182..d35a30d 100644 --- a/src/fault_log_analyzer/api.py +++ b/src/fault_log_analyzer/api.py @@ -1 +1,191 @@ -REST API(已提交至 git) \ No newline at end of file +"""REST API(标准库 http.server 实现)。 + +依据 api-design.md 第 6 节,提供: +- GET /healthz +- GET /api/v1/fault-logs 故障日志列表(分页/过滤) +- GET /api/v1/fault-logs/{id} 故障日志详情 +- GET /api/v1/fault-logs/{id}/root-cause 根因结论 +- GET /api/v1/fault-types 故障类型列表 +- POST /api/v1/fault-types 新建故障类型 +- GET /api/v1/fault-filters 过滤规则列表 +- POST /api/v1/fault-filters 新建过滤规则 + +统一响应包装:{"code": 0, "message": "ok", "data": ...}。 +""" + +from __future__ import annotations + +import json +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from typing import Any, Dict, Optional, Type +from urllib.parse import parse_qs, urlparse + +from .models import FaultFilterRule, FaultType +from .storage import Storage + + +class FaultLogAPIHandler(BaseHTTPRequestHandler): + """绑定 MemoryStorage 的请求处理器。storage 由 create_handler 注入。""" + + storage: Storage = None # type: ignore[assignment] + server_version = "FaultLogAnalyzer/0.1" + + # ------------------------------------------------------------------ + # 工具方法 + # ------------------------------------------------------------------ + def _send_json(self, status: int, payload: Dict[str, Any]) -> None: + body = json.dumps(payload, ensure_ascii=False, default=str).encode("utf-8") + self.send_response(status) + self.send_header("Content-Type", "application/json; charset=utf-8") + self.send_header("Content-Length", str(len(body))) + self.end_headers() + self.wfile.write(body) + + def _ok(self, data: Any, status: int = 200) -> None: + self._send_json(status, {"code": 0, "message": "ok", "data": data}) + + def _error(self, code: int, message: str, status: int) -> None: + self._send_json(status, {"code": code, "message": message, "data": None}) + + def _read_body(self) -> Dict[str, Any]: + length = int(self.headers.get("Content-Length", "0") or "0") + if length <= 0: + return {} + raw = self.rfile.read(length) + try: + data = json.loads(raw.decode("utf-8")) + except (json.JSONDecodeError, UnicodeDecodeError): + return {} + return data if isinstance(data, dict) else {} + + def _query(self) -> Dict[str, str]: + parsed = urlparse(self.path) + qs = parse_qs(parsed.query) + return {k: v[0] for k, v in qs.items() if v} + + def log_message(self, fmt: str, *args: Any) -> None: # noqa: A002 + # 静默访问日志,避免污染测试输出 + return + + # ------------------------------------------------------------------ + # 路由 + # ------------------------------------------------------------------ + def do_GET(self) -> None: # noqa: N802 + parsed = urlparse(self.path) + path = parsed.path.rstrip("/") or "/" + + if path == "/healthz": + return self._send_json(200, {"status": "ok"}) + + if path == "/api/v1/fault-logs": + return self._list_fault_logs() + + if path.startswith("/api/v1/fault-logs/"): + rest = path[len("/api/v1/fault-logs/") :] + if rest.endswith("/root-cause"): + return self._get_root_cause(rest[: -len("/root-cause")]) + return self._get_fault_log(rest) + + if path == "/api/v1/fault-types": + return self._list_fault_types() + + if path == "/api/v1/fault-filters": + return self._list_fault_filters() + + return self._error(40401, "not found", 404) + + def do_POST(self) -> None: # noqa: N802 + parsed = urlparse(self.path) + path = parsed.path.rstrip("/") or "/" + + if path == "/api/v1/fault-types": + return self._create_fault_type() + + if path == "/api/v1/fault-filters": + return self._create_fault_filter() + + return self._error(40401, "not found", 404) + + # ------------------------------------------------------------------ + # 处理器 + # ------------------------------------------------------------------ + def _list_fault_logs(self) -> None: + q = self._query() + try: + page = max(1, int(q.get("page", "1"))) + page_size = max(1, min(int(q.get("page_size", "20")), 200)) + except ValueError: + return self._error(40001, "invalid param: page/page_size", 400) + + result = self.storage.query_fault_logs( + host_id=q.get("host_id") or None, + fault_type=q.get("fault_type") or None, + level=q.get("level") or None, + keyword=q.get("keyword") or None, + page=page, + page_size=page_size, + ) + items = [l.to_dict() for l in result["items"]] + return self._ok({"total": result["total"], "items": items}) + + def _get_fault_log(self, fault_log_id: str) -> None: + log = self.storage.get_fault_log(fault_log_id) + if log is None: + return self._error(40401, "fault log not found", 404) + return self._ok(log.to_dict()) + + def _get_root_cause(self, fault_log_id: str) -> None: + root_cause = self.storage.get_root_cause(fault_log_id) + if root_cause is None: + return self._error(40401, "root cause not found", 404) + return self._ok(root_cause.to_dict()) + + def _list_fault_types(self) -> None: + return self._ok({"items": [ft.to_dict() for ft in self.storage.list_fault_types()]}) + + def _create_fault_type(self) -> None: + body = self._read_body() + if not body.get("fault_type") or not body.get("name"): + return self._error(40002, "fault_type and name are required", 400) + try: + ft = FaultType.from_dict(body) + except (TypeError, ValueError): + return self._error(40001, "invalid param", 400) + self.storage.add_fault_type(ft) + return self._ok(ft.to_dict(), status=201) + + def _list_fault_filters(self) -> None: + return self._ok({"items": [r.to_dict() for r in self.storage.list_filter_rules()]}) + + def _create_fault_filter(self) -> None: + body = self._read_body() + if not body.get("name"): + return self._error(40002, "name is required", 400) + try: + rule = FaultFilterRule.from_dict(body) + except (TypeError, ValueError): + return self._error(40001, "invalid param", 400) + self.storage.add_filter_rule(rule) + return self._ok(rule.to_dict(), status=201) + + +def create_handler(storage: Storage) -> Type[FaultLogAPIHandler]: + """创建绑定存储的处理器类。""" + + class Handler(FaultLogAPIHandler): + pass + + Handler.storage = storage + return Handler + + +def run_server(host: str, port: int, storage: Storage) -> None: + """启动 HTTP 服务(阻塞)。""" + handler = create_handler(storage) + server = ThreadingHTTPServer((host, port), handler) + try: + server.serve_forever() + except KeyboardInterrupt: + pass + finally: + server.server_close() diff --git a/src/fault_log_analyzer/classifier.py b/src/fault_log_analyzer/classifier.py index 3e10fbc..8469fb1 100644 --- a/src/fault_log_analyzer/classifier.py +++ b/src/fault_log_analyzer/classifier.py @@ -1 +1,93 @@ -fault_type 归类(已提交至 git) \ 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 b229e89..843d03e 100644 --- a/src/fault_log_analyzer/cluster.py +++ b/src/fault_log_analyzer/cluster.py @@ -1 +1,178 @@ -TF-IDF + DBSCAN 聚类(已提交至 git) \ 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 ddd7226..1227687 100644 --- a/src/fault_log_analyzer/config.py +++ b/src/fault_log_analyzer/config.py @@ -1 +1,110 @@ -环境变量配置(已提交至 git) \ 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 03b7cb0..5402985 100644 --- a/src/fault_log_analyzer/filters.py +++ b/src/fault_log_analyzer/filters.py @@ -1 +1,54 @@ -FaultFilter 过滤规则(已提交至 git) \ 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 feffeaa..5d6c5c4 100644 --- a/src/fault_log_analyzer/fingerprint.py +++ b/src/fault_log_analyzer/fingerprint.py @@ -1 +1,89 @@ -MinHash 指纹 + Deduplicator(已提交至 git) \ 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/integrations.py b/src/fault_log_analyzer/integrations.py index a9003e7..68526df 100644 --- a/src/fault_log_analyzer/integrations.py +++ b/src/fault_log_analyzer/integrations.py @@ -1 +1,199 @@ -Kafka/ES/MySQL/Redis 适配(已提交至 git) \ No newline at end of file +"""可选真实后端适配(Kafka / Elasticsearch / MySQL / Redis)。 + +所有外部 SDK 均为惰性导入:未安装对应依赖时,仅在实例化时抛出明确异常, +不影响核心逻辑与单元测试运行。 + +- KafkaProducerAdapter / KafkaConsumerAdapter -> kafka-python +- ElasticsearchSink -> elasticsearch +- MySQLFaultRepository -> pymysql +- RedisDedupCache -> redis +""" + +from __future__ import annotations + +import json +from datetime import datetime +from typing import Any, Dict, Iterator, List, Optional + +from .models import Event, FaultFilterRule, FaultLog, FaultType, RootCause +from .storage import Storage + + +# --------------------------------------------------------------------------- +# Kafka +# --------------------------------------------------------------------------- +class KafkaProducerAdapter: + """Kafka 生产者封装(用于生产 logs.fault)。""" + + def __init__(self, bootstrap_servers: str) -> None: + try: + from kafka import KafkaProducer # type: ignore + except ImportError as exc: # pragma: no cover - 依赖未安装 + raise RuntimeError("Kafka 后端需要安装 kafka-python") from exc + self._producer = KafkaProducer( + bootstrap_servers=bootstrap_servers, + value_serializer=lambda v: json.dumps(v, ensure_ascii=False, default=str).encode("utf-8"), + ) + + def send(self, topic: str, value: Dict[str, Any]) -> None: + self._producer.send(topic, value) + + def flush(self) -> None: + self._producer.flush() + + def close(self) -> None: + self._producer.close() + + +class KafkaConsumerAdapter: + """Kafka 消费者封装(用于消费 logs.raw)。""" + + def __init__(self, bootstrap_servers: str, group_id: str, topic: str) -> None: + try: + from kafka import KafkaConsumer # type: ignore + except ImportError as exc: # pragma: no cover - 依赖未安装 + raise RuntimeError("Kafka 后端需要安装 kafka-python") from exc + self._consumer = KafkaConsumer( + topic, + bootstrap_servers=bootstrap_servers, + group_id=group_id, + value_deserializer=lambda raw: json.loads(raw.decode("utf-8")), + auto_offset_reset="earliest", + ) + + def poll(self, timeout_ms: int = 1000) -> List[Dict[str, Any]]: + batch = self._consumer.poll(timeout_ms=timeout_ms) + return [record.value for records in batch.values() for record in records] + + def close(self) -> None: + self._consumer.close() + + +# --------------------------------------------------------------------------- +# Elasticsearch +# --------------------------------------------------------------------------- +class ElasticsearchSink: + """故障日志写入 Elasticsearch(索引 hms-fault-log-{yyyy.MM})。""" + + def __init__(self, hosts: str, index_prefix: str = "hms-fault-log") -> None: + try: + from elasticsearch import Elasticsearch # type: ignore + except ImportError as exc: # pragma: no cover - 依赖未安装 + raise RuntimeError("Elasticsearch 后端需要安装 elasticsearch") from exc + host_list = [h.strip() for h in hosts.split(",") if h.strip()] + self._client = Elasticsearch(host_list or ["http://localhost:9200"]) + self.index_prefix = index_prefix + + @staticmethod + def _index_name(prefix: str, occurred_at: Optional[datetime]) -> str: + dt = occurred_at or datetime.utcnow() + return f"{prefix}-{dt:%Y.%m}" + + def write(self, log: FaultLog) -> None: + index = self._index_name(self.index_prefix, log.occurred_at) + self._client.index(index=index, id=log.fault_log_id, document=log.to_dict()) + + def close(self) -> None: + self._client.close() + + +# --------------------------------------------------------------------------- +# MySQL +# --------------------------------------------------------------------------- +class MySQLFaultRepository: + """故障日志 / 根因结果 MySQL 落库(简化子集,生产可按需扩展)。""" + + def __init__(self, host: str, port: int, user: str, password: str, db: str) -> None: + try: + import pymysql # type: ignore + except ImportError as exc: # pragma: no cover - 依赖未安装 + raise RuntimeError("MySQL 后端需要安装 pymysql") from exc + self._conn = pymysql.connect( + host=host, port=port, user=user, password=password, database=db + ) + + def save_fault_log(self, log: FaultLog) -> None: + sql = ( + "INSERT INTO fault_log (fault_log_id, host_id, fault_type, cluster_id, " + "fingerprint, level, service, message, trace_id, occurred_at, count) " + "VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) " + "ON DUPLICATE KEY UPDATE fault_type=VALUES(fault_type), " + "cluster_id=VALUES(cluster_id), count=VALUES(count)" + ) + with self._conn.cursor() as cursor: + cursor.execute( + sql, + ( + log.fault_log_id, + log.host_id, + log.fault_type, + log.cluster_id, + log.fingerprint, + log.level, + log.service, + log.message, + log.trace_id, + log.occurred_at, + log.count, + ), + ) + self._conn.commit() + + def save_root_cause(self, root_cause: RootCause) -> None: + sql = ( + "INSERT INTO root_cause (fault_log_id, cause_type, evidence, confidence, analysis_at) " + "VALUES (%s, %s, %s, %s, %s) " + "ON DUPLICATE KEY UPDATE cause_type=VALUES(cause_type), " + "evidence=VALUES(evidence), confidence=VALUES(confidence)" + ) + with self._conn.cursor() as cursor: + cursor.execute( + sql, + ( + root_cause.fault_log_id, + root_cause.cause_type, + json.dumps(root_cause.evidence, ensure_ascii=False, default=str), + root_cause.confidence, + root_cause.analysis_at, + ), + ) + self._conn.commit() + + def close(self) -> None: + self._conn.close() + + +# --------------------------------------------------------------------------- +# Redis +# --------------------------------------------------------------------------- +class RedisDedupCache: + """基于 Redis 的指纹去重(对应 Redis 键 hms:log:fp:{fingerprint})。""" + + def __init__(self, url: str, window_seconds: int = 300) -> None: + try: + import redis # type: ignore + except ImportError as exc: # pragma: no cover - 依赖未安装 + raise RuntimeError("Redis 后端需要安装 redis") from exc + self._redis = redis.Redis.from_url(url) + self.window_seconds = window_seconds + + def is_duplicate(self, fingerprint: str) -> bool: + key = f"hms:log:fp:{fingerprint}" + # SET NX EX:成功写入返回 True 表示首次出现;失败表示窗口内重复 + return not bool(self._redis.set(key, "1", nx=True, ex=self.window_seconds)) + + def mark(self, fingerprint: str) -> None: + key = f"hms:log:fp:{fingerprint}" + self._redis.set(key, "1", ex=self.window_seconds) + + def close(self) -> None: + self._redis.close() + + +__all__ = [ + "KafkaProducerAdapter", + "KafkaConsumerAdapter", + "ElasticsearchSink", + "MySQLFaultRepository", + "RedisDedupCache", +] diff --git a/src/fault_log_analyzer/models.py b/src/fault_log_analyzer/models.py index 324cbc8..e134663 100644 --- a/src/fault_log_analyzer/models.py +++ b/src/fault_log_analyzer/models.py @@ -1 +1,183 @@ -数据模型 FaultLog/RootCause/FaultType/FaultFilterRule/Event/ParsedLog(已提交至 git) \ 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。 + + 支持 ISO8601(带 Z 或空格分隔)、Unix 时间戳以及 syslog 风格 + ``Jan 1 10:18:00``(无年份,缺省使用当前年份)。 + """ + 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" + + dt: Optional[datetime] = None + try: + dt = datetime.fromisoformat(text) + except ValueError: + try: + dt = datetime.fromisoformat(text.replace(" ", "T")) + except ValueError: + # syslog 风格:Jan 1 10:18:00(无年份,缺省当前年份) + normalized = " ".join(text.split()) + dt = datetime.strptime(normalized, "%b %d %H:%M:%S") + dt = dt.replace(year=datetime.now().year) + + 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 c9eb679..a4d7003 100644 --- a/src/fault_log_analyzer/parser.py +++ b/src/fault_log_analyzer/parser.py @@ -1 +1,165 @@ -结构化解析 + 模板化 + 分词(已提交至 git) \ 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_PAT = r"\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*(?P" + _TS_PAT + r")\s+" + r"(?P[A-Za-z]+)\s+" + r"(?:\[(?P[^\]]*)\]\s*)?" + r"(?:-\s*)?" + r"(?P.*)$" +) + +_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: + level = m.group("level").upper() + thread = (m.group("thread") or "").strip() + message = (m.group("message") or "").strip() + return ParsedLog( + timestamp=_to_dt(m.group("ts")), + host_id=self._extract_host(message) or thread, + service=thread, + level=level, + message=message, + 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/pipeline.py b/src/fault_log_analyzer/pipeline.py index 68a727a..800bfa2 100644 --- a/src/fault_log_analyzer/pipeline.py +++ b/src/fault_log_analyzer/pipeline.py @@ -1 +1,97 @@ -捕获管道编排(已提交至 git) \ No newline at end of file +"""故障日志捕获管道。 + +依据 architecture.md 5.3.2: +``` +logs.raw → 级别过滤(ERROR/FATAL) → 关键字/正则过滤 → 结构化解析 + → 去重打标 → Kafka(logs.fault) → Elasticsearch(hms-fault-log-*) +``` + +本模块把「过滤 → 解析 → 指纹去重 → 启发式归类 → 落库 → 可选生产」串成 +一个可独立测试的管道。聚类与根因分析由 workers.py 异步完成。 +""" + +from __future__ import annotations + +import uuid +from typing import Any, Callable, List, Optional + +from .classifier import Classifier +from .filters import FaultFilter +from .fingerprint import Deduplicator, MinHash +from .models import FaultLog +from .parser import LogParser +from .storage import Storage + + +class CapturePipeline: + """故障日志捕获管道。 + + 参数: + storage: 故障日志存储。 + capture_filter: 过滤规则,默认 FaultFilter(ERROR/FATAL)。 + parser: 日志解析器,默认 LogParser。 + classifier: 启发式归类器,默认 Classifier。 + dedup: 去重器(实现 is_duplicate 接口),默认进程内 Deduplicator。 + minhash: 指纹器,默认 MinHash。 + producer: 可选回调,捕获成功后调用(生产到 logs.fault / ES)。 + """ + + def __init__( + self, + storage: Storage, + capture_filter: Optional[FaultFilter] = None, + parser: Optional[LogParser] = None, + classifier: Optional[Classifier] = None, + dedup: Optional[Deduplicator] = None, + minhash: Optional[MinHash] = None, + producer: Optional[Callable[[FaultLog], None]] = None, + ) -> None: + self.storage = storage + self.capture_filter = capture_filter or FaultFilter() + self.parser = parser or LogParser() + self.classifier = classifier or Classifier() + self.dedup = dedup or Deduplicator() + self.minhash = minhash or MinHash() + self.producer = producer + self.captured = 0 + self.duplicates = 0 + + def process(self, raw: Any) -> Optional[FaultLog]: + """处理一条原始日志;不是故障日志或重复时返回 None。""" + parsed = self.parser.parse(raw) + + if not self.capture_filter.should_capture(parsed): + return None + + fingerprint = self.minhash.fingerprint(parsed.message or "") + if self.dedup.is_duplicate(fingerprint): + self.duplicates += 1 + return None + + fault_type = self.classifier.classify(parsed.message or "") + log = FaultLog( + fault_log_id=uuid.uuid4().hex, + host_id=parsed.host_id or "", + message=parsed.message or "", + level=parsed.level or "ERROR", + fingerprint=fingerprint, + fault_type=fault_type, + service=parsed.service or None, + trace_id=parsed.trace_id, + occurred_at=parsed.timestamp, + count=1, + ) + self.storage.save_fault_log(log) + if self.producer is not None: + self.producer(log) + self.captured += 1 + return log + + def process_batch(self, raws: List[Any]) -> List[FaultLog]: + """批量处理原始日志,返回捕获到的 FaultLog 列表。""" + results: List[FaultLog] = [] + for raw in raws: + log = self.process(raw) + if log is not None: + results.append(log) + return results diff --git a/src/fault_log_analyzer/root_cause.py b/src/fault_log_analyzer/root_cause.py index b0fffbb..aaada3f 100644 --- a/src/fault_log_analyzer/root_cause.py +++ b/src/fault_log_analyzer/root_cause.py @@ -1 +1,150 @@ -根因分析器(已提交至 git) \ No newline at end of file +"""根因分析(RootCauseAnalyzer)。 + +依据 architecture.md 5.3.4 与 database-design.md root_cause 表: + +- 输入:故障日志(FaultLog)+ 同主机近时间窗口的指标事件(Event)。 +- 策略:规则优先、统计辅助、置信度评分。 + 1. 日志关键字启发式(复用 Classifier.hint)给出候选故障类型。 + 2. 指标-故障规则库:如 disk_used_percent >= 90 -> disk_full。 + 3. 时间窗口 + 主机 + 事件指标关联,匹配则提升置信度。 +- 输出:RootCause{fault_log_id, cause_type, evidence[], confidence}。 + +时间比较统一经过 ``_as_aware`` 归一化,兼容 naive / aware datetime, +避免直接比较时抛出 TypeError(见 tests/test_root_cause.py 鲁棒性用例)。 +""" + +from __future__ import annotations + +import re +from datetime import datetime, timedelta, timezone +from typing import Dict, List, Optional, Sequence, Tuple + +from .classifier import Classifier +from .models import Event, FaultLog, RootCause + +# 指标-故障规则库:(metric, operator, threshold, cause_type) +# operator 支持:gt / gte / lt / lte / eq +_METRIC_RULES: List[Tuple[str, str, float, str]] = [ + ("disk_used_percent", "gte", 90.0, "disk_full"), + ("mem_used_percent", "gte", 90.0, "memory_exhausted"), + ("cpu_usage", "gte", 95.0, "cpu_saturation"), + ("load_1m", "gte", 20.0, "cpu_saturation"), + ("net_pkt_drop", "gte", 100.0, "network_unreachable"), +] + +# OOM 日志关键字(补充 Classifier.hint 已覆盖之外的显式规则,保持可读性) +_OOM_RE = re.compile(r"oomkilled|out of memory|oom", re.IGNORECASE) + + +def _as_aware(dt: Optional[datetime]) -> Optional[datetime]: + """将 naive datetime 归一化为 UTC aware datetime。 + + aware datetime 原样返回;naive datetime 视为 UTC。 + """ + if dt is None: + return None + if dt.tzinfo is None: + return dt.replace(tzinfo=timezone.utc) + return dt + + +def _compare(value: float, operator: str, threshold: float) -> bool: + if operator == "gt": + return value > threshold + if operator == "gte": + return value >= threshold + if operator == "lt": + return value < threshold + if operator == "lte": + return value <= threshold + if operator == "eq": + return value == threshold + return False + + +def _hint_cause(message: str) -> Optional[str]: + """日志消息 -> 候选故障类型(复用分类器启发式规则)。""" + return Classifier.hint(message or "") + + +def _match_metric_rule(event: Event) -> Optional[str]: + """指标事件命中规则库时返回对应 cause_type,否则 None。""" + for metric, operator, threshold, cause_type in _METRIC_RULES: + if event.metric == metric and _compare(event.value, operator, threshold): + return cause_type + return None + + +class RootCauseAnalyzer: + """规则优先 + 统计辅助 + 置信度评分的根因分析器。""" + + def __init__(self, window_minutes: int = 5): + self.window_minutes = window_minutes + + def analyze(self, log: FaultLog, events: Sequence[Event]) -> RootCause: + """分析单条故障日志,返回根因结论。 + + 置信度规则(确定性、可测试): + - 日志关键字命中 +0.6 + - 指标规则命中 +0.25 + - 二者同时命中且结论一致 +0.15(互相印证) + """ + hint_cause = _hint_cause(log.message) + event_cause, event_evidence = self._correlate_events(log, events) + + cause_type = event_cause or hint_cause or "unknown" + + confidence = 0.0 + if hint_cause: + confidence += 0.6 + if event_cause: + confidence += 0.25 + if hint_cause and event_cause and hint_cause == event_cause: + confidence += 0.15 + + evidence: List[Dict] = [{"type": "log", "message": log.message}] + evidence.extend(event_evidence) + + return RootCause( + fault_log_id=log.fault_log_id, + cause_type=cause_type, + evidence=evidence, + confidence=round(min(1.0, confidence), 2), + analysis_at=datetime.now(timezone.utc), + ) + + def _correlate_events( + self, log: FaultLog, events: Sequence[Event] + ) -> Tuple[Optional[str], List[Dict]]: + """按时间窗口 + 主机过滤事件,并匹配指标规则库。""" + window = timedelta(minutes=self.window_minutes) + log_time = _as_aware(log.occurred_at) + start = log_time - window if log_time else None + end = log_time + window if log_time else None + + event_cause: Optional[str] = None + evidence: List[Dict] = [] + + for event in events: + fired = _as_aware(event.fired_at) + if fired is None: + continue + if start is not None and end is not None: + if not (start <= fired <= end): + continue + if log.host_id and event.host_id and event.host_id != log.host_id: + continue + + evidence.append( + { + "type": "event", + "event_id": event.event_id, + "metric": event.metric, + "value": event.value, + } + ) + matched = _match_metric_rule(event) + if matched and event_cause is None: + event_cause = matched + + return event_cause, evidence diff --git a/src/fault_log_analyzer/storage.py b/src/fault_log_analyzer/storage.py index bb247ab..1d34cd6 100644 --- a/src/fault_log_analyzer/storage.py +++ b/src/fault_log_analyzer/storage.py @@ -1 +1,180 @@ -存储抽象 + 内存实现(已提交至 git) \ No newline at end of file +"""存储抽象与内存实现。 + +依据 docs/01-design/database-design.md 中 fault_log / root_cause / +fault_type / fault_filter_rule / event 表结构,提供统一的进程内存储接口。 +生产环境通过 integrations.py 对接 MySQL / Elasticsearch / Redis; +本模块保证零第三方依赖,测试与离线运行使用 MemoryStorage。 +""" + +from __future__ import annotations + +import threading +from datetime import datetime +from typing import Dict, List, Optional + +from .models import Event, FaultFilterRule, FaultLog, FaultType, RootCause + + +class Storage: + """存储接口(抽象基类)。 + + 设计上把五类实体(故障日志 / 根因 / 故障类型 / 过滤规则 / 事件)收敛到 + 一个 Storage 门面,便于 pipeline / workers / api 复用同一份存储。 + """ + + # ---- 故障日志 ------------------------------------------------------ + def save_fault_log(self, log: FaultLog) -> None: + raise NotImplementedError + + def get_fault_log(self, fault_log_id: str) -> Optional[FaultLog]: + raise NotImplementedError + + def list_fault_logs(self) -> List[FaultLog]: + raise NotImplementedError + + def query_fault_logs( + self, + host_id: Optional[str] = None, + fault_type: Optional[str] = None, + level: Optional[str] = None, + keyword: Optional[str] = None, + page: int = 1, + page_size: int = 20, + ) -> Dict: + raise NotImplementedError + + # ---- 根因 ---------------------------------------------------------- + def save_root_cause(self, root_cause: RootCause) -> None: + raise NotImplementedError + + def get_root_cause(self, fault_log_id: str) -> Optional[RootCause]: + raise NotImplementedError + + # ---- 故障类型 ------------------------------------------------------ + def list_fault_types(self) -> List[FaultType]: + raise NotImplementedError + + def get_fault_type(self, fault_type: str) -> Optional[FaultType]: + raise NotImplementedError + + def add_fault_type(self, fault_type: FaultType) -> FaultType: + raise NotImplementedError + + # ---- 过滤规则 ------------------------------------------------------ + def list_filter_rules(self) -> List[FaultFilterRule]: + raise NotImplementedError + + def add_filter_rule(self, rule: FaultFilterRule) -> FaultFilterRule: + raise NotImplementedError + + # ---- 事件(根因分析关联输入) -------------------------------------- + def save_event(self, event: Event) -> None: + raise NotImplementedError + + def list_events_by_host( + self, host_id: str, start: datetime, end: datetime + ) -> List[Event]: + raise NotImplementedError + + +class MemoryStorage(Storage): + """进程内内存实现,用于单元测试与离线运行。""" + + def __init__(self) -> None: + self._lock = threading.RLock() + self._fault_logs: Dict[str, FaultLog] = {} + self._root_causes: Dict[str, RootCause] = {} + self._fault_types: Dict[str, FaultType] = {} + self._filter_rules: Dict[str, FaultFilterRule] = {} + self._events: List[Event] = [] + + # ---- 故障日志 ------------------------------------------------------ + def save_fault_log(self, log: FaultLog) -> None: + with self._lock: + self._fault_logs[log.fault_log_id] = log + + def get_fault_log(self, fault_log_id: str) -> Optional[FaultLog]: + with self._lock: + return self._fault_logs.get(fault_log_id) + + def list_fault_logs(self) -> List[FaultLog]: + with self._lock: + return sorted( + self._fault_logs.values(), + key=lambda l: (l.occurred_at, l.fault_log_id), + ) + + def query_fault_logs( + self, + host_id: Optional[str] = None, + fault_type: Optional[str] = None, + level: Optional[str] = None, + keyword: Optional[str] = None, + page: int = 1, + page_size: int = 20, + ) -> Dict: + with self._lock: + items = list(self._fault_logs.values()) + if host_id: + items = [l for l in items if l.host_id == host_id] + if fault_type: + items = [l for l in items if l.fault_type == fault_type] + if level: + items = [l for l in items if l.level.upper() == level.upper()] + if keyword: + needle = keyword.lower() + items = [l for l in items if needle in (l.message or "").lower()] + items.sort(key=lambda l: (l.occurred_at, l.fault_log_id), reverse=True) + total = len(items) + page = max(1, page) + page_size = max(1, min(page_size, 200)) + start = (page - 1) * page_size + return {"total": total, "items": items[start : start + page_size]} + + # ---- 根因 ---------------------------------------------------------- + def save_root_cause(self, root_cause: RootCause) -> None: + with self._lock: + self._root_causes[root_cause.fault_log_id] = root_cause + + def get_root_cause(self, fault_log_id: str) -> Optional[RootCause]: + with self._lock: + return self._root_causes.get(fault_log_id) + + # ---- 故障类型 ------------------------------------------------------ + def list_fault_types(self) -> List[FaultType]: + with self._lock: + return list(self._fault_types.values()) + + def get_fault_type(self, fault_type: str) -> Optional[FaultType]: + with self._lock: + return self._fault_types.get(fault_type) + + def add_fault_type(self, fault_type: FaultType) -> FaultType: + with self._lock: + self._fault_types[fault_type.fault_type] = fault_type + return fault_type + + # ---- 过滤规则 ------------------------------------------------------ + def list_filter_rules(self) -> List[FaultFilterRule]: + with self._lock: + return list(self._filter_rules.values()) + + def add_filter_rule(self, rule: FaultFilterRule) -> FaultFilterRule: + with self._lock: + self._filter_rules[rule.name] = rule + return rule + + # ---- 事件 ---------------------------------------------------------- + def save_event(self, event: Event) -> None: + with self._lock: + self._events.append(event) + + def list_events_by_host( + self, host_id: str, start: datetime, end: datetime + ) -> List[Event]: + with self._lock: + return [ + e + for e in self._events + if e.host_id == host_id and start <= e.fired_at <= end + ] diff --git a/src/fault_log_analyzer/workers.py b/src/fault_log_analyzer/workers.py index dbc94ce..c765f56 100644 --- a/src/fault_log_analyzer/workers.py +++ b/src/fault_log_analyzer/workers.py @@ -1 +1,108 @@ -聚类/根因分析 worker(已提交至 git) \ No newline at end of file +"""聚类与根因分析 worker。 + +依据 architecture.md 5.3.3 / 5.3.4: +- ClusterWorker:批量对未归类故障日志做 TF-IDF + DBSCAN 聚类,映射到 fault_type。 +- RootCauseWorker:批量关联同主机近期指标事件,产出 root_cause。 +""" + +from __future__ import annotations + +from datetime import timedelta +from typing import Dict, List, Optional, Sequence + +from .classifier import Classifier +from .cluster import Clusterer +from .models import FaultLog, RootCause +from .root_cause import RootCauseAnalyzer +from .storage import Storage + + +def _cluster_label(cluster_id: int) -> str: + return f"c-{cluster_id}" + + +class ClusterWorker: + """批量聚类并归类故障日志。""" + + def __init__( + self, + storage: Storage, + clusterer: Optional[Clusterer] = None, + classifier: Optional[Classifier] = None, + ) -> None: + self.storage = storage + self.clusterer = clusterer or Clusterer() + self.classifier = classifier or Classifier() + + def run(self, logs: Optional[Sequence[FaultLog]] = None) -> Dict: + """对故障日志聚类并持久化归类结果。 + + 返回:{"clustered": n, "clusters": m, "noise": k} + """ + items: List[FaultLog] = list(logs) if logs is not None else self.storage.list_fault_logs() + if not items: + return {"clustered": 0, "clusters": 0, "noise": 0} + + messages = [l.message or "" for l in items] + labels = self.clusterer.fit_predict(messages) + + # 按簇收集代表消息,用于簇级归类 + cluster_messages: Dict[int, List[str]] = {} + for label, message in zip(labels, messages): + if label >= 0: + cluster_messages.setdefault(label, []).append(message) + + fault_type_by_cluster: Dict[int, Optional[str]] = {} + for label, msgs in cluster_messages.items(): + fault_type_by_cluster[label] = self.classifier.classify_cluster(label, msgs) + + clustered = 0 + for log, label in zip(items, labels): + if label >= 0: + log.cluster_id = _cluster_label(label) + log.fault_type = fault_type_by_cluster[label] or log.fault_type + clustered += 1 + self.storage.save_fault_log(log) + + noise = sum(1 for lbl in labels if lbl < 0) + return { + "clustered": clustered, + "clusters": self.clusterer.n_clusters, + "noise": noise, + } + + +class RootCauseWorker: + """批量执行根因分析并持久化结论。""" + + def __init__( + self, + storage: Storage, + analyzer: Optional[RootCauseAnalyzer] = None, + window_minutes: int = 5, + ) -> None: + self.storage = storage + self.window_minutes = window_minutes + self.analyzer = analyzer or RootCauseAnalyzer(window_minutes=window_minutes) + + def run(self, logs: Optional[Sequence[FaultLog]] = None) -> List[RootCause]: + """对尚未分析过的故障日志执行根因分析,返回新产出的结论列表。""" + items: List[FaultLog] = list(logs) if logs is not None else self.storage.list_fault_logs() + results: List[RootCause] = [] + window = timedelta(minutes=self.window_minutes) + + for log in items: + if self.storage.get_root_cause(log.fault_log_id) is not None: + continue + start = log.occurred_at - window + end = log.occurred_at + window + events = ( + self.storage.list_events_by_host(log.host_id, start, end) + if log.host_id + else [] + ) + root_cause = self.analyzer.analyze(log, events) + self.storage.save_root_cause(root_cause) + results.append(root_cause) + + return results