
如果说AIOps 1.0时代我们做到了“看”(监控可视化)和“知”(告警通知),那么AIOps 2.0时代(2025-2027)的核心命题是“动”——让系统在故障发生时自动完成从发现到修复的全流程。
核心愿景:打造一个运维“自动驾驶系统”,让80%的已知故障在无人介入的情况下自动闭环,20%的未知故障由AI辅助人工快速决策。
本文交付物:一套完整的自动化故障处置平台架构设计方案 + 核心模块实现代码 + 落地路线图。
层级 | 组件 | 选型 | 理由 |
|---|---|---|---|
数据总线 | 消息队列 | Apache Kafka | 高吞吐、持久化、支持回放 |
时序存储 | 监控数据库 | VictoriaMetrics | 兼容PromQL,存储成本低 |
日志存储 | 日志数据库 | Elasticsearch 8.x | 全文检索 + 聚合分析能力强 |
图谱存储 | 知识图谱 | Neo4j | 服务依赖关系、影响面分析 |
AI编排 | 大模型应用框架 | LangChain4j / LangGraph | 支持复杂Agent工作流 |
模型服务 | LLM网关 | 自研 + DeepSeek API | 混合部署(云端+私有化) |
执行引擎 | 自动化工具 | Terraform + Ansible + K8s API | 声明式配置,幂等性好 |
工作流引擎 | 故障处置编排 | Temporal / Cadence | 支持长时运行、重试、补偿 |
// 统一事件模型
@Data
@Builder
public class OpsEvent {
private String eventId; // 全局唯一ID
private EventType type; // METRIC / LOG / TRACE / CHANGE
private String source; // 来源系统(如prometheus-01)
private String cluster; // 集群名称
private String service; // 服务名称
private String instance; // 实例IP
private Map<String, String> tags; // 标签(如 env=prod, region=us-east)
private String content; // 事件内容(JSON格式)
private String rawData; // 原始数据
private Instant timestamp; // 事件发生时间
private Instant receiveTime; // 平台接收时间
}
// 事件标准化处理器
@Component
public class EventNormalizer {
public OpsEvent normalize(Object source) {
if (source instanceof PrometheusAlert) {
PrometheusAlert alert = (PrometheusAlert) source;
return OpsEvent.builder()
.eventId(UUID.randomUUID().toString())
.type(EventType.METRIC)
.source("prometheus")
.service(alert.getLabels().get("service"))
.tags(alert.getLabels())
.content(alert.getAnnotations().get("summary"))
.timestamp(alert.getStartsAt())
.build();
}
// ELK日志、JaegerTrace、ChangeRecord等其他类型处理...
throw new UnsupportedOperationException("Unknown source type");
}
}# 基于K8s API构建服务依赖图
from kubernetes import client, config
import networkx as nx
class ServiceTopologyBuilder:
def __init__(self):
config.load_incluster_config()
self.core_v1 = client.CoreV1Api()
self.network_v1 = client.NetworkingV1Api()
def build_topology(self):
G = nx.DiGraph()
# 1. 获取所有Service
services = self.core_v1.list_service_for_all_namespaces()
for svc in services.items:
G.add_node(svc.metadata.name,
type='service',
namespace=svc.metadata.namespace,
cluster_ip=svc.spec.cluster_ip,
labels=svc.metadata.labels)
# 2. 获取所有Pod,建立Service->Pod关系
pods = self.core_v1.list_pod_for_all_namespaces()
for pod in pods.items:
if pod.metadata.labels and 'app' in pod.metadata.labels:
# 通过标签匹配找到所属Service
for svc in services.items:
if self._match_labels(svc.spec.selector, pod.metadata.labels):
G.add_edge(svc.metadata.name, pod.metadata.name, type='runs')
break
# 3. 通过NetworkPolicy或Istio判断服务间调用关系
# (实际生产环境中可通过APM数据增强)
return G纯阈值告警误报率高,我们采用 “三阶段联合检测”:
class AnomalyDetector:
def __init__(self):
self.statistical_detector = StatisticalDetector() # 3-sigma/IQR
self.ml_detector = MLDetector() # Isolation Forest
self.llm_verifier = LLMVerifier() # 大模型二次确认
def detect(self, metric_series, context):
# 阶段1:统计检测(快速筛选)
stat_anomalies = self.statistical_detector.detect(metric_series)
if not stat_anomalies:
return []
# 阶段2:机器学习检测(降噪)
ml_anomalies = self.ml_detector.detect(metric_series, stat_anomalies)
if not ml_anomalies:
return []
# 阶段3:大模型验证(减少误报)
# 传入该服务的历史模式、变更记录等上下文
confirmed = []
for anomaly in ml_anomalies:
is_real = self.llm_verifier.verify(
anomaly=anomaly,
historical_pattern=context['historical_pattern'],
recent_changes=context['changes'],
related_services=context['related_services']
)
if is_real:
confirmed.append(anomaly)
return confirmed当多个告警同时发生时,系统需要快速定位哪个是“因”、哪些是“果”。
# 根因定位引擎
class RootCauseEngine:
def __init__(self, neo4j_driver, llm):
self.driver = neo4j_driver
self.llm = llm
def locate_root_cause(self, anomalies: List[Anomaly]) -> RootCauseResult:
# 1. 从知识图谱中提取受影响节点的上下游依赖
affected_services = [a.service for a in anomalies]
with self.driver.session() as session:
# Cypher查询:查找所有受影响节点的共同上游
result = session.run("""
MATCH (s:Service)-[:DEPENDS_ON*1..5]->(upstream:Service)
WHERE s.name IN $services
RETURN upstream.name AS candidate, COUNT(*) AS affected_count
ORDER BY affected_count DESC
LIMIT 5
""", services=affected_services)
candidates = [record['candidate'] for record in result]
# 2. 结合时间窗口:最早出现的异常更有可能是根因
# 3. 结合变更记录:最近变更的服务风险最高
candidates_with_metadata = []
for candidate in candidates:
first_anomaly_time = min([a.timestamp for a in anomalies if a.service == candidate], default=None)
recent_changes = get_recent_changes(candidate, hours=1)
candidates_with_metadata.append({
'service': candidate,
'first_occurrence': first_anomaly_time,
'recent_changes': recent_changes,
'impacted_downstream': [a.service for a in anomalies if candidate in a.dependencies]
})
# 4. 使用大模型进行综合推理
root_cause = self.llm.invoke(
f"根据以下候选根因分析,结合运维经验判断最可能的故障源:{candidates_with_metadata}"
)
return RootCauseResult(
root_service=root_cause['service'],
confidence=root_cause['confidence'],
reasoning=root_cause['explanation'],
evidence=root_cause['evidence']
)大模型Agent根据根因、历史案例、可用工具,自动生成处置方案。
# 处置方案生成Agent(使用LangGraph实现状态机)
from langgraph.graph import StateGraph, END
class RemediationAgent:
def __init__(self, llm, tools, knowledge_base):
self.llm = llm
self.tools = tools # restart_pod, scale_up, rollback_config, etc.
self.kb = knowledge_base
def build_graph(self):
workflow = StateGraph(RemediationState)
# 定义节点
workflow.add_node("analyze_root_cause", self.analyze_root_cause)
workflow.add_node("search_history", self.search_history)
workflow.add_node("generate_plan", self.generate_plan)
workflow.add_node("assess_risk", self.assess_risk)
workflow.add_node("execute_plan", self.execute_plan)
workflow.add_node("verify_result", self.verify_result)
workflow.add_node("escalate", self.escalate_to_human)
# 定义边(条件路由)
workflow.set_entry_point("analyze_root_cause")
workflow.add_edge("analyze_root_cause", "search_history")
workflow.add_edge("search_history", "generate_plan")
workflow.add_edge("generate_plan", "assess_risk")
workflow.add_conditional_edges(
"assess_risk",
lambda state: "execute_plan" if state['risk_score'] < 0.7 else "escalate"
)
workflow.add_edge("execute_plan", "verify_result")
workflow.add_conditional_edges(
"verify_result",
lambda state: END if state['success'] else "escalate"
)
return workflow.compile()
def generate_plan(self, state):
"""生成具体处置步骤"""
prompt = f"""
根因:{state['root_cause']}
历史相似案例:{state['history_cases']}
可用工具列表:{[t.name for t in self.tools]}
请生成一个处置计划,包含:
1. 操作步骤(按顺序)
2. 每一步调用的工具和参数
3. 预期结果
4. 回滚条件
以JSON格式输出。
"""
plan = self.llm.invoke(prompt, response_format='json')
state['plan'] = json.loads(plan)
return state
def assess_risk(self, state):
"""风险评估:操作是否可能造成更大影响"""
# 检查是否有"删除"、"重启生产主库"等高危操作
plan_text = str(state['plan'])
high_risk_keywords = ['delete', 'drop', 'reboot', 'shutdown', 'force']
for keyword in high_risk_keywords:
if keyword in plan_text.lower():
# 如果是高危操作,进行影响面分析
impact = self.impact_analysis(state['plan'])
state['risk_score'] = impact['score']
state['risk_detail'] = impact['detail']
return state
state['risk_score'] = 0.3 # 默认低风险
return state@Service
public class RemediationExecutor {
@Autowired
private K8sClient k8sClient;
@Autowired
private TerraformClient tfClient;
@Autowired
private AnsibleClient ansibleClient;
/**
* 执行处置计划,支持暂停点(等待人工确认高风险步骤)
*/
public ExecutionResult execute(RemediationPlan plan, ExecutionContext ctx) {
List<ExecutionStep> steps = plan.getSteps();
ExecutionResult result = new ExecutionResult();
for (int i = 0; i < steps.size(); i++) {
ExecutionStep step = steps.get(i);
// 检查是否需要人工审批
if (step.requiresApproval() && !ctx.isApproved(step.getId())) {
// 暂停执行,等待外部审批信号
result.setStatus(ExecutionStatus.PENDING_APPROVAL);
result.setPendingStep(step);
return result;
}
try {
// 执行操作并记录快照(用于回滚)
StateSnapshot snapshot = executeStep(step);
result.getSnapshots().add(snapshot);
result.getExecutedSteps().add(step);
} catch (Exception e) {
// 执行失败,触发自动回滚
log.error("Step {} failed: {}", step.getId(), e.getMessage());
rollback(result.getSnapshots());
result.setStatus(ExecutionStatus.ROLLED_BACK);
result.setError(e.getMessage());
return result;
}
}
result.setStatus(ExecutionStatus.SUCCESS);
return result;
}
/**
* 回滚:按执行顺序的反向执行回滚操作
*/
private void rollback(List<StateSnapshot> snapshots) {
Collections.reverse(snapshots);
for (StateSnapshot snapshot : snapshots) {
if (snapshot.getType() == SnapshotType.POD_RESTART) {
// 回滚重启:无操作(重启本身不可逆,但可以检查状态)
continue;
} else if (snapshot.getType() == SnapshotType.SCALE_CHANGE) {
// 回滚扩容:缩容回原数量
k8sClient.scaleDeployment(snapshot.getDeployment(), snapshot.getOriginalReplicas());
} else if (snapshot.getType() == SnapshotType.CONFIG_CHANGE) {
// 回滚配置变更
ansibleClient.rollbackConfig(snapshot.getConfigVersion());
}
}
}
}# 执行权限配置(OPA策略)
apiVersion: v1
kind: ConfigMap
metadata:
name: remediation-policies
data:
policy.rego: |
package remediation
# 默认拒绝所有操作
default allow = false
# 允许只读操作(日志查询、状态查看)
allow = true {
input.operation in ["query_logs", "get_pod_status", "get_metrics"]
}
# 允许指定服务在非生产环境执行操作
allow = true {
input.environment == "staging"
input.service in authorized_staging_services
input.operation in allowed_staging_operations
}
# 生产环境高风险操作需要双人审批
allow = true {
input.environment == "production"
input.operation in ["restart_pod", "scale_deployment"]
input.approvals == ["sre-oncall", "service-owner"]
}python
class VerificationEngine:
def __init__(self, prometheus_client):
self.prometheus = prometheus_client
def verify_remediation(self, incident, executed_plan):
"""验证处置措施是否有效"""
# 1. 检查根因指标是否恢复正常
root_metric = incident['root_cause_metric']
current_value = self.prometheus.query(
f"{root_metric['name']}{{service='{root_metric['service']}'}}"
)
expected_value = root_metric['baseline']
metric_ok = abs(current_value - expected_value) / expected_value < 0.1
# 2. 检查副作用:其他服务是否受影响
side_effects = self.check_side_effects(
affected_services=incident['related_services'],
duration=executed_plan['execution_time']
)
# 3. 检查用户体验指标
error_rate = self.prometheus.query(
f"sum(rate(http_requests_total{{status=~'5..'}}[5m])) / sum(rate(http_requests_total[5m]))"
)
ux_ok = error_rate < incident['sla_threshold']
return VerificationResult(
success=metric_ok and not side_effects and ux_ok,
metric_ok=metric_ok,
side_effects=side_effects,
ux_ok=ux_ok
)class IncidentLearning:
def __init__(self, knowledge_base, llm):
self.kb = knowledge_base
self.llm = llm
def archive_incident(self, incident, resolution):
"""将成功处置的故障归档到知识库"""
# 使用大模型生成结构化摘要
summary = self.llm.invoke(f"""
请将以下故障处置过程总结为一个可复用的知识条目:
故障描述:{incident['description']}
症状:{incident['symptoms']}
根因:{incident['root_cause']}
处置过程:{resolution['steps']}
耗时:{resolution['duration']}
验证结果:{resolution['verification']}
输出格式:
- 故障类型:[类型标签]
- 典型症状:[关键症状指标]
- 修复步骤:[操作步骤]
- 适用场景:[条件]
- 风险提示:[注意事项]
""")
# 存入向量库,供后续相似故障检索
self.kb.add_incident(
incident_id=incident['id'],
summary=summary,
symptoms=incident['symptoms'],
resolution=resolution['steps'],
tags=extract_tags(summary)
)阶段 | 周期 | 目标 | 核心交付 | 成功标准 |
|---|---|---|---|---|
POC阶段 | 1-2月 | 验证核心能力 | 单场景异常检测+告警降噪 | 告警减少50%,根因定位准确率>70% |
试点阶段 | 3-4月 | 扩展到2-3个核心业务 | 自动化处置能力(只读+重启) | 平均修复时间(MTTR)降低40% |
扩展阶段 | 5-8月 | 覆盖全业务线 | 完整闭环(含回滚+验证) | 30%故障完全自动修复,人工介入减少60% |
优化阶段 | 9-12月 | 持续迭代 | 知识库自动更新+模型微调 | MTTR降低70%,SLA达标率99.99% |
陷阱 | 表现 | 应对策略 |
|---|---|---|
过度自动化 | 自动操作导致更大范围故障 | 设置“围栏”:只在已验证的故障模式上启用自动修复,新场景强制人工审批 |
数据质量差 | 告警不准确、日志缺失 | 先做数据治理:统一日志格式、完善标签体系、清理无效告警 |
大模型幻觉 | 生成错误的处置方案 | 强制“方案审查”环节:高风险操作必须经过规则引擎二次校验 |
权限扩散 | 系统权限过大 | 使用临时凭据(STS)、操作审计、权限隔离 |
回滚失败 | 无法恢复到故障前状态 | 所有变更前强制快照,回滚预案随执行计划一起生成 |
指标 | 实施前 | 实施后 | 改善幅度 |
|---|---|---|---|
MTTR(平均修复时间) | 45分钟 | 12分钟 | ↓73% |
MTTD(平均发现时间) | 15分钟 | 3分钟(AI预测) | ↓80% |
月度告警数量 | 8,700条 | 2,100条 | ↓76% |
自动修复比例 | 0% | 35% | +35% |
SRE人力占用 | 60% | 25% | ↓58% |
重大故障次数 | 年均8次 | 年均2次 | ↓75% |
关键成功因素:
构建云环境自动化故障处置平台,本质上是将SRE的经验、直觉和决策能力,逐步编码为可运行的系统。这需要:
未来的运维工程师,不再需要凌晨3点被电话吵醒。他们将有更多时间去做更有价值的事情——设计更优雅的系统架构、编写更好的混沌实验、训练更智能的故障预测模型。
是时候让你的运维平台,学会自己“看病”了。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。