第5讲:故障转移与高可用

📅 发布时间:2026/8/22 0:11:13
第5讲:故障转移与高可用
前四讲我们构建了一个完整的分布式任务调度系统但还没有考虑故障情况。在生产环境中节点宕机、网络分区、磁盘故障都是常态。这一讲我们来给系统穿上防弹衣——实现故障转移和高可用机制。一、高可用架构总览1.1 故障类型与应对┌─────────────────────────────────────────────────────────────┐ │ 故障类型矩阵 │ ├─────────────┬─────────────────────┬─────────────────────────┤ │ 故障类型 │ 影响范围 │ 应对措施 │ ├─────────────┼─────────────────────┼─────────────────────────┤ │ Worker宕机 │ 部分任务中断 │ 任务迁移 重新调度 │ │ Scheduler宕机│ 调度中断 │ Leader选举 热备 │ │ 网络分区 │ 脑裂风险 │ 多数派决策 fencing │ │ 慢Worker │ 任务延迟 │ 熔断 降级 │ │ 磁盘满 │ 日志丢失 │ 限流 告警 │ │ 内存泄漏 │ OOM Kill │ 资源隔离 重启 │ ├─────────────┴─────────────────────┴─────────────────────────┤ │ │ │ 核心原则 │ │ • 假设一切都会失败Design for Failure │ │ • 优雅降级Graceful Degradation │ │ • 快速恢复Fast Recovery │ └─────────────────────────────────────────────────────────────┘1.2 高可用架构正常状态 ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │ Scheduler │────▶│ Worker A │────▶│ Worker B │ │ (Leader) │ │ Task 1,2,3 │ │ Task 4,5,6 │ └─────────────┘ └─────────────┘ └─────────────┘ │ ▼ ┌─────────────┐ │ Scheduler │ │ (Follower) │ └─────────────┘ Worker A 宕机后 ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │ Scheduler │────▶│ Worker B │────▶│ Worker C │ │ (Leader) │ │ Task 4,5,6 │ │ Task 1,2,3 │ │ │ │ (接管A的任务)│ │ (新加入) │ └─────────────┘ └─────────────┘ └─────────────┘ │ ▼ ┌─────────────┐ │ Scheduler │ │ (Follower) │ └─────────────┘二、心跳检测与故障发现2.1 健康检查器# fault_tolerance/heartbeat.py 心跳检测与故障发现 通过定期心跳检测Worker的健康状态。 from __future__ import annotations from typing import Dict, List, Optional, Callable, Set from dataclasses import dataclass, field import asyncio import time import logging from collections import defaultdict from coordination.registry import WorkerInfo logger logging.getLogger(__name__) dataclass class HeartbeatRecord: 心跳记录 worker_id: str last_heartbeat: float 0.0 missed_count: int 0 is_alive: bool True consecutive_failures: int 0 property def is_stale(self, threshold: float 10.0) - bool: 是否过期 return time.time() - self.last_heartbeat threshold class HealthChecker: 健康检查器 定期检测Worker健康状态发现故障并触发恢复。 def __init__(self, check_interval: float 5.0, miss_threshold: int 3, recovery_interval: float 30.0): Args: check_interval: 检查间隔秒 miss_threshold: 连续缺失次数阈值 recovery_interval: 恢复检查间隔秒 self.check_interval check_interval self.miss_threshold miss_threshold self.recovery_interval recovery_interval # 心跳记录 self._records: Dict[str, HeartbeatRecord] {} # 故障节点列表 self._failed_workers: Set[str] set() # 回调 self.on_worker_failed: Optional[Callable] None self.on_worker_recovered: Optional[Callable] None # 控制 self._running False self._check_task: Optional[asyncio.Task] None async def record_heartbeat(self, worker_id: str): 记录心跳 Args: worker_id: Worker ID if worker_id not in self._records: self._records[worker_id] HeartbeatRecord(worker_idworker_id) record self._records[worker_id] record.last_heartbeat time.time() record.missed_count 0 record.consecutive_failures 0 # 如果之前是故障状态标记为恢复 if not record.is_alive: record.is_alive True self._failed_workers.discard(worker_id) logger.info(f❤️ Worker recovered: {worker_id}) if self.on_worker_recovered: await self.on_worker_recovered(worker_id) async def _check_health(self): 健康检查循环 while self._running: now time.time() for worker_id, record in self._records.items(): # 跳过已经确认失败的节点 if worker_id in self._failed_workers: continue # 检查心跳是否过期 if now - record.last_heartbeat self.check_interval: record.missed_count 1 record.consecutive_failures 1 logger.debug(fWorker {worker_id} missed heartbeat f({record.missed_count}/{self.miss_threshold})) # 超过阈值判定为故障 if record.missed_count self.miss_threshold: await self._mark_failed(worker_id) await asyncio.sleep(self.check_interval) async def _mark_failed(self, worker_id: str): 标记Worker为故障 Args: worker_id: Worker ID record self._records.get(worker_id) if not record: return record.is_alive False self._failed_workers.add(worker_id) logger.warning(f Worker marked as failed: {worker_id} f(missed {record.missed_count} heartbeats)) if self.on_worker_failed: await self.on_worker_failed(worker_id) async def _check_recovery(self): 恢复检查循环 while self._running: for worker_id in list(self._failed_workers): record self._records.get(worker_id) if not record: continue # 如果收到新的心跳会自动恢复 # 这里只是定期清理过期的故障记录 if time.time() - record.last_heartbeat self.recovery_interval * 2: logger.info(fRemoving stale failure record: {worker_id}) self._failed_workers.discard(worker_id) self._records.pop(worker_id, None) await asyncio.sleep(self.recovery_interval) def is_alive(self, worker_id: str) - bool: 检查Worker是否存活 record self._records.get(worker_id) if not record: return False return record.is_alive and worker_id not in self._failed_workers def get_failed_workers(self) - List[str]: 获取故障Worker列表 return list(self._failed_workers) def get_healthy_workers(self) - List[str]: 获取健康Worker列表 return [ wid for wid, record in self._records.items() if self.is_alive(wid) ] async def start(self): 启动健康检查 self._running True self._check_task asyncio.create_task(self._check_health()) asyncio.create_task(self._check_recovery()) logger.info(Health checker started) async def stop(self): 停止健康检查 self._running False if self._check_task: self._check_task.cancel() logger.info(Health checker stopped)三、任务重试与幂等性3.1 重试管理器# fault_tolerance/retry.py 任务重试与幂等性 管理任务的重试逻辑确保幂等执行。 from __future__ import annotations from typing import Dict, List, Optional, Callable from dataclasses import dataclass, field import asyncio import time import logging from enum import Enum from scheduler.models.task import Task, TaskStatus logger logging.getLogger(__name__) class RetryStrategy(Enum): 重试策略 FIXED fixed # 固定间隔 EXPONENTIAL exponential # 指数退避 LINEAR linear # 线性增长 IMMEDIATE immediate # 立即重试 dataclass class RetryPolicy: 重试策略配置 定义任务的重试规则。 max_retries: int 3 strategy: RetryStrategy RetryStrategy.EXPONENTIAL base_delay_ms: float 1000 # 基础延迟毫秒 max_delay_ms: float 60000 # 最大延迟毫秒 jitter: bool True # 是否添加抖动 def calculate_delay(self, attempt: int) - float: 计算本次重试的延迟 Args: attempt: 当前重试次数从1开始 Returns: 延迟时间毫秒 if self.strategy RetryStrategy.IMMEDIATE: delay 0 elif self.strategy RetryStrategy.FIXED: delay self.base_delay_ms elif self.strategy RetryStrategy.LINEAR: delay self.base_delay_ms * attempt elif self.strategy RetryStrategy.EXPONENTIAL: delay self.base_delay_ms * (2 ** (attempt - 1)) else: delay self.base_delay_ms # 限制最大值 delay min(delay, self.max_delay_ms) # 添加抖动±25% if self.jitter and delay 0: import random jitter_range delay * 0.25 delay random.uniform(-jitter_range, jitter_range) delay max(0, delay) return delay class IdempotencyGuard: 幂等性守卫 确保任务不会被重复执行。 使用唯一ID和去重表。 def __init__(self): # 已执行的任务ID集合 self._executed_ids: set set() # 执行中的任务ID集合 self._in_progress_ids: set set() # 去重窗口秒 self.dedup_window: float 3600 # 1小时 def can_execute(self, task_id: str) - bool: 检查任务是否可以执行 Args: task_id: 任务ID Returns: 是否可以执行 if task_id in self._executed_ids: logger.warning(fTask already executed: {task_id}) return False if task_id in self._in_progress_ids: logger.warning(fTask already in progress: {task_id}) return False return True def mark_in_progress(self, task_id: str): 标记任务为执行中 self._in_progress_ids.add(task_id) def mark_executed(self, task_id: str): 标记任务为已执行 self._in_progress_ids.discard(task_id) self._executed_ids.add(task_id) def clear(self): 清理去重表 self._executed_ids.clear() self._in_progress_ids.clear() class RetryManager: 重试管理器 管理任务的重试生命周期。 def __init__(self, default_policy: RetryPolicy None): self.default_policy default_policy or RetryPolicy() self.idempotency IdempotencyGuard() # 自定义策略 self._custom_policies: Dict[str, RetryPolicy] {} # 重试统计 self.stats { total_retries: 0, successful_retries: 0, failed_retries: 0, skipped_dedup: 0 } def set_policy(self, task_type: str, policy: RetryPolicy): 为特定任务类型设置重试策略 Args: task_type: 任务类型 policy: 重试策略 self._custom_policies[task_type] policy def get_policy(self, task: Task) - RetryPolicy: 获取任务的重试策略 Args: task: 任务 Returns: 重试策略 return self._custom_policies.get( task.task_type, self.default_policy ) async def should_retry(self, task: Task) - bool: 判断是否应该重试 Args: task: 任务 Returns: 是否应该重试 # 检查幂等性 if not self.idempotency.can_execute(task.task_id): self.stats[skipped_dedup] 1 return False # 检查重试次数 retry_count task.result.retry_count if task.result else 0 policy self.get_policy(task) return retry_count policy.max_retries async def get_retry_delay(self, task: Task) - float: 获取下次重试的延迟 Args: task: 任务 Returns: 延迟时间毫秒 retry_count task.result.retry_count if task.result else 0 policy self.get_policy(task) return policy.calculate_delay(retry_count 1) async def prepare_retry(self, task: Task) - Optional[float]: 准备重试任务 Args: task: 任务 Returns: 延迟时间毫秒None表示不重试 if not await self.should_retry(task): return None delay await self.get_retry_delay(task) self.stats[total_retries] 1 logger.info(fPreparing retry for {task.name}: fattempt {(task.result.retry_count if task.result else 0) 1}, fdelay{delay:.0f}ms) return delay def on_retry_success(self, task: Task): 重试成功回调 self.stats[successful_retries] 1 self.idempotency.mark_executed(task.task_id) def on_retry_failure(self, task: Task): 重试失败回调 self.stats[failed_retries] 1 def get_stats(self) - dict: 获取统计信息 return { **self.stats, dedup_table_size: len(self.idempotency._executed_ids) }四、故障转移引擎4.1 故障转移实现# fault_tolerance/failover.py 故障转移引擎 当Worker发生故障时自动迁移任务到健康的Worker。 from __future__ import annotations from typing import Dict, List, Optional, Set, Callable from dataclasses import dataclass, field import asyncio import time import logging from collections import defaultdict from scheduler.models.task import Task, TaskStatus from coordination.registry import WorkerInfo from fault_tolerance.heartbeat import HealthChecker from fault_tolerance.retry import RetryManager logger logging.getLogger(__name__) dataclass class TaskMigration: 任务迁移记录 task_id: str source_worker: str target_worker: Optional[str] None migrated_at: float 0.0 status: str pending # pending/migrating/completed/failed class FailoverEngine: 故障转移引擎 监控Worker健康状态在故障时自动迁移任务。 def __init__(self, health_checker: HealthChecker, retry_manager: RetryManager): self.health_checker health_checker self.retry_manager retry_manager # Worker → 任务映射 self._worker_tasks: Dict[str, Set[str]] defaultdict(set) # 任务 → Worker映射 self._task_workers: Dict[str, str] {} # 迁移记录 self._migrations: Dict[str, TaskMigration] {} # 回调 self.on_task_migrating: Optional[Callable] None self.on_task_migrated: Optional[Callable] None # 统计 self.stats { total_failovers: 0, successful_failovers: 0, failed_failovers: 0, tasks_moved: 0 } # 设置健康检查回调 self.health_checker.on_worker_failed self._on_worker_failed self.health_checker.on_worker_recovered self._on_worker_recovered def register_task(self, task: Task, worker_id: str): 注册任务到Worker Args: task: 任务 worker_id: Worker ID self._worker_tasks[worker_id].add(task.task_id) self._task_workers[task.task_id] worker_id def unregister_task(self, task_id: str): 注销任务 Args: task_id: 任务ID worker_id self._task_workers.pop(task_id, None) if worker_id: self._worker_tasks[worker_id].discard(task_id) async def _on_worker_failed(self, worker_id: str): Worker故障回调 Args: worker_id: 故障Worker ID logger.warning(f Initiating failover for worker: {worker_id}) # 获取该Worker上的所有任务 affected_tasks list(self._worker_tasks.get(worker_id, set())) if not affected_tasks: logger.info(fNo tasks to migrate from {worker_id}) return self.stats[total_failovers] len(affected_tasks) # 迁移任务 for task_id in affected_tasks: migration TaskMigration( task_idtask_id, source_workerworker_id, migrated_attime.time() ) self._migrations[task_id] migration if self.on_task_migrating: await self.on_task_migrating(task_id, worker_id) logger.info(fFailover initiated: {len(affected_tasks)} tasks ffrom {worker_id}) async def migrate_task(self, task_id: str, available_workers: List[WorkerInfo]) - bool: 迁移单个任务 Args: task_id: 任务ID available_workers: 可用Worker列表 Returns: 是否成功 migration self._migrations.get(task_id) if not migration: return False # 选择目标Worker排除源Worker candidates [ w for w in available_workers if w.worker_id ! migration.source_worker ] if not candidates: logger.warning(fNo candidate workers for task {task_id}) migration.status failed self.stats[failed_failovers] 1 return False # 选择负载最低的Worker target min(candidates, keylambda w: w.current_load) migration.target_worker target.worker_id migration.status migrating logger.info(fMigrating task {task_id}: f{migration.source_worker} → {target.worker_id}) # 更新映射 self._worker_tasks[migration.source_worker].discard(task_id) self._worker_tasks[target.worker_id].add(task_id) self._task_workers[task_id] target.worker_id migration.status completed self.stats[successful_failovers] 1 self.stats[tasks_moved] 1 if self.on_task_migrated: await self.on_task_migrated(task_id, target.worker_id) return True async def _on_worker_recovered(self, worker_id: str): Worker恢复回调 Args: worker_id: 恢复的Worker ID logger.info(f Worker recovered: {worker_id}) # 可以选择将部分任务迁回原Worker # 这里简化处理不做自动迁回 def get_worker_tasks(self, worker_id: str) - List[str]: 获取Worker上的任务列表 return list(self._worker_tasks.get(worker_id, set())) def get_task_worker(self, task_id: str) - Optional[str]: 获取任务所在的Worker return self._task_workers.get(task_id) def get_stats(self) - dict: 获取统计信息 return { **self.stats, active_migrations: len([ m for m in self._migrations.values() if m.status migrating ]), total_workers: len(self._worker_tasks) }五、熔断与降级5.1 断路器# fault_tolerance/circuit_breaker.py 断路器模式 防止故障扩散保护系统稳定性。 from __future__ import annotations from typing import Dict, Optional, Callable from dataclasses import dataclass import asyncio import time import logging from enum import Enum logger logging.getLogger(__name__) class CircuitState(Enum): 断路器状态 CLOSED closed # 正常工作 OPEN open # 断开 HALF_OPEN half_open # 半开试探 dataclass class CircuitBreakerConfig: 断路器配置 failure_threshold: int 5 # 失败阈值 success_threshold: int 2 # 半开状态下成功阈值 open_timeout_ms: float 30000 # 断开超时毫秒 half_open_timeout_ms: float 5000 # 半开超时 class CircuitBreaker: 断路器 监控失败率当达到阈值时断开电路避免雪崩。 def __init__(self, name: str, config: CircuitBreakerConfig None): self.name name self.config config or CircuitBreakerConfig() self.state CircuitState.CLOSED self.failure_count 0 self.success_count 0 self.last_state_change time.time() # 回调 self.on_open: Optional[Callable] None self.on_close: Optional[Callable] None self.on_half_open: Optional[Callable] None async def call(self, func: Callable, *args, **kwargs): 执行受保护的调用 Args: func: 要执行的函数 Returns: 函数返回值 Raises: CircuitBreakerOpenError: 断路器打开时 if self.state CircuitState.OPEN: if self._should_attempt_reset(): self._set_state(CircuitState.HALF_OPEN) else: raise CircuitBreakerOpenError( fCircuit breaker {self.name} is OPEN ) try: result await func(*args, **kwargs) self._on_success() return result except Exception as e: self._on_failure() raise def _on_success(self): 成功回调 if self.state CircuitState.HALF_OPEN: self.success_count 1 if self.success_count self.config.success_threshold: self._set_state(CircuitState.CLOSED) else: self.failure_count 0 def _on_failure(self): 失败回调 self.failure_count 1 if self.state CircuitState.HALF_OPEN: self._set_state(CircuitState.OPEN) elif self.failure_count self.config.failure_threshold: self._set_state(CircuitState.OPEN) def _set_state(self, state: CircuitState): 设置状态 old_state self.state self.state state self.last_state_change time.time() logger.info(fCircuit breaker {self.name}: f{old_state.value} → {state.value}) if state CircuitState.OPEN and self.on_open: self.on_open() elif state CircuitState.CLOSED and self.on_close: self.on_close() elif state CircuitState.HALF_OPEN and self.on_half_open: self.on_half_open() def _should_attempt_reset(self) - bool: 是否应该尝试重置 elapsed (time.time() - self.last_state_change) * 1000 return elapsed self.config.open_timeout_ms def reset(self): 手动重置断路器 self._set_state(CircuitState.CLOSED) self.failure_count 0 self.success_count 0 class CircuitBreakerOpenError(Exception): 断路器打开异常 pass class CircuitBreakerRegistry: 断路器注册表 管理多个断路器实例。 def __init__(self): self._breakers: Dict[str, CircuitBreaker] {} def get_or_create(self, name: str, config: CircuitBreakerConfig None) - CircuitBreaker: 获取或创建断路器 if name not in self._breakers: self._breakers[name] CircuitBreaker(name, config) return self._breakers[name] def get_all_open(self) - list: 获取所有打开的断路器 return [ b for b in self._breakers.values() if b.state CircuitState.OPEN ] def reset_all(self): 重置所有断路器 for breaker in self._breakers.values(): breaker.reset()六、集成与演示6.1 完整故障转移演示# examples/fault_tolerance_demo.py 故障转移与高可用演示 import asyncio import logging import sys import time import random sys.path.insert(0, ..) from scheduler.models.task import Task, TaskStatus from fault_tolerance.heartbeat import HealthChecker from fault_tolerance.retry import RetryManager, RetryPolicy, RetryStrategy from fault_tolerance.failover import FailoverEngine from fault_tolerance.circuit_breaker import CircuitBreaker, CircuitBreakerConfig logging.basicConfig(levellogging.INFO) async def demo_heartbeat(): 演示心跳检测 print( * 71) print( 心跳检测演示) print( * 71) checker HealthChecker(check_interval1.0, miss_threshold3) failed_workers [] recovered_workers [] async def on_fail(worker_id): failed_workers.append(worker_id) print(f 检测到故障: {worker_id}) async def on_recover(worker_id): recovered_workers.append(worker_id) print(f ❤️ Worker恢复: {worker_id}) checker.on_worker_failed on_fail checker.on_worker_recovered on_recover await checker.start() # 模拟Worker心跳 print(\nWorker正常心跳...) for i in range(5): await checker.record_heartbeat(worker-1) await asyncio.sleep(0.5) # 模拟Worker宕机 print(\nWorker-1 宕机停止心跳...) await asyncio.sleep(5) # 模拟恢复 print(\nWorker-1 恢复...) await checker.record_heartbeat(worker-1) await asyncio.sleep(1) await checker.stop() print(f\n检测结果:) print(f 故障Worker: {failed_workers}) print(f 恢复Worker: {recovered_workers}) async def demo_retry(): 演示重试机制 print(\n * 71) print( 重试机制演示) print( * 71) manager RetryManager() # 设置指数退避策略 policy RetryPolicy( max_retries3, strategyRetryStrategy.EXPONENTIAL, base_delay_ms500, jitterFalse ) print(f\n重试策略: {policy.strategy.value}) print(f最大重试: {policy.max_retries}) print(f基础延迟: {policy.base_delay_ms}ms) # 模拟失败任务 task Task(nameflaky-task) task.mark_running(worker-1) task.mark_failed(Connection timeout) print(f\n任务初始状态: {task.status.value}) print(f重试次数: {task.result.retry_count if task.result else 0}) # 模拟重试 for attempt in range(1, 4): delay await manager.get_retry_delay(task) print(f\n第{attempt}次重试:) print(f 延迟: {delay:.0f}ms) # 模拟执行 await asyncio.sleep(delay / 1000) task.mark_running(worker-1) task.mark_failed(fAttempt {attempt} failed) print(f 结果: {task.status.value}) print(f 已重试: {task.result.retry_count}) print(f\n最终状态: {task.status.value}) print(f统计: {manager.get_stats()}) async def demo_failover(): 演示故障转移 print(\n * 71) print( 故障转移演示) print( * 71) checker HealthChecker(check_interval1.0, miss_threshold2) retry_manager RetryManager() failover FailoverEngine(checker, retry_manager) # 注册任务到Worker print(\n注册任务到Worker...) for i in range(5): task Task(nameftask-{i1}) failover.register_task(task, worker-1) print(f {task.name} → worker-1) # 模拟Worker心跳 print(\nWorker正常心跳...) for i in range(3): await checker.record_heartbeat(worker-1) await asyncio.sleep(0.5) # 模拟Worker宕机 print(\n⚠️ Worker-1 宕机!) await asyncio.sleep(3) # 查看故障转移 print(f\n故障Worker: {checker.get_failed_workers()}) print(f待迁移任务: {failover.get_worker_tasks(worker-1)}) # 模拟迁移 print(\n执行任务迁移...) from coordination.registry import WorkerInfo available [ WorkerInfo(worker_idworker-2, host, port9002), WorkerInfo(worker_idworker-3, host, port9003), ] for task_id in list(failover.get_worker_tasks(worker-1)): success await failover.migrate_task(task_id, available) print(f {✅ if success else ❌} {task_id}) print(f\n转移统计: {failover.get_stats()}) async def demo_circuit_breaker(): 演示断路器 print(\n * 71) print(⚡ 断路器演示) print( * 71) breaker CircuitBreaker( worker-api, CircuitBreakerConfig( failure_threshold3, success_threshold2, open_timeout_ms5000 ) ) # 模拟失败请求 print(\n模拟连续失败...) call_count 0 async def failing_call(): nonlocal call_count call_count 1 raise ConnectionError(fRequest failed #{call_count}) for i in range(5): try: await breaker.call(failing_call) except CircuitBreakerOpenError: print(f ⛔ 断路器打开请求被拒绝 (第{i1}次)) except ConnectionError as e: print(f ❌ 请求失败: {e}) if breaker.state.value open: print(f 断路器已打开!) # 等待恢复 print(f\n等待断路器半开...) await asyncio.sleep(5) # 模拟成功请求 print(f\n尝试恢复...) async def success_call(): return OK for i in range(3): try: result await breaker.call(success_call) print(f ✅ 请求成功: {result} (状态: {breaker.state.value})) except CircuitBreakerOpenError: print(f ⛔ 断路器仍打开) print(f\n最终状态: {breaker.state.value}) async def main(): await demo_heartbeat() await demo_retry() await demo_failover() await demo_circuit_breaker() if __name__ __main__: asyncio.run(main())七、测试7.1 故障转移测试# tests/test_fault_tolerance.py import pytest import asyncio from fault_tolerance.heartbeat import HealthChecker from fault_tolerance.retry import RetryManager, RetryPolicy, RetryStrategy from fault_tolerance.failover import FailoverEngine from fault_tolerance.circuit_breaker import CircuitBreaker, CircuitBreakerOpenError from scheduler.models.task import Task, TaskStatus class TestHealthChecker: 健康检查测试 pytest.mark.asyncio async def test_detect_failure(self): checker HealthChecker(check_interval0.1, miss_threshold3) failures [] async def on_fail(wid): failures.append(wid) checker.on_worker_failed on_fail await checker.start() # 发送一次心跳然后停止 await checker.record_heartbeat(worker-1) await asyncio.sleep(0.5) await checker.stop() assert worker-1 in failures def test_is_alive(self): checker HealthChecker() assert not checker.is_alive(unknown) class TestRetryManager: 重试测试 pytest.mark.asyncio async def test_exponential_backoff(self): policy RetryPolicy( max_retries3, strategyRetryStrategy.EXPONENTIAL, base_delay_ms100, jitterFalse ) delays [] for i in range(3): delay policy.calculate_delay(i 1) delays.append(delay) assert delays [100, 200, 400] def test_idempotency(self): guard RetryManager().idempotency assert guard.can_execute(task-1) guard.mark_executed(task-1) assert not guard.can_execute(task-1) class TestCircuitBreaker: 断路器测试 pytest.mark.asyncio async def test_open_on_failures(self): breaker CircuitBreaker( test, config{failure_threshold: 3, open_timeout_ms: 10000} ) async def fail(): raise ValueError(fail) for i in range(3): with pytest.raises(ValueError): await breaker.call(fail) assert breaker.state.value open pytest.mark.asyncio async def test_half_open(self): breaker CircuitBreaker( test, config{failure_threshold: 2, open_timeout_ms: 100} ) async def fail(): raise ValueError(fail) async def succeed(): return ok # 触发打开 for i in range(2): with pytest.raises(ValueError): await breaker.call(fail) assert breaker.state.value open # 等待半开 await asyncio.sleep(0.2) # 应该进入半开状态 result await breaker.call(succeed) assert result ok if __name__ __main__: pytest.main([__file__, -v])八、总结8.1 本讲成果组件文件功能HealthChecker​fault_tolerance/heartbeat.py心跳检测与故障发现RetryManager​fault_tolerance/retry.py重试管理与幂等性FailoverEngine​fault_tolerance/failover.py故障转移引擎CircuitBreaker​fault_tolerance/circuit_breaker.py断路器模式8.2 核心知识点心跳检测定期探测连续失败阈值判定重试策略指数退避 抖动防止惊群效应幂等性唯一ID去重确保Exactly-Once语义故障转移自动迁移任务到健康节点断路器快速失败防止雪崩效应8.3 下一讲预告第6讲持久化与历史我们将实现数据的持久化存储任务存储PostgreSQL执行日志历史记录查询数据归档策略准备好了吗让我们在第6讲再见开发之余的小工具推荐​处理 Base64、JWT 解析、JSON 格式化、Crontab 计算、PDF 合并压缩这些碎片需求我常用一个纯前端本地工具箱zz365.top子页 PDF 大师PDF 大师 - zz365工具箱。所有计算在浏览器完成文件不上传服务器关页即清。免费、无登录、无广告适合开发者当常驻标签页。