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

发布时间:2026/8/22 0:10:52

第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工具箱。所有计算在浏览器完成文件不上传服务器关页即清。免费、无登录、无广告适合开发者当常驻标签页。

相关新闻

BEHAVIOR-1K具身智能仿真平台:1000项家务任务基准的完整剖析

BEHAVIOR-1K具身智能仿真平台:1000项家务任务基准的完整剖析

2026/8/22 0:10:52

BEHAVIOR-1K具身智能仿真平台:1000项家务任务基准的完整剖析 【免费下载链接】BEHAVIOR-1K BEHAVIOR-1K: a platform for accelerating Embodied AI research. Join our Discord for support: https://discord.gg/bccR5vGFEx 项目地址: https://gitcode.com/gh_mi…

降AIGC神器实测!AI率92%暴降至5%!实测10款AI智能降重工具!免费降AIGC额度薅到爽!

降AIGC神器实测!AI率92%暴降至5%!实测10款AI智能降重工具!免费降AIGC额度薅到爽!

2026/8/22 0:10:52

2026 年各大高校和期刊平台的 AI 检测系统又升级了,知网 AIGC、维普 AI、万方智能检测三大平台的算法迭代速度越来越快,上个月能蒙混过关的改写方式,这个月直接就会被标红预警。单纯的同义词替换、语序调整早就不管用了,想要有效降…

消除AI代码的“AI味”:Claude Code设计优化技能配置与实战指南

消除AI代码的“AI味”:Claude Code设计优化技能配置与实战指南

2026/8/22 0:00:52

大家好,我是专注于前端开发与AI工具实践的技术博主。在日常使用 Claude Code 等AI编程助手时,你是否也遇到过这样的困扰:生成的代码功能上没问题,但代码风格、组件设计、交互逻辑总透着一股“AI味”——布局单调、样式简陋、交互生…

WindowsWorld:以进程为中心的GUI智能体基准测试设计与实践

WindowsWorld:以进程为中心的GUI智能体基准测试设计与实践

2026/8/22 1:50:57

1. 项目概述:为什么我们需要一个“以进程为中心”的GUI智能体基准测试? 如果你最近关注AI Agent领域,尤其是那些号称能“像人一样操作电脑”的自主GUI智能体,可能会发现一个有趣的现象:大多数演示和基准测试都聚焦在单…

Stacking集成学习建模用户体验影响因素

Stacking集成学习建模用户体验影响因素

2026/8/22 1:50:57

1. 项目概述:从真实业务痛点出发的建模实践北京移动作为华北地区用户规模最大的通信运营商之一,每天产生数千万级的网络信令、投诉工单、APP行为日志和基站性能指标。2022年MathorCup大赛B题抛出一个直击行业核心的问题:哪些因素真正影响了用…

Java全栈开发工程师面试核心考点与实战解析

Java全栈开发工程师面试核心考点与实战解析

2026/8/22 1:50:57

1. Java全栈开发工程师面试的核心考察维度Java全栈开发工程师的面试从来不是简单的技术问答,而是一场对候选人综合能力的全面检验。作为面试过数百名Java开发者的技术负责人,我发现大多数候选人失败的原因往往不是技术深度不够,而是对全栈工程…

牛客刷题进度追踪器:算法学习与面试准备利器

牛客刷题进度追踪器:算法学习与面试准备利器

2026/8/22 1:50:57

1. 项目概述:牛客刷题进度追踪器牛客tracker是一款专为算法学习者设计的刷题进度管理工具,它深度整合牛客网的题库资源,通过可视化的方式帮助用户追踪每日刷题进度、记录成就并分析学习轨迹。这个工具特别适合准备技术面试或参与编程竞赛的开…

FitGirl Repack Launcher 完全指南:从搜索游戏到一键启动

FitGirl Repack Launcher 完全指南:从搜索游戏到一键启动

2026/8/22 1:50:57

FitGirl Repack Launcher 完全指南:从搜索游戏到一键启动 【免费下载链接】Fitgirl-Repack-Launcher An Electron launcher designed specifically for FitGirl Repacks, utilizing pure vanilla JavaScript, HTML, and CSS for optimal performance and customizat…

多智能体AI编程协调能力评估框架:从原理到实践

多智能体AI编程协调能力评估框架:从原理到实践

2026/8/22 1:40:56

这次我们来看一个关于多智能体AI编程协调能力评估的研究项目。它不是一个新的代码生成工具,而是一个评估框架,旨在解决一个核心问题:当多个AI智能体协作完成一个复杂的编程任务时,它们之间的“协调”能力如何衡量?这个…

【文章复现】非线性值迭代自适应动态规划(ADP):离散时间非线性系统的策略迭代自适应动态规划算法研究附Matlab代码

【文章复现】非线性值迭代自适应动态规划(ADP):离散时间非线性系统的策略迭代自适应动态规划算法研究附Matlab代码

2026/8/21 21:41:19

✅作者简介:热爱科研的Matlab仿真开发者,擅长毕业设计辅导、数学建模、数据处理、建模仿真、程序设计、完整代码获取、论文复现及科研仿真。🍎 往期回顾关注个人主页:Matlab科研工作室👇 关注我领取海量matlab电子书和…

【双层规划,节点出清价,绿证交易,CVaR方法】两级电力市场环境下计及风险的省间交易商最优购电模型附Matlab代码

【双层规划,节点出清价,绿证交易,CVaR方法】两级电力市场环境下计及风险的省间交易商最优购电模型附Matlab代码

2026/8/20 21:07:35

✅作者简介:热爱科研的Matlab仿真开发者,擅长毕业设计辅导、数学建模、数据处理、建模仿真、程序设计、完整代码获取、论文复现及科研仿真。🍎 往期回顾关注个人主页:Matlab科研工作室👇 关注我领取海量matlab电子书和…

隐式mpc+自适应mpc+时变mpc,线性时变模型预测控制附Simulink仿真

隐式mpc+自适应mpc+时变mpc,线性时变模型预测控制附Simulink仿真

2026/8/19 8:02:16

✅作者简介:热爱科研的Matlab仿真开发者,擅长毕业设计辅导、数学建模、数据处理、建模仿真、程序设计、完整代码获取、论文复现及科研仿真。🍎 往期回顾关注个人主页:Matlab科研工作室👇 关注我领取海量matlab电子书和…

多尺度智能体控制:从宏观密度场到微观决策的架构与实践

多尺度智能体控制:从宏观密度场到微观决策的架构与实践

2026/8/22 0:00:52

1. 从宏观到微观:多尺度智能体控制的核心挑战在智能体(Agent)技术日益普及的今天,我们面临着一个越来越普遍的难题:如何同时管理成千上万个,甚至百万级别的智能体?无论是城市交通中的自动驾驶车…

CUBE标准:统一AI智能体评测的度量衡与架构解析

CUBE标准:统一AI智能体评测的度量衡与架构解析

2026/8/22 0:00:52

1. 项目概述:为什么我们需要一个统一的智能体评测标准?最近在折腾各种AI智能体项目,从简单的自动化脚本到复杂的多模态交互系统,我发现了一个让人头疼的共性问题:评测。每次开发完一个智能体,想看看它到底行…

沉金PCB工艺实战指南:从设计到SMT焊接的可靠性保障

沉金PCB工艺实战指南:从设计到SMT焊接的可靠性保障

2026/8/22 0:00:52

在电子硬件开发领域,PCB(印制电路板)的沉金工艺是提升产品可靠性和焊接质量的关键环节。对于需要高密度互连、长期稳定运行或高频信号传输的板卡,如“黍姐仿通行证”这类可能涉及身份识别、数据交互的硬件项目,选择正确…

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

2026/8/17 12:00:53

一天写完毕业论文在2026年已不再是天方夜谭。2026年最炸裂、实测能大幅提速的AI论文写作工具,覆盖选题构思、文献整理、内容生成、格式排版等核心场景,真正帮你高效搞定论文难题。 一、全流程王者:一站式搞定论文全链路(一天定稿首…

导师推荐!2026最新AI论文工具测评与实用推荐

导师推荐!2026最新AI论文工具测评与实用推荐

2026/8/15 10:10:27

2026年真正好用的AI论文工具,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。 一、…

告别游戏崩溃:XCOM 2模组管理器的智能革命

告别游戏崩溃:XCOM 2模组管理器的智能革命

2026/8/22 1:32:34

告别游戏崩溃:XCOM 2模组管理器的智能革命 【免费下载链接】xcom2-launcher The Alternative Mod Launcher (AML) is a replacement for the default game launchers from XCOM 2 and XCOM Chimera Squad. 项目地址: https://gitcode.com/gh_mirrors/xc/xcom2-lau…