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

相关新闻

BETAFPV Configurator:5 分钟配好你的 Whoop 飞控

BETAFPV Configurator:5 分钟配好你的 Whoop 飞控

BETAFPV Configurator:5 分钟配好你的 Whoop 飞控 【免费下载链接】BETAFPV_Configurator 项目地址: https://gitcode.com/gh_mirrors/be/BETAFPV_Configurator 拿到 Cetus 的第一关不是飞,是连不上电脑:COM 口找不到、固件刷不进、E…

2026/8/22 0:08:13 阅读更多 →
QuickCut:视频压缩、切段、合并一步到位的免费工具——新手指南

QuickCut:视频压缩、切段、合并一步到位的免费工具——新手指南

QuickCut:视频压缩、切段、合并一步到位的免费工具——新手指南 【免费下载链接】QuickCut Your most handy video processing software 项目地址: https://gitcode.com/gh_mirrors/qu/QuickCut QuickCut 是一款免费开源的视频处理软件,底层引擎是…

2026/8/22 0:08:13 阅读更多 →
Cellpose-SAM 细胞分割完全指南:从安装到自定义模型一次讲清

Cellpose-SAM 细胞分割完全指南:从安装到自定义模型一次讲清

Cellpose-SAM 细胞分割完全指南:从安装到自定义模型一次讲清 【免费下载链接】cellpose a generalist algorithm for cellular segmentation with human-in-the-loop capabilities 项目地址: https://gitcode.com/gh_mirrors/ce/cellpose Cellpose 是一款通用…

2026/8/22 0:08:13 阅读更多 →

最新新闻

AI写作工具助力毕业生高效应对求职与学术挑战

AI写作工具助力毕业生高效应对求职与学术挑战

1. 项目背景与需求分析2025届毕业生正面临前所未有的就业竞争压力,学术论文、求职简历、实习报告等各类文书写作需求激增。在这个时间节点上,AI辅助写作工具已经从早期的简单语法检查,发展到能够提供内容生成、结构优化、风格调整等全方位支持…

2026/8/23 3:20:09 阅读更多 →
片上网络仿真器Noxim:从零入门到架构探索与性能分析

片上网络仿真器Noxim:从零入门到架构探索与性能分析

1. 项目概述:从零认识片上网络仿真器 如果你正在研究芯片设计,尤其是多核处理器或大规模片上系统(SoC),那么“仿真”这个词对你来说一定不陌生。在硬件真正流片之前,我们如何验证一个由几十甚至上百个核心…

2026/8/23 3:20:09 阅读更多 →
Markdown样式定制全攻略:从CSS基础到工具链实战

Markdown样式定制全攻略:从CSS基础到工具链实战

1. 从“能用”到“好看”:为什么我们需要折腾Markdown样式?如果你和我一样,是个重度Markdown使用者,从写技术文档、记笔记到写博客,都离不开它,那你肯定也经历过这个阶段:一开始,你满…

2026/8/23 3:20:09 阅读更多 →
支持向量机(SVM)原理详解:从间隔最大化到核技巧实战

支持向量机(SVM)原理详解:从间隔最大化到核技巧实战

1. 项目概述:从“分界线”到“最优解”的思维跃迁如果你在数据科学或者机器学习领域摸爬滚打过一阵子,肯定对“分类”这个任务不陌生。我们手里有一堆数据点,每个点都带着一堆特征,然后被贴上了不同的标签,比如“垃圾邮…

2026/8/23 3:20:09 阅读更多 →
图像预处理中的CenterCrop:原理、实现与实战避坑指南

图像预处理中的CenterCrop:原理、实现与实战避坑指南

1. 从“为什么需要CenterCrop”说起在图像处理或者深度学习模型训练的前期,我们经常会遇到一个看似简单却至关重要的步骤:图像裁剪。你可能已经用过torchvision.transforms.CenterCrop,或者OpenCV里类似的函数,感觉就是“把图片中…

2026/8/23 3:20:09 阅读更多 →
Linux磁盘管理:安全卸载与格式化分区的完整指南

Linux磁盘管理:安全卸载与格式化分区的完整指南

1. 为什么“卸载分区”不等于“删除文件”?在Linux系统管理中,“卸载分区”和“格式化分区”是两个经常被混淆,但内核逻辑完全不同的操作。很多新手,甚至一些有经验的用户,在命令行前敲下umount或mkfs时,心…

2026/8/23 3:19:09 阅读更多 →

日新闻

[光学原理与应用-521]:对光的错误理解与纠偏

[光学原理与应用-521]:对光的错误理解与纠偏

首先光是一种能量的载体和形态,宏观上观察到的光是由无数个微观的光量子组成的,每个光子在产生的瞬间,其在真空的空间中以确定不变的速度沿着一个初始的方向一直向前,在微观层面,每个光量子的运动轨迹是以波函数所展现…

2026/8/23 0:00:50 阅读更多 →
SIP通话转接原理与REFER方法实战解析

SIP通话转接原理与REFER方法实战解析

1. 通话转接不是“挂断再拨号”,而是SIP会话的动态重定向你有没有遇到过这样的场景:客服坐席A正在和客户通电话,突然需要把这通对话无缝转给专家坐席B,客户完全感知不到中间的断连——既没听到忙音,也没被要求重新拨号…

2026/8/23 0:00:50 阅读更多 →
Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

1. 为什么选择Kolla-ansible来部署单节点OpenStack?如果你正在寻找一种能把OpenStack从“概念”快速变成“可用的实验环境”的方法,那么Kolla-ansible几乎是当前最主流、最省心的选择。我见过太多人卡在手动编译依赖、配置服务、处理版本冲突的泥潭里&am…

2026/8/23 0:00:50 阅读更多 →

周新闻

[光学原理与应用-521]:对光的错误理解与纠偏

[光学原理与应用-521]:对光的错误理解与纠偏

首先光是一种能量的载体和形态,宏观上观察到的光是由无数个微观的光量子组成的,每个光子在产生的瞬间,其在真空的空间中以确定不变的速度沿着一个初始的方向一直向前,在微观层面,每个光量子的运动轨迹是以波函数所展现…

2026/8/23 0:00:50 阅读更多 →
SIP通话转接原理与REFER方法实战解析

SIP通话转接原理与REFER方法实战解析

1. 通话转接不是“挂断再拨号”,而是SIP会话的动态重定向你有没有遇到过这样的场景:客服坐席A正在和客户通电话,突然需要把这通对话无缝转给专家坐席B,客户完全感知不到中间的断连——既没听到忙音,也没被要求重新拨号…

2026/8/23 0:00:50 阅读更多 →
Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

1. 为什么选择Kolla-ansible来部署单节点OpenStack?如果你正在寻找一种能把OpenStack从“概念”快速变成“可用的实验环境”的方法,那么Kolla-ansible几乎是当前最主流、最省心的选择。我见过太多人卡在手动编译依赖、配置服务、处理版本冲突的泥潭里&am…

2026/8/23 0:00:50 阅读更多 →

月新闻

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南 【免费下载链接】BaiduNetdiskPlugin-macOS For macOS.百度网盘 破解SVIP、下载速度限制~ 项目地址: https://gitcode.com/gh_mirrors/ba/BaiduNetdiskPlugin-macOS 还在为百度网盘macOS版的龟速下…

2026/8/22 18:08:39 阅读更多 →
终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换 【免费下载链接】ncmdump 项目地址: https://gitcode.com/gh_mirrors/ncmd/ncmdump 还在为网易云音乐下载的NCM格式文件无法在其他播放器播放而烦恼吗?ncmdump解密工具帮你轻松解决这个困…

2026/8/22 7:31:03 阅读更多 →
HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

AgentCard 智能体卡片:为英语学习 App 打造桌面级学习助手适用平台:HarmonyOS 7.0 (API 26 Beta)一、引言 HarmonyOS 7.0(API 26 Beta)新增了 AgentCard 智能体卡片能力,这是继 HMAF(鸿蒙智能体框架&#x…

2026/8/22 3:22:48 阅读更多 →