ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

第8讲:高级特性

2026/8/25 14:12:44 拓冰建站 浏览量
第8讲:高级特性 前七讲我们构建了一个完整的分布式任务调度系统具备了调度、执行、高可用、持久化和监控能力。这一讲我们来为系统添加一些高级特性让它更加强大和灵活。一、高级特性总览1.1 特性清单┌─────────────────────────────────────────────────────────────┐ │ 高级特性矩阵 │ ├─────────────┬─────────────────────┬─────────────────────────┤ │ 特性 │ 用途 │ 实现难度 │ ├─────────────┼─────────────────────┼─────────────────────────┤ │ Cron表达式 │ 周期性任务调度 │ ⭐⭐ │ │ DAG工作流 │ 复杂任务编排 │ ⭐⭐⭐ │ │ 动态分片 │ 大数据量并行处理 │ ⭐⭐⭐⭐ │ │ 任务依赖 │ 条件触发 │ ⭐⭐ │ │ 动态优先级 │ 运行时调整优先级 │ ⭐ │ │ 任务标签 │ 分类与路由 │ ⭐ │ │ 灰度发布 │ 逐步放量 │ ⭐⭐⭐ │ │ 任务克隆 │ 快速复制 │ ⭐ │ └─────────────┴─────────────────────┴─────────────────────────┘1.2 高级调度场景场景1Cron周期性任务 ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │ 每天凌晨2点 │───▶│ 每周一上午 │───▶│ 每月1号 │ │ 数据备份 │ │ 报表生成 │ │ 账单结算 │ └─────────────┘ └─────────────┘ └─────────────┘ 场景2DAG工作流 数据抽取 ──▶ 数据清洗 ──▶ 数据转换 ──▶ 数据加载 │ │ ▼ ▼ 发送通知 ◀──────────────── 校验质量 场景3动态分片 ┌─────────────────────────────────────────────┐ │ 100万条数据 │ │ ┌────┬────┬────┬────┬────┬────┬────┬────┐ │ │ │ S1 │ S2 │ S3 │ S4 │ S5 │ S6 │ S7 │ S8 │ │ │ └─┬──┴─┬──┴─┬──┴─┬──┴─┬──┴─┬──┴─┬──┴─┬──┘ │ │ │ │ │ │ │ │ │ │ │ │ ▼ ▼ ▼ ▼ ▼ ▼ ▼ ▼ │ │ W1 W2 W3 W4 W5 W6 W7 W8 │ └─────────────────────────────────────────────┘二、Cron表达式支持2.1 Cron解析器# advanced/cron.py Cron表达式解析器 支持标准Cron表达式和增强语法。 from __future__ import annotations from typing import List, Optional, Set, Tuple from dataclasses import dataclass, field import calendar import re from datetime import datetime, timedelta, timezone class CronField: Cron字段解析基类 def __init__(self, expression: str, min_val: int, max_val: int): self.min_val min_val self.max_val max_val self.values: Set[int] set() self._parse(expression) def _parse(self, expression: str): 解析表达式 if expression *: self.values set(range(self.min_val, self.max_val 1)) return parts expression.split(,) for part in parts: self._parse_part(part) def _parse_part(self, part: str): 解析单个部分 # 步进: */5, 1/10 if / in part: base, step part.split(/) step int(step) if base *: start self.min_val elif - in base: start, end map(int, base.split(-)) self.values.update(range(start, end 1, step)) return else: start int(base) self.values.update(range(start, self.max_val 1, step)) return # 范围: 1-5 if - in part: start, end map(int, part.split(-)) self.values.update(range(start, end 1)) return # 具体值: 3 self.values.add(int(part)) def contains(self, value: int) - bool: return value in self.values class MinuteField(CronField): def __init__(self, expr: str): super().__init__(expr, 0, 59) class HourField(CronField): def __init__(self, expr: str): super().__init__(expr, 0, 23) class DayOfMonthField(CronField): def __init__(self, expr: str): super().__init__(expr, 1, 31) class MonthField(CronField): def __init__(self, expr: str): super().__init__(expr, 1, 12) class DayOfWeekField(CronField): def __init__(self, expr: str): super().__init__(expr, 0, 6) # 0Sunday dataclass class CronExpression: Cron表达式 标准格式: 分 时 日 月 周 示例: 0 2 * * 1 每周一凌晨2点 minute: str * hour: str * day_of_month: str * month: str * day_of_week: str * # 解析后的字段 _minute_field: MinuteField field(initFalse, reprFalse) _hour_field: HourField field(initFalse, reprFalse) _day_field: DayOfMonthField field(initFalse, reprFalse) _month_field: MonthField field(initFalse, reprFalse) _dow_field: DayOfWeekField field(initFalse, reprFalse) def __post_init__(self): self._minute_field MinuteField(self.minute) self._hour_field HourField(self.hour) self._day_field DayOfMonthField(self.day_of_month) self._month_field MonthField(self.month) self._dow_field DayOfWeekField(self.day_of_week) classmethod def parse(cls, expression: str) - CronExpression: 解析Cron表达式字符串 Args: expression: Cron表达式 (如 0 2 * * 1) Returns: CronExpression实例 parts expression.strip().split() if len(parts) ! 5: raise ValueError(fInvalid cron expression: {expression}) return cls( minuteparts[0], hourparts[1], day_of_monthparts[2], monthparts[3], day_of_weekparts[4] ) def get_next_run(self, from_time: datetime None) - Optional[datetime]: 计算下一次执行时间 Args: from_time: 起始时间默认当前时间 Returns: 下一次执行时间 if from_time is None: from_time datetime.now() # 从下一分钟开始搜索 current from_time.replace(second0, microsecond0) timedelta(minutes1) # 最多搜索2年 end current timedelta(days730) while current end: if self._matches(current): return current current timedelta(minutes1) return None def get_multiple_runs(self, count: int 5, from_time: datetime None) - List[datetime]: 获取多次执行时间 Args: count: 需要的次数 from_time: 起始时间 Returns: 执行时间列表 runs [] current from_time or datetime.now() while len(runs) count: next_run self.get_next_run(current) if next_run is None: break runs.append(next_run) current next_run timedelta(minutes1) return runs def _matches(self, dt: datetime) - bool: 检查时间是否匹配Cron表达式 if not self._minute_field.contains(dt.minute): return False if not self._hour_field.contains(dt.hour): return False if not self._month_field.contains(dt.month): return False # 日和星期是OR关系任一匹配即可 day_match self._day_field.contains(dt.day) dow_match self._dow_field.contains(dt.weekday()) # 如果日和星期都有具体值则两个都要匹配 if self.day_of_month ! * and self.day_of_week ! *: return day_match and dow_match return day_match or dow_match def __str__(self) - str: return f{self.minute} {self.hour} {self.day_of_month} {self.month} {self.day_of_week} class CronScheduler: Cron调度器 管理周期性任务的调度。 def __init__(self): self._jobs: dict {} # job_id - (cron_expr, task_template) def add_job(self, job_id: str, cron_expr: str, task_template: dict): 添加Cron任务 Args: job_id: 任务ID cron_expr: Cron表达式 task_template: 任务模板 cron CronExpression.parse(cron_expr) self._jobs[job_id] (cron, task_template) def remove_job(self, job_id: str): 移除Cron任务 self._jobs.pop(job_id, None) def get_due_tasks(self, current_time: datetime None) - list: 获取到期的任务 Args: current_time: 当前时间 Returns: 到期任务列表 if current_time is None: current_time datetime.now() due_tasks [] for job_id, (cron, template) in self._jobs.items(): next_run cron.get_next_run(current_time - timedelta(minutes1)) if next_run and next_run current_time: task {**template, job_id: job_id, scheduled_time: next_run} due_tasks.append(task) return due_tasks def list_jobs(self) - list: 列出所有Cron任务 return [ { job_id: jid, cron: str(cron), template: tmpl } for jid, (cron, tmpl) in self._jobs.items() ]三、DAG工作流引擎3.1 增强版DAG# advanced/dag_workflow.py DAG工作流引擎 支持复杂的任务编排包括条件分支、并行执行和超时控制。 from __future__ import annotations from typing import Dict, List, Optional, Set, Callable, Any from dataclasses import dataclass, field from enum import Enum import asyncio import time import logging from collections import defaultdict, deque from scheduler.models.task import Task, TaskStatus logger logging.getLogger(__name__) class WorkflowStatus(Enum): 工作流状态 PENDING pending RUNNING running PAUSED paused COMPLETED completed FAILED failed CANCELLED cancelled class NodeType(Enum): 节点类型 TASK task # 普通任务 CONDITION cond # 条件分支 PARALLEL parallel # 并行网关 JOIN join # 汇聚网关 SUB_WORKFLOW sub # 子工作流 dataclass class WorkflowNode: 工作流节点 代表DAG中的一个步骤。 node_id: str name: str node_type: NodeType NodeType.TASK task_template: dict field(default_factorydict) # 条件分支 condition: Optional[Callable] None # 返回True/False # 超时 timeout_seconds: Optional[float] None # 重试 max_retries: int 0 retry_delay: float 1.0 # 状态 status: str pending result: Any None error: Optional[str] None started_at: Optional[float] None completed_at: Optional[float] None dataclass class WorkflowEdge: 工作流边 source_id: str target_id: str label: str # 条件标签用于条件分支 class WorkflowDefinition: 工作流定义 描述一个DAG工作流的拓扑结构。 def __init__(self, workflow_id: str, name: str): self.workflow_id workflow_id self.name name self.nodes: Dict[str, WorkflowNode] {} self.edges: List[WorkflowEdge] [] # 邻接表 self._outgoing: Dict[str, List[WorkflowEdge]] defaultdict(list) self._incoming: Dict[str, List[WorkflowEdge]] defaultdict(list) def add_node(self, node: WorkflowNode): 添加节点 self.nodes[node.node_id] node def add_edge(self, source_id: str, target_id: str, label: str ): 添加边 edge WorkflowEdge(source_id, target_id, label) self.edges.append(edge) self._outgoing[source_id].append(edge) self._incoming[target_id].append(edge) def get_upstream(self, node_id: str) - List[WorkflowNode]: 获取上游节点 return [ self.nodes[e.source_id] for e in self._incoming.get(node_id, []) if e.source_id in self.nodes ] def get_downstream(self, node_id: str) - List[WorkflowNode]: 获取下游节点 return [ self.nodes[e.target_id] for e in self._outgoing.get(node_id, []) if e.target_id in self.nodes ] def get_roots(self) - List[WorkflowNode]: 获取根节点没有上游的节点 upstream_ids {e.source_id for e in self.edges} return [ n for n in self.nodes.values() if n.node_id not in upstream_ids ] def validate(self) - List[str]: 验证工作流定义 Returns: 错误信息列表 errors [] # 检查是否有环 visited set() rec_stack set() def dfs(node_id): visited.add(node_id) rec_stack.add(node_id) for edge in self._outgoing.get(node_id, []): if edge.target_id not in visited: if dfs(edge.target_id): return True elif edge.target_id in rec_stack: return True rec_stack.discard(node_id) return False for node_id in self.nodes: if node_id not in visited: if dfs(node_id): errors.append(fCycle detected involving node: {node_id}) break # 检查孤立节点 connected set() for edge in self.edges: connected.add(edge.source_id) connected.add(edge.target_id) for node_id in self.nodes: if node_id not in connected and len(self.nodes) 1: errors.append(fIsolated node: {node_id}) return errors class WorkflowInstance: 工作流实例 工作流的一次执行。 def __init__(self, definition: WorkflowDefinition, inputs: dict None): self.definition definition self.instance_id fwf-{int(time.time())}-{definition.workflow_id} self.inputs inputs or {} self.status WorkflowStatus.PENDING # 节点状态快照 self.node_states: Dict[str, WorkflowNode] { nid: WorkflowNode( node_idn.node_id, namen.name, node_typen.node_type, task_templaten.task_template.copy() ) for nid, n in definition.nodes.items() } # 上下文数据 self.context: dict {} # 回调 self.on_node_completed: Optional[Callable] None self.on_workflow_completed: Optional[Callable] None def get_ready_nodes(self) - List[WorkflowNode]: 获取可以执行的节点 条件所有上游节点已完成且满足条件。 ready [] for node_id, node in self.node_states.items(): if node.status ! pending: continue upstream self.definition.get_upstream(node_id) # 检查所有上游是否完成 all_completed all( self.node_states[u.node_id].status completed for u in upstream ) if not all_completed: continue # 检查条件分支 conditions_met True for edge in self.definition._incoming.get(node_id, []): if edge.label: source_result self.node_states[edge.source_id].result if source_result ! edge.label: conditions_met False break if conditions_met: ready.append(node) return ready def mark_node_completed(self, node_id: str, result: Any None): 标记节点完成 node self.node_states[node_id] node.status completed node.result result node.completed_at time.time() self.context[f{node_id}.result] result if self.on_node_completed: self.on_node_completed(node) def mark_node_failed(self, node_id: str, error: str): 标记节点失败 node self.node_states[node_id] node.status failed node.error error node.completed_at time.time() def is_completed(self) - bool: 检查工作流是否完成 return all( n.status in (completed, failed, skipped) for n in self.node_states.values() ) class WorkflowEngine: 工作流引擎 管理工作流的注册和执行。 def __init__(self): self._definitions: Dict[str, WorkflowDefinition] {} self._instances: Dict[str, WorkflowInstance] {} self._running False def register_workflow(self, definition: WorkflowDefinition): 注册工作流定义 errors definition.validate() if errors: raise ValueError(fInvalid workflow: {errors}) self._definitions[definition.workflow_id] definition logger.info(fWorkflow registered: {definition.name} ({definition.workflow_id})) async def execute_workflow(self, workflow_id: str, inputs: dict None) - str: 执行工作流 Args: workflow_id: 工作流ID inputs: 输入参数 Returns: 实例ID definition self._definitions.get(workflow_id) if not definition: raise ValueError(fWorkflow not found: {workflow_id}) instance WorkflowInstance(definition, inputs) instance.status WorkflowStatus.RUNNING self._instances[instance.instance_id] instance logger.info(fWorkflow started: {definition.name} f(instance: {instance.instance_id})) # 异步执行 asyncio.create_task(self._execute_instance(instance)) return instance.instance_id async def _execute_instance(self, instance: WorkflowInstance): 执行工作流实例 try: while not instance.is_completed(): ready_nodes instance.get_ready_nodes() if not ready_nodes: if not instance.is_completed(): # 死锁检测 logger.warning(fWorkflow deadlocked: {instance.instance_id}) instance.status WorkflowStatus.FAILED break # 并行执行所有就绪节点 tasks [ self._execute_node(instance, node) for node in ready_nodes ] await asyncio.gather(*tasks) # 检查最终状态 failed any( n.status failed for n in instance.node_states.values() ) instance.status WorkflowStatus.FAILED if failed else WorkflowStatus.COMPLETED logger.info(fWorkflow completed: {instance.instance_id} f({instance.status.value})) if instance.on_workflow_completed: instance.on_workflow_completed(instance) except Exception as e: logger.error(fWorkflow execution failed: {e}) instance.status WorkflowStatus.FAILED async def _execute_node(self, instance: WorkflowInstance, node: WorkflowNode): 执行单个节点 logger.info(fExecuting node: {node.name} ({node.node_id})) node.status running node.started_at time.time() try: # 模拟任务执行 # 实际项目中会调用调度器 await asyncio.sleep(0.1) instance.mark_node_completed(node.node_id, {status: ok}) except Exception as e: instance.mark_node_failed(node.node_id, str(e)) def get_instance(self, instance_id: str) - Optional[WorkflowInstance]: 获取工作流实例 return self._instances.get(instance_id) def get_instance_status(self, instance_id: str) - dict: 获取实例状态 instance self._instances.get(instance_id) if not instance: return {error: Instance not found} return { instance_id: instance.instance_id, workflow_id: instance.definition.workflow_id, status: instance.status.value, nodes: { nid: { name: n.name, status: n.status, result: n.result, error: n.error } for nid, n in instance.node_states.items() } }四、动态分片4.1 分片引擎# advanced/sharding.py 动态分片引擎 将大数据集拆分为多个小分片并行处理。 from __future__ import annotations from typing import Dict, List, Optional, Any, Callable from dataclasses import dataclass, field from enum import Enum import asyncio import math import time import logging import hashlib from scheduler.models.task import Task, TaskStatus logger logging.getLogger(__name__) class ShardingStrategy(Enum): 分片策略 RANGE range # 范围分片 HASH hash # 哈希分片 LIST list # 列表分片 DYNAMIC dynamic # 动态分片自适应 dataclass class Shard: 分片 数据的一个子集。 shard_id: str index: int total_shards: int data_range: tuple # (start, end) 或 (key_pattern,) task_template: dict field(default_factorydict) def to_task(self) - Task: 转换为任务 return Task( namefshard-{self.shard_id}, task_typeshard, handlerprocess.shard, params{ shard_id: self.shard_id, index: self.index, total: self.total_shards, data_range: self.data_range, **self.task_template } ) class ShardGenerator: 分片生成器 根据策略和数据规模生成分片。 staticmethod def range_shard(total_items: int, shard_size: int, task_template: dict None) - List[Shard]: 范围分片 Args: total_items: 总数据量 shard_size: 每个分片的大小 task_template: 任务模板 Returns: 分片列表 num_shards math.ceil(total_items / shard_size) shards [] for i in range(num_shards): start i * shard_size end min((i 1) * shard_size, total_items) shard Shard( shard_idfrange-{i:04d}, indexi, total_shardsnum_shards, data_range(start, end), task_templatetask_template or {} ) shards.append(shard) return shards staticmethod def hash_shard(keys: List[str], num_shards: int, task_template: dict None) - Dict[str, Shard]: 哈希分片 Args: keys: 数据键列表 num_shards: 分片数量 task_template: 任务模板 Returns: 键到分片的映射 shards {} key_to_shard {} for i, key in enumerate(keys): shard_index int(hashlib.md5(key.encode()).hexdigest(), 16) % num_shards shard_id fhash-{shard_index:04d} if shard_id not in shards: shards[shard_id] Shard( shard_idshard_id, indexshard_index, total_shardsnum_shards, data_range(shard_index, shard_index), task_templatetask_template or {} ) key_to_shard[key] shard_id return key_to_shard staticmethod def dynamic_shard(total_items: int, min_shard_size: int, max_shard_size: int, task_template: dict None) - List[Shard]: 动态分片自适应 根据数据量自动调整分片大小。 Args: total_items: 总数据量 min_shard_size: 最小分片大小 max_shard_size: 最大分片大小 task_template: 任务模板 Returns: 分片列表 # 计算理想分片数 ideal_num max(1, total_items // ((min_shard_size max_shard_size) // 2)) # 调整分片大小 shard_size max(min_shard_size, min(max_shard_size, total_items // ideal_num)) num_shards math.ceil(total_items / shard_size) return ShardGenerator.range_shard(total_items, shard_size, task_template) class ShardExecutor: 分片执行器 管理分片的创建、分发和结果汇总。 def __init__(self): self._shards: Dict[str, List[Shard]] {} self._results: Dict[str, Dict[str, Any]] {} async def execute_sharded(self, name: str, shards: List[Shard], process_func: Callable, concurrency: int 5) - List[Any]: 执行分片任务 Args: name: 任务名称 shards: 分片列表 process_func: 处理函数 (shard) - result concurrency: 并发数 Returns: 结果列表 self._shards[name] shards semaphore asyncio.Semaphore(concurrency) async def process_with_limit(shard): async with semaphore: logger.info(fProcessing shard: {shard.shard_id} f({shard.index 1}/{shard.total_shards})) try: result await process_func(shard) self._results.setdefault(name, {})[shard.shard_id] result return result except Exception as e: logger.error(fShard failed: {shard.shard_id}: {e}) raise tasks [process_with_limit(s) for s in shards] results await asyncio.gather(*tasks, return_exceptionsTrue) # 检查是否有失败 errors [r for r in results if isinstance(r, Exception)] if errors: logger.warning(f{len(errors)}/{len(shards)} shards failed) return [r for r in results if not isinstance(r, Exception)] def aggregate_results(self, name: str, merge_func: Callable None) - Any: 汇总分片结果 Args: name: 任务名称 merge_func: 合并函数 Returns: 汇总结果 results self._results.get(name, {}) if not results: return None if merge_func: return merge_func(results) return results def get_progress(self, name: str) - dict: 获取执行进度 shards self._shards.get(name, []) results self._results.get(name, {}) return { total: len(shards), completed: len(results), progress: len(results) / max(len(shards), 1) * 100 }五、高级特性演示5.1 综合演示# examples/advanced_features_demo.py 高级特性演示 import asyncio import logging import sys from datetime import datetime, timedelta sys.path.insert(0, ..) from advanced.cron import CronExpression, CronScheduler from advanced.dag_workflow import ( WorkflowDefinition, WorkflowNode, WorkflowEngine, NodeType ) from advanced.sharding import ShardGenerator, ShardExecutor, Shard logging.basicConfig(levellogging.INFO) async def demo_cron(): 演示Cron表达式 print( * 79) print(⏰ Cron表达式演示) print( * 79) # 解析各种Cron表达式 expressions [ 0 2 * * *, # 每天凌晨2点 */15 * * * *, # 每15分钟 0 9-17 * * 1-5, # 工作日每小时 0 0 1 * *, # 每月1号 30 4 * * 0, # 每周日凌晨4:30 ] print(\nCron表达式解析:) for expr in expressions: cron CronExpression.parse(expr) next_runs cron.get_multiple_runs(count3) print(f\n {expr}) for i, run in enumerate(next_runs): print(f 第{i1}次: {run.strftime(%Y-%m-%d %H:%M)}) # Cron调度器 print(\n\nCron调度器演示:) scheduler CronScheduler() scheduler.add_job( daily-backup, 0 2 * * *, {name: 数据备份, type: backup, target: /data} ) scheduler.add_job( hourly-report, 0 * * * *, {name: 小时报表, type: report, period: hour} ) print(\n已注册的Cron任务:) for job in scheduler.list_jobs(): print(f {job[job_id]}: {job[cron]} → {job[template][name]}) async def demo_dag_workflow(): 演示DAG工作流 print(\n * 79) print( DAG工作流演示) print( * 79) # 创建工作流定义 workflow WorkflowDefinition(data-pipeline, 数据处理流水线) # 创建节点 extract WorkflowNode(extract, 数据抽取) clean WorkflowNode(clean, 数据清洗) transform WorkflowNode(transform, 数据转换) validate WorkflowNode(validate, 数据校验) notify_success WorkflowNode(notify-success, 发送成功通知) notify_fail WorkflowNode(notify-fail, 发送失败通知) # 添加节点 for node in [extract, clean, transform, validate, notify_success, notify_fail]: workflow.add_node(node) # 添加边 workflow.add_edge(extract, clean) workflow.add_edge(clean, transform) workflow.add_edge(transform, validate) workflow.add_edge(validate, notify-success, labelsuccess) workflow.add_edge(validate, notify-fail, labelfailure) # 验证 errors workflow.validate() if errors: print(f验证错误: {errors}) return print(f\n工作流: {workflow.name}) print(f节点数: {len(workflow.nodes)}) print(f边数: {len(workflow.edges)}) # 执行工作流 engine WorkflowEngine() engine.register_workflow(workflow) instance_id await engine.execute_workflow(data-pipeline) await asyncio.sleep(1) status engine.get_instance_status(instance_id) print(f\n执行状态: {status[status]}) print(节点执行情况:) for nid, info in status[nodes].items(): icon ✅ if info[status] completed else ⏳ print(f {icon} {info[name]}: {info[status]}) async def demo_sharding(): 演示动态分片 print(\n * 79) print( 动态分片演示) print( * 79) # 模拟100万条数据 total_items 1_000_000 shard_size 100_000 print(f\n总数据量: {total_items:,}) print(f分片大小: {shard_size:,}) # 范围分片 print(\n1. 范围分片:) shards ShardGenerator.range_shard(total_items, shard_size) print(f 生成 {len(shards)} 个分片) for s in shards[:3]: print(f {s.shard_id}: {s.data_range}) print(f ... 共 {len(shards)} 个分片) # 哈希分片 print(\n2. 哈希分片:) keys [fuser-{i} for i in range(100)] key_map ShardGenerator.hash_shard(keys, 8) # 统计分布 distribution {} for key, shard_id in key_map.items(): distribution[shard_id] distribution.get(shard_id, 0) 1 print( 键分布情况:) for shard_id, count in sorted(distribution.items()): bar █ * (count // 2) print(f {shard_id}: {bar} ({count})) # 动态分片 print(\n3. 动态分片自适应:) dynamic_shards ShardGenerator.dynamic_shard( total_items, min_shard_size50_000, max_shard_size200_000 ) print(f 自适应生成 {len(dynamic_shards)} 个分片) # 模拟执行 print(\n4. 模拟分片执行:) executor ShardExecutor() async def process_shard(shard): await asyncio.sleep(0.1) return {shard: shard.shard_id, processed: True} results await executor.execute_sharded( demo, shards[:5], process_shard, concurrency3 ) progress executor.get_progress(demo) print(f 进度: {progress[completed]}/{progress[total]} f({progress[progress]:.0f}%)) async def main(): await demo_cron() await demo_dag_workflow() await demo_sharding() if __name__ __main__: asyncio.run(main())六、测试6.1 高级特性测试# tests/test_advanced.py import pytest from datetime import datetime, timedelta from advanced.cron import CronExpression, CronScheduler from advanced.dag_workflow import WorkflowDefinition, WorkflowNode, WorkflowEngine from advanced.sharding import ShardGenerator, ShardExecutor class TestCron: Cron测试 def test_parse_standard(self): cron CronExpression.parse(0 2 * * *) assert cron.minute 0 assert cron.hour 2 def test_next_run(self): cron CronExpression.parse(0 * * * *) # 每小时 now datetime(2024, 1, 1, 10, 30) next_run cron.get_next_run(now) assert next_run.hour 11 assert next_run.minute 0 def test_multiple_runs(self): cron CronExpression.parse(*/30 * * * *) # 每30分钟 runs cron.get_multiple_runs(count4) assert len(runs) 4 class TestDAG: DAG测试 def test_workflow_validation(self): wf WorkflowDefinition(test, Test) n1 WorkflowNode(n1, Node 1) n2 WorkflowNode(n2, Node 2) wf.add_node(n1) wf.add_node(n2) wf.add_edge(n1, n2) errors wf.validate() assert len(errors) 0 def test_cycle_detection(self): wf WorkflowDefinition(test, Test) n1 WorkflowNode(n1, Node 1) n2 WorkflowNode(n2, Node 2) wf.add_node(n1) wf.add_node(n2) wf.add_edge(n1, n2) wf.add_edge(n2, n1) # 制造环 errors wf.validate() assert len(errors) 0 class TestSharding: 分片测试 def test_range_shard(self): shards ShardGenerator.range_shard(1000, 100) assert len(shards) 10 assert shards[0].data_range (0, 100) def test_dynamic_shard(self): shards ShardGenerator.dynamic_shard(1000, 50, 200) assert len(shards) 0 assert all(s.total_shards len(shards) for s in shards) if __name__ __main__: pytest.main([__file__, -v])七、总结7.1 本讲成果组件文件功能CronExpression​advanced/cron.pyCron表达式解析CronScheduler​advanced/cron.py周期性调度WorkflowEngine​advanced/dag_workflow.pyDAG工作流引擎ShardGenerator​advanced/sharding.py分片生成器ShardExecutor​advanced/sharding.py分片执行器7.2 核心知识点Cron表达式标准的5字段格式支持步进、范围和列表DAG工作流有向无环图支持条件分支和并行执行动态分片范围分片、哈希分片、自适应分片结果聚合分片结果的收集和合并7.3 下一讲预告第9讲性能优化我们将对系统进行性能调优异步IO优化批量处理资源隔离连接池优化缓存策略准备好了吗让我们在第9讲再见开发之余的小工具推荐​处理 Base64、JWT 解析、JSON 格式化、Crontab 计算、PDF 合并压缩这些碎片需求我常用一个纯前端本地工具箱zz365.top。所有计算在浏览器完成文件不上传服务器关页即清。免费、无登录、无广告适合开发者当常驻标签页。