develop: 重做 fault-log-analyzer 分析模块(遵守模块开发规范)
This commit is contained in:
parent
178167ef9d
commit
661fa2f02b
@ -1 +1,8 @@
|
||||
(已提交至远程 main 分支)
|
||||
"""fault-log-analyzer:HMS 故障日志捕获与分析模块。
|
||||
|
||||
提供故障日志捕获管道、归类聚类、根因分析以及 REST API。
|
||||
核心逻辑零强制第三方依赖(纯 Python 标准库实现),生产环境可
|
||||
按需接入 Kafka / Elasticsearch / MySQL / Redis 等真实后端。
|
||||
"""
|
||||
|
||||
__version__ = "0.1.0"
|
||||
|
||||
@ -1 +1,93 @@
|
||||
(已提交至远程 main 分支)
|
||||
"""簇 -> 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
|
||||
|
||||
@ -1 +1,178 @@
|
||||
(已提交至远程 main 分支)
|
||||
"""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
|
||||
|
||||
@ -1 +1,110 @@
|
||||
(已提交至远程 main 分支)
|
||||
"""配置模型与环境变量加载。
|
||||
|
||||
零第三方依赖,通过环境变量(带合理默认值)与可选显式参数配置。
|
||||
"""
|
||||
|
||||
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)
|
||||
|
||||
@ -1 +1,54 @@
|
||||
(已提交至远程 main 分支)
|
||||
"""故障日志过滤规则。
|
||||
|
||||
依据 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]
|
||||
|
||||
@ -1 +1,89 @@
|
||||
(已提交至远程 main 分支)
|
||||
"""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()
|
||||
|
||||
@ -1 +1,169 @@
|
||||
(已提交至远程 main 分支)
|
||||
"""数据模型定义。
|
||||
|
||||
依据 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
|
||||
|
||||
@ -1 +1,164 @@
|
||||
(已提交至远程 main 分支)
|
||||
"""日志结构化解析与模板化。
|
||||
|
||||
支持常见日志形态:
|
||||
- JSON 对象(含 timestamp/level/message/host_id/service/trace_id)
|
||||
- log4j 风格:`2025-01-01 10:18:00 ERROR [app] message`
|
||||
- syslog 风格:`<PRI>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<ts>\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<pri>\d{1,3})>\d?\s*(?P<ts>\w{3}\s+\d{1,2}\s+\d{2}:\d{2}:\d{2})\s+"
|
||||
r"(?P<host>\S+)\s+(?P<app>\S+?)(?:\[(?P<pid>\d+)\])?:\s*(?P<message>.*)$"
|
||||
)
|
||||
|
||||
|
||||
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", "<IP>", 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", "<UUID>", text)
|
||||
text = re.sub(r"\b0x[0-9a-fA-F]+\b", "<HEX>", text)
|
||||
text = re.sub(r"\b\d{10,}\b", "<LONG>", text)
|
||||
text = re.sub(r"\b\d+(?:\.\d+)?\b", "<NUM>", 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]
|
||||
|
||||
@ -1 +1,124 @@
|
||||
(已提交至远程 main 分支)
|
||||
"""根因分析(辅助性)。
|
||||
|
||||
依据 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,
|
||||
)
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user