共计 2209 个字符,预计需要花费 6 分钟才能阅读完成。
背景痛点
在传统的单体应用中,任务调度通常采用简单的定时任务或队列机制。但随着业务复杂度提升,这种架构暴露出几个严重问题:

- 扩展性瓶颈:单机任务队列难以应对突发流量,垂直扩展成本高
- 状态管理困难:任务依赖关系需要手动维护,出错时难以恢复
- 可观测性差:缺乏统一视图查看任务执行链路和资源占用情况
Agent+Workflow 架构通过将任务执行单元 (Agent) 与编排逻辑 (Workflow) 分离,实现了:
- 水平扩展能力:Agent 可动态扩缩容
- 显式状态管理:Workflow 引擎维护 DAG 执行状态
- 可视化监控:完整记录任务生命周期
技术对比
主流方案横向评测
| 方案 | 最大 QPS | 容错机制 | 学习成本 | 适用场景 |
|---|---|---|---|---|
| Celery | 5k | 重试 + 死信队列 | 低 | 简单异步任务 |
| Airflow | 1k | 任务状态持久化 | 中 | 批处理 ETL |
| 自建 Workflow | 10k+ | 自定义 checkpoint | 高 | 复杂业务流程编排 |
选型建议
- 短期需求:直接使用 Celery+Flower 组合
- 长期演进:基于 Redis Stream 自建 Workflow 引擎
- 特殊场景:金融级事务要求建议采用 Cadence/Temporal
核心实现
有状态 Agent 设计
from dataclasses import dataclass
from enum import Enum, auto
class TaskStatus(Enum):
PENDING = auto()
RUNNING = auto()
COMPLETED = auto()
@dataclass
class Task:
task_id: str
params: dict
status: TaskStatus = TaskStatus.PENDING
result: Optional[Any] = None
class ProcessingAgent:
def __init__(self, redis_conn):
self.redis = redis_conn
self.current_task = None
def consume_task(self, queue_name):
"""幂等消费任务,确保重复消费不会导致状态不一致"""
task_data = self.redis.xread({queue_name: '>'}, block=1000)
if not task_data:
return None
task = Task(**json.loads(task_data[0][1]))
if task.status != TaskStatus.PENDING:
return None
task.status = TaskStatus.RUNNING
self.current_task = task
return task
Workflow 调度算法
- DAG 构建:使用邻接表存储任务依赖关系
dag = {'task_a': ['task_b', 'task_c'],
'task_b': ['task_d'],
'task_c': ['task_d'],
'task_d': []}
- 拓扑排序:Kahn 算法实现
def topological_sort(dag):
in_degree = {node: 0 for node in dag}
for node in dag:
for neighbor in dag[node]:
in_degree[neighbor] += 1
queue = deque([node for node in dag if in_degree[node] == 0])
sorted_nodes = []
while queue:
node = queue.popleft()
sorted_nodes.append(node)
for neighbor in dag[node]:
in_degree[neighbor] -= 1
if in_degree[neighbor] == 0:
queue.append(neighbor)
if len(sorted_nodes) != len(dag):
raise ValueError("存在循环依赖")
return sorted_nodes
避坑指南
典型故障场景
- 僵尸任务:
- 现象:任务状态长期处于 RUNNING 但实际已停滞
-
解决方案:实现心跳机制 + 超时回收
-
循环依赖:
- 现象:拓扑排序抛出异常
-
解决方案:设计时使用可视化工具验证 DAG
-
雪崩效应:
- 现象:某个任务失败导致级联故障
- 解决方案:实现熔断机制和降级策略
性能优化
压测数据(Locust)
| 连接池大小 | 平均响应时间 | 吞吐量(req/s) | 错误率 |
|---|---|---|---|
| 10 | 120ms | 850 | 0.1% |
| 50 | 85ms | 4200 | 0.3% |
| 100 | 72ms | 6800 | 1.2% |
Redis 优化建议
- 使用
connection_pool替代直接连接 - 设置合理的
max_connections(建议物理核心数的 2 - 3 倍) - 启用
health_check_interval检测断连
代码规范
必须遵守的 PEP8 要点
- 类名使用
PascalCase - 函数名使用
snake_case - 类型注解必须完整
- 文档字符串使用 Google 风格
示例:
def retry_task(task_id: str, max_retries: int = 3) -> bool:
""" 实现任务幂等重试逻辑
Args:
task_id: 需要重试的任务 ID
max_retries: 最大重试次数
Returns:
bool: 是否重试成功
"""
# ... 实现代码
互动讨论
开放性问题:如何设计跨数据中心的 Workflow 容灾方案?建议从以下几个维度思考:
- 数据同步策略(同步 / 异步)
- 故障检测机制
- 切换时的事务一致性保证
- 回滚方案设计
欢迎在评论区分享你的架构设计方案!
正文完
发表至: 技术分享
近一天内
