错误处理哲学:三个进程如何优雅处理故障
在分布式系统中,网络分区、进程崩溃、资源耗尽都是日常。智策平台由三个独立进程组成——onto-control(Control Layer)、onto-data(Data Layer)、onto-intelligence(Reasoning & Decision Layer + Agent Runtime Layer),它们通过 gRPC 互相调用。任何一次跨进程调用都可能失败,任何一个进程都可能宕机。
错误处理哲学:三个进程如何优雅处理故障
“系列:S2 架构全景 · 第 11 篇 | 难度:中级 | 阅读时间:18 分钟
#TL;DR
- 智策平台采用统一 gRPC 错误码体系,将所有服务间错误映射到标准 gRPC Status Code,前端通过 API Gateway 转换为 HTTP 状态码,实现错误语义一致性。
- 三个进程分别实现了分层重试策略:onto-control 使用 Spring Retry 指数退避、onto-data 使用 Quarkus Fault Tolerance 断路器、onto-intelligence 使用 Temporal 工作流级重试,确保故障不会级联扩散。
- 死信队列(DLQ)+ 人工审批回退 + 优雅降级三重保障,使平台在部分故障时仍能提供有损但可用的服务。
#引言:故障是常态,不是意外
在分布式系统中,网络分区、进程崩溃、资源耗尽都是日常。智策平台由三个独立进程组成——onto-control(Control Layer)、onto-data(Data Layer)、onto-intelligence(Reasoning & Decision Layer + Agent Runtime Layer),它们通过 gRPC 互相调用。任何一次跨进程调用都可能失败,任何一个进程都可能宕机。
故障概率模型(简化):
单次 gRPC 调用失败概率: ~0.1%(网络抖动)
单个进程 1 小时内重启概率:~0.5%(OOM、GC 停顿)
一次业务操作涉及调用次数:3-7 次
一次业务操作失败概率: ~0.3%-0.7%
日均请求量 10 万次时:
预期每日失败请求数:300-700 次
用户可感知的错误数:目标 < 10 次(通过重试和降级消化)
本文将系统讲解智策平台如何通过 gRPC 错误码规范、多层重试策略、断路器模式、优雅降级和死信队列五大机制,将故障对用户的影响降到最低。
#1. gRPC 错误码体系
#1.1 标准错误码映射
智策平台强制所有服务间通信使用 gRPC。gRPC 定义了 16 个标准状态码,平台约定了每个码的使用场景:
gRPC Status Code 使用规范:
Code 值 使用场景 是否可重试
──────────────────────────────────────────────────────────────────
OK 0 成功 -
CANCELLED 1 客户端主动取消 否
UNKNOWN 2 未分类错误(兜底) 视情况
INVALID_ARGUMENT 3 参数校验失败 否(修正后可)
DEADLINE_EXCEEDED 4 超时 是
NOT_FOUND 5 资源不存在 否
ALREADY_EXISTS 6 资源已存在(幂等检查) 否
PERMISSION_DENIED 7 权限不足 否
RESOURCE_EXHAUSTED 8 配额或限流 是(退避后)
FAILED_PRECONDITION 9 前置条件不满足 否(修正后可)
ABORTED 10 事务冲突/乐观锁失败 是
OUT_OF_RANGE 11 分页越界等 否
UNIMPLEMENTED 12 未实现的方法 否
INTERNAL 13 服务内部错误 是(有限次)
UNAVAILABLE 14 服务不可用 是(退避后)
DATA_LOSS 15 不可恢复的数据丢失 否
UNAUTHENTICATED 16 未认证 否(重新认证后可)
#1.2 错误详情(Error Details)
裸露的状态码不足以定位问题。平台要求所有 gRPC 错误必须携带结构化的错误详情:
// 平台统一错误详情
message PlatformErrorDetail {
string error_id = 1; // 唯一错误 ID(用于日志关联)
string service = 2; // 发生错误的服务名
string Layer = 3; // 所属 Layer(B/C/D/E)
int64 timestamp = 4; // 错误发生时间戳
string trace_id = 5; // 分布式追踪 ID
string world_id = 6; // 关联的 World ID
repeated string context = 7; // 上下文信息链
RetryAdvice retry_advice = 8; // 重试建议
}
message RetryAdvice {
bool retryable = 1; // 是否可重试
int32 max_retries = 2; // 建议最大重试次数
int64 retry_after_ms = 3; // 建议等待时间(毫秒)
string strategy = 4; // "exponential" | "linear" | "fixed"
}
#1.3 错误传播链
当一次请求跨多个进程时,错误需要沿调用链传播,同时保留每一层的上下文:
错误传播示例(创建 Action 失败):
用户请求 → API Gateway → onto-control → onto-intelligence
↓
规则评估超时
gRPC DEADLINE_EXCEEDED
↓
onto-control 捕获,追加上下文
"Action pre-check failed: rule evaluation timeout"
gRPC DEADLINE_EXCEEDED + PlatformErrorDetail
↓
API Gateway 转换
HTTP 504 Gateway Timeout
{
"error": "ACTION_PRECHECK_TIMEOUT",
"message": "Action 前置检查超时",
"trace_id": "abc-123-def",
"retry_after": 5000
}
#1.4 三个进程的错误拦截器
每个进程都实现了 gRPC 拦截器,统一处理错误的捕获、包装和上报:
onto-control(Java / Spring Boot):
GlobalGrpcExceptionInterceptor
├── 捕获 Spring 异常 → 映射 gRPC 错误码
├── 注入 trace_id 和 world_id
├── 记录到 structured log(JSON 格式)
└── 上报 Prometheus metrics(grpc_server_errors_total)
onto-data(Java / Quarkus):
QuarkusGrpcErrorInterceptor
├── 捕获 Iceberg/Nessie 异常 → 映射 gRPC 错误码
├── NessieConflictException → ABORTED(可重试)
├── IcebergCommitFailedException → ABORTED(可重试)
└── 上报 Micrometer metrics
onto-intelligence(Python / FastAPI):
grpc_error_interceptor(异步拦截器)
├── 捕获 Python 异常 → 映射 gRPC 错误码
├── TimeoutError → DEADLINE_EXCEEDED
├── RuleEvaluationError → FAILED_PRECONDITION
└── 上报 Prometheus metrics(prometheus_client)
#2. 重试策略
#2.1 重试决策树
不是所有错误都应该重试。平台定义了清晰的重试决策树:
收到 gRPC 错误
│
├─ 检查 RetryAdvice.retryable
│ ├─ false → 直接返回错误
│ └─ true → 继续
│
├─ 检查错误码
│ ├─ UNAVAILABLE → 重试(服务可能正在重启)
│ ├─ DEADLINE_EXCEEDED → 重试(可能是暂时性超时)
│ ├─ ABORTED → 重试(乐观锁冲突)
│ ├─ RESOURCE_EXHAUSTED → 重试(限流,需要退避)
│ ├─ INTERNAL → 有限重试(可能是 bug,也可能是暂时性)
│ └─ 其他 → 不重试
│
├─ 检查重试次数
│ ├─ < max_retries → 执行重试
│ └─ >= max_retries → 放弃,触发降级逻辑
│
└─ 计算退避时间
├─ exponential: base * 2^attempt + jitter
├─ linear: base * attempt + jitter
└─ fixed: base + jitter
#2.2 onto-control:Spring Retry
onto-control 使用 Spring Retry 框架实现声明式重试:
// onto-control 中的 gRPC 客户端重试配置
@Configuration
public class GrpcRetryConfig {
@Bean
public RetryTemplate grpcRetryTemplate() {
return RetryTemplate.builder()
.maxAttempts(3)
.exponentialBackoff(
100, // initialInterval: 100ms
2.0, // multiplier
2000, // maxInterval: 2s
true // withJitter
)
.retryOn(StatusRuntimeException.class)
.traversingCauses()
.build();
}
}
// 使用示例:调用 onto-intelligence 的规则评估
@GrpcService
public class ActionServiceImpl {
@Retryable(
retryFor = {StatusRuntimeException.class},
maxAttempts = 3,
backoff = @Backoff(delay = 200, multiplier = 2)
)
public ActionResult executeAction(ActionRequest request) {
// 调用 onto-intelligence
var result = intelligenceClient.evaluateRules(request);
return result;
}
@Recover
public ActionResult executeActionFallback(
StatusRuntimeException e, ActionRequest request) {
// 降级:记录到待处理队列
pendingActionQueue.enqueue(request, e);
return ActionResult.pending("规则评估暂不可用,Action 已加入待处理队列");
}
}
#2.3 onto-data:Quarkus Fault Tolerance
onto-data 使用 MicroProfile Fault Tolerance 注解:
// onto-data 中的 Nessie 操作重试
@ApplicationScoped
public class NessieCommitService {
@Retry(
maxRetries = 5,
delay = 100,
maxDuration = 10000, // 最长 10 秒
jitter = 50,
retryOn = {NessieConflictException.class}
)
@Fallback(fallbackMethod = "commitWithMerge")
public CommitResult commitChanges(Branch branch, List<Operation> ops) {
return nessieApi.commitMultipleOperations()
.branch(branch)
.operations(ops)
.commit();
}
// 降级:尝试先合并再提交
public CommitResult commitWithMerge(Branch branch, List<Operation> ops) {
var latestBranch = nessieApi.getReference()
.refName(branch.getName())
.get();
// 三路合并 + 重试提交
return mergeAndCommit(latestBranch, branch, ops);
}
}
#2.4 onto-intelligence:Temporal 工作流重试
onto-intelligence 中的长时间运行任务使用 Temporal 工作流,自带重试机制:
# onto-intelligence 中的 Temporal Activity 重试配置
from temporalio import activity, workflow
from temporalio.common import RetryPolicy
from datetime import timedelta
STANDARD_RETRY_POLICY = RetryPolicy(
initial_interval=timedelta(seconds=1),
backoff_coefficient=2.0,
maximum_interval=timedelta(minutes=5),
maximum_attempts=10,
non_retryable_error_types=[
"InvalidArgumentError",
"PermissionDeniedError",
"DataValidationError",
],
)
@workflow.defn
class RuleEvaluationWorkflow:
@workflow.run
async def run(self, request: RuleEvalRequest) -> RuleEvalResult:
# Activity 调用自动重试
result = await workflow.execute_activity(
evaluate_rules,
request,
start_to_close_timeout=timedelta(minutes=2),
retry_policy=STANDARD_RETRY_POLICY,
)
return result
@activity.defn
async def evaluate_rules(request: RuleEvalRequest) -> RuleEvalResult:
"""评估规则 — 失败时 Temporal 自动按策略重试"""
engine = get_rule_engine()
return await engine.evaluate(request.rules, request.context)
#3. 断路器模式
#3.1 为什么需要断路器
重试能处理瞬时故障,但如果下游服务持续不可用,无限重试只会加剧问题。断路器(Circuit Breaker)在检测到连续失败后,会"断开"对下游的调用,直接返回降级结果。
断路器状态机:
┌─────────┐ 连续失败 ┌────────┐
│ CLOSED │ ──────────────→ │ OPEN │
│ (正常) │ > 阈值 │ (熔断) │
└─────────┘ └────────┘
↑ │
│ 成功请求 │ 超时后允许
│ > 阈值 │ 少量请求探测
│ ↓
┌─────────────┐
│ HALF-OPEN │
│ (半开/探测) │
└─────────────┘
参数配置:
失败阈值:连续 5 次失败 → 打开断路器
断路时间:30 秒后进入半开状态
探测数量:半开状态允许 3 个请求通过
恢复阈值:3 个请求全部成功 → 关闭断路器
#3.2 各进程的断路器实现
onto-control(Spring Boot):
使用 Resilience4j CircuitBreaker
配置文件:application.yml
──────────────────────────────
resilience4j:
circuitbreaker:
instances:
intelligence-service:
slidingWindowSize: 10
failureRateThreshold: 50
waitDurationInOpenState: 30s
permittedNumberOfCallsInHalfOpenState: 3
data-service:
slidingWindowSize: 20
failureRateThreshold: 60
waitDurationInOpenState: 20s
onto-data(Quarkus):
使用 MicroProfile Fault Tolerance @CircuitBreaker
──────────────────────────────
@CircuitBreaker(
requestVolumeThreshold = 10,
failureRatio = 0.5,
delay = 30000, // 30s
successThreshold = 3
)
public QueryResult queryObjects(QueryRequest request) { ... }
onto-intelligence(Python):
使用 pybreaker 库
──────────────────────────────
control_breaker = CircuitBreaker(
fail_max=5,
reset_timeout=30,
state_storage=RedisCircuitBreakerStorage(
state_key="cb:control-service",
redis_client=redis_client,
),
)
@control_breaker
async def call_control_service(request):
async with grpc.aio.insecure_channel(CONTROL_ADDR) as channel:
stub = ControlServiceStub(channel)
return await stub.GetOntologySchema(request)
#3.3 断路器监控
所有断路器状态都暴露到 Prometheus,可在 Grafana 中实时查看:
Prometheus 指标:
# onto-control
resilience4j_circuitbreaker_state{name="intelligence-service"} 0|1|2
resilience4j_circuitbreaker_failure_rate{name="intelligence-service"} 23.5
resilience4j_circuitbreaker_calls_total{name="intelligence-service",kind="successful"} 9823
resilience4j_circuitbreaker_calls_total{name="intelligence-service",kind="failed"} 42
# onto-intelligence
circuit_breaker_state{service="control-service"} closed|open|half-open
circuit_breaker_failure_count{service="control-service"} 3
#4. 优雅降级
#4.1 降级等级
平台定义了四个降级等级,不同等级下提供不同的服务能力:
降级等级矩阵:
Level 名称 触发条件 影响
─────────────────────────────────────────────────────────
L0 全部正常 所有服务正常 无
L1 部分降级 单个下游服务断路器打开 功能受限
L2 严重降级 2+ 下游服务不可用 只读模式
L3 最小服务 核心进程异常 仅返回缓存
各等级可用功能:
L0 L1 L2 L3
Ontology 查询 ✅ ✅ ✅ ⚠️(缓存)
对象实例 CRUD ✅ ✅ ❌ ❌
Action 执行 ✅ ⚠️ ❌ ❌
规则评估 ✅ ⚠️ ❌ ❌
派生属性 ✅ ❌ ❌ ❌
Dashboard 展示 ✅ ✅ ⚠️ ⚠️(快照)
#4.2 降级策略实现
# onto-intelligence 的降级服务
class DegradedReasoningService:
"""当 onto-control 不可用时的降级实现"""
def __init__(self, cache: OntologyCache):
self.cache = cache
self.degradation_level = DegradationLevel.L0
async def evaluate_rule(self, rule_id: str, context: dict) -> RuleResult:
if self.degradation_level >= DegradationLevel.L2:
# L2: 返回缓存的上次评估结果
cached = await self.cache.get_last_result(rule_id)
if cached:
return RuleResult(
value=cached.value,
confidence=0.5, # 降低置信度
source="cache",
stale=True,
cached_at=cached.timestamp,
)
raise ServiceDegradedError("规则评估不可用,且无缓存结果")
if self.degradation_level == DegradationLevel.L1:
# L1: 尝试评估,但使用缓存的 Schema
try:
schema = await self.cache.get_schema(rule_id)
return await self._evaluate_with_cached_schema(rule_id, schema, context)
except Exception:
cached = await self.cache.get_last_result(rule_id)
if cached:
return RuleResult(value=cached.value, confidence=0.3, source="cache")
raise
# L0: 正常评估
return await self._evaluate_normal(rule_id, context)
#4.3 自动降级与恢复
自动降级控制器(每个进程内运行):
每 5 秒检查一次 →
├─ 检查各断路器状态
│ ├─ 全部 CLOSED → L0
│ ├─ 1 个 OPEN → L1
│ ├─ 2+ 个 OPEN → L2
│ └─ 自身健康检查失败 → L3
│
├─ 发送降级事件到 Kafka
│ topic: platform.degradation.events
│ {
│ "service": "onto-intelligence",
│ "level": "L1",
│ "reason": "control-service circuit breaker open",
│ "timestamp": "2026-03-24T10:30:00Z"
│ }
│
└─ 当所有断路器恢复 CLOSED →
延迟 60 秒确认稳定 → 恢复到 L0
发送恢复事件
#5. 死信队列(DLQ)
#5.1 消息消费失败处理
Kafka 消息消费失败时,不能简单丢弃。平台使用死信队列保证消息不丢失:
Kafka 消息处理流程:
正常 Topic 重试 Topic 死信 Topic
platform.actions → platform.actions.retry → platform.actions.dlq
│ │ │
├─ 消费成功 → ACK ├─ 重试成功 → ACK ├─ 人工处理
├─ 消费失败 → ├─ 重试 3 次仍失败 → │
│ 发送到 retry topic │ 发送到 dlq topic └─ 运维告警
└─ 反序列化失败 → └─
直接进 dlq
重试策略:
第 1 次重试:延迟 1 秒
第 2 次重试:延迟 5 秒
第 3 次重试:延迟 30 秒
第 3 次仍失败 → 进入 DLQ
#5.2 DLQ 消息结构
{
"original_topic": "platform.actions",
"original_partition": 3,
"original_offset": 12847,
"original_key": "world-001:action-execute",
"original_value": "<base64 encoded original message>",
"error_history": [
{
"attempt": 1,
"timestamp": "2026-03-24T10:30:01Z",
"error": "NessieConflictException: commit conflict on branch main",
"service": "onto-data",
"trace_id": "abc-123"
},
{
"attempt": 2,
"timestamp": "2026-03-24T10:30:06Z",
"error": "NessieConflictException: commit conflict on branch main",
"service": "onto-data",
"trace_id": "abc-124"
},
{
"attempt": 3,
"timestamp": "2026-03-24T10:30:36Z",
"error": "NessieConflictException: commit conflict on branch main",
"service": "onto-data",
"trace_id": "abc-125"
}
],
"dlq_timestamp": "2026-03-24T10:30:36Z",
"status": "PENDING_REVIEW"
}
#5.3 DLQ 管理与回放
DLQ 管理 API(onto-control 提供):
GET /api/v1/dlq/messages — 查看 DLQ 中的消息列表
GET /api/v1/dlq/messages/{id} — 查看单条消息详情
POST /api/v1/dlq/messages/{id}/replay — 重放单条消息
POST /api/v1/dlq/messages/replay-all — 重放所有消息
DELETE /api/v1/dlq/messages/{id} — 确认并删除消息
GET /api/v1/dlq/stats — DLQ 统计信息
DLQ 告警规则(Prometheus AlertManager):
- DLQ 消息数 > 10 → Warning 级别告警
- DLQ 消息数 > 100 → Critical 级别告警
- DLQ 消息积压时间 > 1 小时 → Warning
- DLQ 消息积压时间 > 24 小时 → Critical
#6. Temporal 工作流的错误处理
#6.1 工作流级别的容错
onto-intelligence 中的复杂业务逻辑(决策执行、批量规则评估)使用 Temporal 工作流编排。Temporal 提供了比普通重试更强大的容错能力:
Temporal 容错特性:
1. Activity 自动重试
├─ 可配置重试策略(次数、间隔、退避系数)
├─ 可指定不可重试的异常类型
└─ 重试对工作流代码透明
2. 工作流超时控制
├─ WorkflowExecutionTimeout:整个工作流的最大运行时间
├─ WorkflowRunTimeout:单次 run 的最大时间
└─ ActivityStartToCloseTimeout:单个 Activity 的最大时间
3. 心跳检测
├─ 长时间 Activity 必须定期报告心跳
├─ 心跳超时 → Temporal 认为 Activity 失败 → 触发重试
└─ 心跳可携带进度信息,重试时从断点恢复
4. 补偿(Saga 模式)
├─ 每个 Activity 可定义对应的补偿 Activity
├─ 工作流失败时按逆序执行补偿
└─ 保证最终一致性
#6.2 Saga 补偿模式
对于涉及多个服务的操作(如执行 Action),平台使用 Saga 模式保证最终一致性:
@workflow.defn
class ExecuteActionWorkflow:
"""Action 执行工作流 — Saga 模式"""
@workflow.run
async def run(self, request: ActionExecuteRequest) -> ActionResult:
compensations: list[Callable] = []
try:
# Step 1: 前置规则评估
pre_check = await workflow.execute_activity(
evaluate_preconditions,
request,
start_to_close_timeout=timedelta(seconds=30),
retry_policy=STANDARD_RETRY_POLICY,
)
if not pre_check.passed:
return ActionResult.rejected(pre_check.reason)
# Step 2: 锁定相关对象(乐观锁)
lock = await workflow.execute_activity(
acquire_object_locks,
request.object_ids,
start_to_close_timeout=timedelta(seconds=10),
)
compensations.append(lambda: release_object_locks(lock))
# Step 3: 执行数据变更
mutation = await workflow.execute_activity(
apply_mutations,
request.mutations,
start_to_close_timeout=timedelta(minutes=1),
retry_policy=STANDARD_RETRY_POLICY,
)
compensations.append(lambda: rollback_mutations(mutation))
# Step 4: 后置规则评估
post_check = await workflow.execute_activity(
evaluate_postconditions,
request,
start_to_close_timeout=timedelta(seconds=30),
)
if not post_check.passed:
# 后置检查失败 → 触发补偿回滚
raise PostconditionFailedError(post_check.reason)
# Step 5: 提交并发布事件
await workflow.execute_activity(
commit_and_publish,
mutation,
start_to_close_timeout=timedelta(seconds=15),
)
return ActionResult.success(mutation.summary)
except Exception as e:
# 逆序执行补偿
for compensate in reversed(compensations):
try:
await workflow.execute_activity(
compensate,
start_to_close_timeout=timedelta(seconds=30),
)
except Exception as comp_error:
workflow.logger.error(
f"补偿失败: {comp_error},需要人工介入"
)
raise
#6.3 Temporal 工作流监控
Temporal 仪表板关键指标:
工作流统计:
├─ 活跃工作流数量
├─ 工作流完成率(成功 / 总数)
├─ 平均工作流执行时间
└─ 工作流失败原因分布
Activity 统计:
├─ Activity 执行次数(含重试)
├─ Activity 平均延迟
├─ Activity 失败率
└─ 重试次数分布
告警规则:
├─ 工作流失败率 > 5% → Warning
├─ 工作流失败率 > 15% → Critical
├─ 活跃工作流数 > 1000 → Warning(可能有积压)
└─ Activity 平均延迟 > 30s → Warning
#7. 错误恢复与自愈
#7.1 健康检查机制
每个进程都实现了多层健康检查:
健康检查层次:
Layer 1: Liveness(存活检查)
检查内容:进程是否运行、是否能响应请求
频率:每 5 秒
失败处理:Docker/K8s 重启容器
端点:
onto-control: GET /actuator/health/liveness
onto-data: GET /q/health/live
onto-intelligence: GET /health/live
Layer 2: Readiness(就绪检查)
检查内容:依赖服务是否可用(数据库、Kafka、下游服务)
频率:每 10 秒
失败处理:从负载均衡中摘除,不再接收新请求
端点:
onto-control: GET /actuator/health/readiness
onto-data: GET /q/health/ready
onto-intelligence: GET /health/ready
Layer 3: Startup(启动检查)
检查内容:初始化是否完成(Schema 加载、缓存预热)
频率:每 2 秒,最多等 120 秒
失败处理:标记为启动失败,触发告警
#7.2 自动恢复策略
自动恢复矩阵:
故障类型 检测方式 恢复策略
────────────────────────────────────────────────────────
进程 OOM Docker 健康检查 自动重启(最多 3 次/小时)
数据库连接断开 连接池心跳 自动重连 + 指数退避
Kafka 消费者掉线 Consumer group rebalance 自动 rejoin
gRPC 连接中断 Channel 状态监控 自动重建连接
Nessie 冲突 ABORTED 错误码 自动合并重试
Schema 缓存过期 TTL + 版本号 自动刷新
Temporal Worker 崩溃 Temporal Server 检测 自动重新分配任务
#7.3 手动恢复工具
当自动恢复无法解决问题时,运维人员可以使用以下工具:
运维命令(通过 onto-control Admin API):
# 强制刷新所有 Schema 缓存
POST /admin/cache/schema/refresh
# 重新初始化 gRPC 连接
POST /admin/grpc/reconnect?target=onto-intelligence
# 手动触发断路器状态转换
POST /admin/circuit-breaker/intelligence-service/force-close
# 重放 DLQ 消息
POST /admin/dlq/replay?topic=platform.actions&count=10
# 导出错误报告
GET /admin/errors/report?from=2026-03-24T00:00:00Z&to=2026-03-24T23:59:59Z
#8. 实战案例:一次级联故障的处理
#8.1 故障场景
时间线:
T+0s: onto-data 的 Doris 连接池耗尽(100 个连接全部占用)
T+1s: 新的查询请求开始排队,10s 后超时
T+10s: onto-control 调用 onto-data 的 gRPC 请求开始超时
错误码:DEADLINE_EXCEEDED
T+15s: onto-control 的 data-service 断路器:失败 3/10
T+30s: onto-control 的 data-service 断路器:失败 5/10 → 打开
T+30s: 降级控制器检测到断路器打开 → 切换到 L1
T+30s: Dashboard 查询走缓存,Action 执行被暂停
T+45s: Doris 连接池开始恢复(慢查询被 kill)
T+60s: onto-data 健康检查恢复正常
T+90s: onto-control 断路器进入 HALF-OPEN,放行 3 个请求
T+91s: 3 个请求全部成功
T+91s: 断路器关闭 → L0 恢复延迟倒计时 60s
T+151s: 降级控制器确认稳定 → 恢复 L0
T+151s: 暂停的 Action 从待处理队列中恢复执行
#8.2 事后分析
故障根因:
一个用户提交了全表扫描查询,占用了 30 个连接 × 30 秒
连接池配置 max=100,正常负载已使用 75 个
剩余 25 个连接在 3 秒内被新请求耗尽
修复措施:
1. 查询超时限制:单查询最长 10 秒(之前是 60 秒)
2. 连接池扩容:max=200,min-idle=50
3. 查询复杂度检查:估算行数 > 100 万的查询需要审批
4. 连接池使用率告警:> 80% 时 Warning
故障影响统计:
故障持续时间:151 秒
受影响请求数:~240 个
用户可见错误:12 个(其余被降级和重试消化)
数据丢失:0(DLQ 保证消息不丢)
Action 延迟执行:37 个(全部在恢复后 5 分钟内完成)
#Key Takeaways
- 错误码标准化:统一的 gRPC 错误码 + 结构化 ErrorDetail,使错误在跨服务传播时不丢失语义。
- 分层重试:不同技术栈使用不同的重试实现(Spring Retry / MicroProfile FT / Temporal),但遵循相同的重试决策树。
- 断路器保护:防止故障级联扩散,每个进程独立管理对下游的断路器状态。
- 优雅降级:四级降级模型确保平台在部分故障时仍能提供有限服务,而不是完全不可用。
- 死信队列兜底:消息消费失败不丢弃,进入 DLQ 等待人工处理或重放。
- Saga 补偿:跨服务的长事务通过 Temporal 工作流 + Saga 模式保证最终一致性。
- 自愈优先:大多数故障通过自动重启、重连、重试解决,人工介入是最后手段。
#Next Article
下一篇 S2-12 配置管理:从 YAML 到运行时的配置链路 将讲解三个进程如何管理配置——从本地 YAML 文件到 Docker Compose 环境变量注入,再到运行时配置热更新和 Feature Flag 控制。
tags: gRPC, error-handling, circuit-breaker, retry, dead-letter-queue, Temporal, saga, fault-tolerance, graceful-degradation