智能体工作流引擎:任务拆解与并发调度架构
处理复杂智能体目标的挑战
当面对诸如"生成行业深度分析报告"这类宏观指令时,单一的智能体模型往往难以直接输出高质量结果。此类目标通常隐含了多个子步骤,例如信息检索、数据清洗、竞品对比以及最终的内容合成。为了有效管理这一过程,我们需要构建一个能够理解意图、拆解步骤并协调执行的工作流引擎。
核心模块实现
1. 智能工作流规划
利用大语言模型的推理能力,将宏观目标转化为可执行的步骤列表。以下是一个规划器类的实现,它负责生成带有依赖关系的活动图:
from typing import Any, Mapping, List
import asyncio
class WorkflowArchitect:
def __init__(self, model_client):
self.model = model_client
async def plan_workflow(self, objective: str) -> Mapping[str, Any]:
"""将宏观目标转化为结构化工作流"""
system_prompt = f"""
目标:{objective}
请规划执行步骤,返回 JSON 格式:
{{
"steps": [
{{
"step_id": "step_01",
"action": "动作名称",
"payload": "具体参数",
"prerequisites": ["依赖的 step_id"],
"timeout": 300
}}
]
}}
注意:
1. 确保依赖关系无环
2. 步骤应具备原子性
3. 标记可并行执行的节点
"""
raw_output = await self.model.chat(system_prompt)
return self._sanitize_plan(raw_output)
def _sanitize_plan(self, plan_data: dict) -> dict:
"""校验计划合法性"""
self._validate_dependency_graph(plan_data["steps"])
return self._topological_sort(plan_data["steps"])
2. 并发控制与执行
为了最大化资源利用率,执行器需要支持异步并发,同时限制最大并发数以防止系统过载。使用信号量来控制活跃任务的数量:
class ConcurrencyController:
def __init__(self, concurrency_limit: int = 10):
self.limit = concurrency_limit
self.semaphore = asyncio.Semaphore(concurrency_limit)
self.state_store = {}
async def run_graph(self, workflow_plan: Mapping):
"""基于依赖图执行任务"""
pending_tasks = {step["step_id"]: step for step in workflow_plan["steps"]}
completed_ids = set()
while len(completed_ids) < len(pending_tasks):
# 找出当前可执行的任务
executable = []
for sid, step in pending_tasks.items():
if sid not in completed_ids:
deps = set(step.get("prerequisites", []))
if deps.issubset(completed_ids):
executable.append(self._run_step(sid, step))
if not executable:
raise RuntimeError("检测到死锁或依赖缺失")
# 并发执行当前批次
results = await asyncio.gather(*executable)
for sid, res in results:
self.state_store[sid] = res
completed_ids.add(sid)
async def _run_step(self, step_id: str, step_config: dict):
async with self.semaphore:
# 模拟实际业务逻辑执行
result = await self._invoke_action(step_config)
return step_id, result
3. 状态持久化层
在长运行任务中,中间状态需要持久化以防止进程崩溃导致数据丢失。使用 Redis 作为高速缓存层存储任务状态:
import redis.asyncio as redis
import msgpack
class StateStore:
def __init__(self, redis_url: str):
self.client = redis.from_url(redis_url)
async def persist_state(self, workflow_id: str, step_id: str, data: Any):
"""异步保存状态"""
key = f"wf:{workflow_id}:step:{step_id}"
packed_data = msgpack.packb(data)
await self.client.setex(key, 3600, packed_data)
async def restore_state(self, workflow_id: str, step_id: str) -> Any:
"""恢复状态"""
key = f"wf:{workflow_id}:step:{step_id}"
raw = await self.client.get(key)
if not raw:
return None
return msgpack.unpackb(raw)
4. 调度策略模式
系统应支持多种调度模式以适应不同场景,例如串行管道、全并行或基于依赖图的混合模式:
class FlowEngine:
def __init__(self):
self.planner = WorkflowArchitect()
self.runner = ConcurrencyController()
self.store = StateStore()
async def execute_serial(self, steps: List[dict]):
"""串行执行模式"""
context = {}
for step in steps:
output = await self.runner._run_step(step["id"], step)
context[step["id"]] = output
return context
async def execute_fan_out(self, steps: List[dict]):
"""扇出并行模式"""
tasks = [self.runner._run_step(s["id"], s) for s in steps]
return await asyncio.gather(*tasks)
async def execute_dynamic(self, goal: str):
"""动态规划执行模式"""
plan = await self.planner.plan_workflow(goal)
return await self.runner.run_graph(plan)
5. 资源与性能调优
通过缓存命中检查和资源预估来提升整体吞吐量:
from cachetools import TTLCache
class ResourceTuner:
def __init__(self):
self.cache = TTLCache(maxsize=500, ttl=600)
async def tune_step(self, step_config: dict) -> dict:
"""优化步骤配置"""
step_hash = hash(step_config["action"] + str(step_config["payload"]))
# 缓存命中直接返回
if step_hash in self.cache:
return self.cache[step_hash]
# 动态调整资源配额
config = step_config.copy()
config["quota"] = self._calculate_quota(config)
self.cache[step_hash] = config
return config
def _calculate_quota(self, config: dict) -> dict:
return {
"memory_mb": 512 if config["type"] == "compute" else 256,
"timeout_sec": 60
}
工程化落地案例:自动化分析管道
以下是一个完整的分析管道实现,展示了如何组合上述模块来完成具体的业务需求:
class AnalysisPipeline:
def __init__(self):
self.engine = FlowEngine()
self.tuner = ResourceTuner()
async def process_topic(self, topic: str):
# 1. 动态生成工作流
workflow = await self.engine.planner.plan_workflow(
f"针对 {topic} 进行多维数据分析"
)
# 2. 预优化每个步骤
optimized_steps = []
for step in workflow["steps"]:
opt_step = await self.tuner.tune_step(step)
optimized_steps.append(opt_step)
workflow["steps"] = optimized_steps
# 3. 执行并收集结果
final_state = await self.engine.runner.run_graph(workflow)
# 4. 聚合输出
return self._aggregate_insights(final_state)
def _aggregate_insights(self, state: dict):
# 整合各步骤产出物
report_parts = []
for key, value in state.items():
report_parts.append(f"## {key}\n{value}")
return "\n".join(report_parts)
稳定性与可靠性指南
在分布式或高负载环境下,需重点关注以下几个工程问题:
- 依赖循环检测:在任务提交前必须进行拓扑排序检查,防止因循环依赖导致执行挂起。
- 资源隔离:不同优先级的任务应进入不同的队列,避免低优先级任务阻塞关键路径。
- 状态幂等性:任务重试机制必须保证幂等,避免重复执行导致数据污染,可通过唯一请求 ID 实现。
- 异常熔断:当连续任务失败率达到阈值时,应触发熔断机制,暂停新任务提交并报警。
- 可观测性:每个步骤的执行耗时、资源消耗及状态变更都应接入监控系统的 Trace 链路。
