返回博客

错误处理哲学:三个进程如何优雅处理故障

在分布式系统中,网络分区、进程崩溃、资源耗尽都是日常。智策平台由三个独立进程组成——onto-control(Control Layer)、onto-data(Data Layer)、onto-intelligence(Reasoning & Decision Layer + Agent Runtime Layer),它们通过 gRPC 互相调用。任何一次跨进程调用都可能失败,任何一个进程都可能宕机。

Coomia发布于 2025年7月4日20 分钟阅读
分享本文Twitter / X

错误处理哲学:三个进程如何优雅处理故障

系列: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 互相调用。任何一次跨进程调用都可能失败,任何一个进程都可能宕机。

Code
故障概率模型(简化):

单次 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 个标准状态码,平台约定了每个码的使用场景:

Code
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 错误必须携带结构化的错误详情:

PROTOBUF
// 平台统一错误详情
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 错误传播链

当一次请求跨多个进程时,错误需要沿调用链传播,同时保留每一层的上下文:

Code
错误传播示例(创建 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 拦截器,统一处理错误的捕获、包装和上报:

Code
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 重试决策树

不是所有错误都应该重试。平台定义了清晰的重试决策树:

Code
收到 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 框架实现声明式重试:

Java
// 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 注解:

Java
// 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 工作流,自带重试机制:

Python
# 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)在检测到连续失败后,会"断开"对下游的调用,直接返回降级结果。

Code
断路器状态机:

    ┌─────────┐     连续失败     ┌────────┐
    │  CLOSED  │ ──────────────→ │  OPEN   │
    │ (正常)   │   > 阈值        │ (熔断)  │
    └─────────┘                  └────────┘
         ↑                           │
         │    成功请求                │ 超时后允许
         │    > 阈值                 │ 少量请求探测
         │                           ↓
    ┌─────────────┐
    │  HALF-OPEN   │
    │ (半开/探测)   │
    └─────────────┘

参数配置:
  失败阈值:连续 5 次失败 → 打开断路器
  断路时间:30 秒后进入半开状态
  探测数量:半开状态允许 3 个请求通过
  恢复阈值:3 个请求全部成功 → 关闭断路器

#3.2 各进程的断路器实现

Code
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 中实时查看:

Code
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 降级等级

平台定义了四个降级等级,不同等级下提供不同的服务能力:

Code
降级等级矩阵:

Level  名称          触发条件                    影响
─────────────────────────────────────────────────────────
L0     全部正常       所有服务正常                 无
L1     部分降级       单个下游服务断路器打开        功能受限
L2     严重降级       2+ 下游服务不可用            只读模式
L3     最小服务       核心进程异常                 仅返回缓存

各等级可用功能:
                        L0    L1    L2    L3
Ontology 查询           ✅    ✅    ✅    ⚠️(缓存)
对象实例 CRUD           ✅    ✅    ❌    ❌
Action 执行             ✅    ⚠️    ❌    ❌
规则评估                ✅    ⚠️    ❌    ❌
派生属性                ✅    ❌    ❌    ❌
Dashboard 展示          ✅    ✅    ⚠️    ⚠️(快照)

#4.2 降级策略实现

Python
# 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 自动降级与恢复

Code
自动降级控制器(每个进程内运行):

每 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 消息消费失败时,不能简单丢弃。平台使用死信队列保证消息不丢失:

Code
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 消息结构

JSON
{
  "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 管理与回放

Code
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 提供了比普通重试更强大的容错能力:

Code
Temporal 容错特性:

1. Activity 自动重试
   ├─ 可配置重试策略(次数、间隔、退避系数)
   ├─ 可指定不可重试的异常类型
   └─ 重试对工作流代码透明

2. 工作流超时控制
   ├─ WorkflowExecutionTimeout:整个工作流的最大运行时间
   ├─ WorkflowRunTimeout:单次 run 的最大时间
   └─ ActivityStartToCloseTimeout:单个 Activity 的最大时间

3. 心跳检测
   ├─ 长时间 Activity 必须定期报告心跳
   ├─ 心跳超时 → Temporal 认为 Activity 失败 → 触发重试
   └─ 心跳可携带进度信息,重试时从断点恢复

4. 补偿(Saga 模式)
   ├─ 每个 Activity 可定义对应的补偿 Activity
   ├─ 工作流失败时按逆序执行补偿
   └─ 保证最终一致性

#6.2 Saga 补偿模式

对于涉及多个服务的操作(如执行 Action),平台使用 Saga 模式保证最终一致性:

Python
@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 工作流监控

Code
Temporal 仪表板关键指标:

工作流统计:
  ├─ 活跃工作流数量
  ├─ 工作流完成率(成功 / 总数)
  ├─ 平均工作流执行时间
  └─ 工作流失败原因分布

Activity 统计:
  ├─ Activity 执行次数(含重试)
  ├─ Activity 平均延迟
  ├─ Activity 失败率
  └─ 重试次数分布

告警规则:
  ├─ 工作流失败率 > 5% → Warning
  ├─ 工作流失败率 > 15% → Critical
  ├─ 活跃工作流数 > 1000 → Warning(可能有积压)
  └─ Activity 平均延迟 > 30s → Warning

#7. 错误恢复与自愈

#7.1 健康检查机制

每个进程都实现了多层健康检查:

Code
健康检查层次:

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 自动恢复策略

Code
自动恢复矩阵:

故障类型              检测方式              恢复策略
────────────────────────────────────────────────────────
进程 OOM             Docker 健康检查        自动重启(最多 3 次/小时)
数据库连接断开        连接池心跳             自动重连 + 指数退避
Kafka 消费者掉线      Consumer group rebalance  自动 rejoin
gRPC 连接中断         Channel 状态监控        自动重建连接
Nessie 冲突          ABORTED 错误码          自动合并重试
Schema 缓存过期       TTL + 版本号            自动刷新
Temporal Worker 崩溃  Temporal Server 检测    自动重新分配任务

#7.3 手动恢复工具

当自动恢复无法解决问题时,运维人员可以使用以下工具:

Code
运维命令(通过 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 故障场景

Code
时间线:

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 事后分析

Code
故障根因:
  一个用户提交了全表扫描查询,占用了 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

  1. 错误码标准化:统一的 gRPC 错误码 + 结构化 ErrorDetail,使错误在跨服务传播时不丢失语义。
  2. 分层重试:不同技术栈使用不同的重试实现(Spring Retry / MicroProfile FT / Temporal),但遵循相同的重试决策树。
  3. 断路器保护:防止故障级联扩散,每个进程独立管理对下游的断路器状态。
  4. 优雅降级:四级降级模型确保平台在部分故障时仍能提供有限服务,而不是完全不可用。
  5. 死信队列兜底:消息消费失败不丢弃,进入 DLQ 等待人工处理或重放。
  6. Saga 补偿:跨服务的长事务通过 Temporal 工作流 + Saga 模式保证最终一致性。
  7. 自愈优先:大多数故障通过自动重启、重连、重试解决,人工介入是最后手段。

#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