项目04-手搓Agent之DAG任务编排(按顺序执行任务)
·
介绍
基于 DAG(有向无环图)的团队任务编排主要实现了以下三个内容:
- 多 Agent(Coder / Researcher/ Testor)协作执行任务 :如何具体地完成任务;
- 生产者 - 消费者模型 :任务队列 + 多线程异步执行;
- DAG 任务依赖编排 :保证任务按先后顺序执行(比如必须先分析需求→再设计架构→最后编码);
1 示例:顶层接口
def demo_dag_orchestration():
"""演示 3: DAG 任务编排
功能:创建带依赖的任务,查看执行顺序和可执行任务
"""
print("\n=== 演示 3: DAG 任务编排 ===")
# 初始化OpenAI客户端(需提前定义API_KEY和BASE_URL)
client = OpenAI(api_key=API_KEY, base_url=BASE_URL)
# 创建AI智能体团队
team = AgentTeam(client, MODEL)
# 创建DAG任务编排器(绑定智能体团队)
orchestrator = DAGTaskOrchestrator(team)
# 创建3个有依赖关系的任务
# 任务1:分析需求(无依赖)
t1 = orchestrator.add_task_with_deps({"content": "分析需求"})
# 任务2:设计架构(依赖任务1)
t2 = orchestrator.add_task_with_deps({"content": "设计架构"}, [t1])
# 任务3:编写代码(依赖任务2)
t3 = orchestrator.add_task_with_deps({"content": "编写代码"}, [t2])
# 打印演示信息
print(f"✅ 创建了 3 个任务")
print(f"📋 执行顺序: {orchestrator.get_execution_order()}")
print(f"🎯 可执行任务: {orchestrator.get_ready_tasks()}")
- 首先先创建
Agent团队; - 基于团队构建任务编排器
DAG; - 创建任务依赖关系;
2 创建 Agent 团队
虽然前一章节,通过 TeamCollaboration 创建了 Agent 的邮箱,但是邮箱是用来通信的。具体 Agent 的工作信息没有存储,因此需要 AgentTeam 类。
class AgentTeam:
"""🤖 Agent 团队 - 生产者消费者 + DAG
核心:管理多个AI智能体,提供任务队列,启动工作线程
"""
def __init__(self, api_client, model: str = "gpt-5-nano"):
self.api_client = api_client # OpenAI客户端
self.model = model # 使用的大模型名称
# 初始化智能体集群:1个主管 + 3个子智能体(编码/研究/测试)
self.lead = LeadAgent("lead", model, api_client)
self.sub_agents = {
"coder": SubAgent("coder", "coding", model, api_client),
"researcher": SubAgent("researcher", "research", model, api_client),
"tester": SubAgent("tester", "testing", model, api_client),
}
# 生产者-消费者模型核心:任务队列(存储待执行任务)
self.task_queue = queue.Queue()
# DAG调度器:管理任务依赖关系
self.dag_scheduler = DAGScheduler()
# 存储运行中的智能体线程
self.running_agents = {}
# 团队协作协议(版本、角色、通信方式)
self.team_protocol = {
"version": "1.0",
"roles": list(self.sub_agents.keys()),
"communication": "mailbox"
}
def submit_task(self, task: Dict) -> str:
"""提交任务到团队队列(生产者)
:param task: 任务字典
:return: 任务ID
"""
# 生成唯一任务ID(时间戳保证不重复)
task_id = task.get("id", f"task_{datetime.now().timestamp()}")
# 将任务放入队列
self.task_queue.put({"id": task_id, **task})
return task_id
def start_team(self, num_workers: int = 3):
"""启动智能体团队,创建工作线程(消费者)
:param num_workers: 工作线程数量
"""
for i in range(num_workers):
# 轮询分配子智能体(coder→researcher→tester循环)
agent_name = list(self.sub_agents.keys())[i % len(self.sub_agents)]
agent = self.sub_agents[agent_name]
# 创建守护线程:主线程退出,子线程自动退出
t = threading.Thread(
target=self._worker_loop,
args=(agent,),
daemon=True
)
t.start()
# 记录运行中的线程
self.running_agents[agent_name] = t
def _worker_loop(self, agent: SubAgent):
"""智能体工作循环(消费者核心逻辑)
持续从队列取任务 → 执行任务 → 标记完成
"""
while True:
try:
# 从队列取任务,超时2秒(避免无限阻塞)
task = self.task_queue.get(timeout=2)
# 子智能体执行任务
result = agent.process_task(task)
# 标记任务完成(队列任务计数-1)
self.task_queue.task_done()
except queue.Empty:
# 队列为空,继续循环等待新任务
continue
except Exception as e:
# 捕获异常,避免线程崩溃
print(f"❌ [{agent.name}] Error: {e}")
def get_team_status(self) -> Dict:
"""获取团队运行状态"""
return {
"protocol": self.team_protocol,
"queue_size": self.task_queue.qsize(), # 待执行任务数量
"agents": list(self.sub_agents.keys()) # 智能体列表
}
- 第一步:创建 lead 和 sub_agent. 实例化这几个 Agent;
- 第二步:定义任务队列. 存储待执行的任务;
- 第三步:DAG 调度器实例化. 管理任务的依赖关系;
- 第四步:实例化团队协议. 用json形式进行存储,协议为
mailbox;
2.1 提交任务-生产者
def submit_task(self, task: Dict) -> str:
"""提交任务到团队队列(生产者)
:param task: 任务字典
:return: 任务ID
"""
# 生成唯一任务ID(时间戳保证不重复)
task_id = task.get("id", f"task_{datetime.now().timestamp()}")
# 将任务放入队列
self.task_queue.put({"id": task_id, **task})
return task_id
2.2 Agent 工作逻辑-消费者
消费逻辑为:(1)通过轮询的方式给每一个 Agent 赋予 agent_name;(2)通过线程的方式启动 sub_agent (本质就是守护线程);(3)具体处理任务的逻辑:从队列取任务,进行执行得到 result 后,任务状态标记为完成;
def start_team(self, num_workers: int = 3):
"""启动智能体团队,创建工作线程(消费者)
:param num_workers: 工作线程数量
"""
for i in range(num_workers):
# 轮询分配子智能体(coder→researcher→tester循环)
agent_name = list(self.sub_agents.keys())[i % len(self.sub_agents)]
agent = self.sub_agents[agent_name]
# 创建守护线程:主线程退出,子线程自动退出
t = threading.Thread(
target=self._worker_loop,
args=(agent,),
daemon=True
)
t.start()
# 记录运行中的线程
self.running_agents[agent_name] = t
def _worker_loop(self, agent: SubAgent):
"""智能体工作循环(消费者核心逻辑)
持续从队列取任务 → 执行任务 → 标记完成
"""
while True:
try:
# 从队列取任务,超时2秒(避免无限阻塞)
task = self.task_queue.get(timeout=2)
# 子智能体执行任务
result = agent.process_task(task)
# 标记任务完成(队列任务计数-1)
self.task_queue.task_done()
except queue.Empty:
# 队列为空,继续循环等待新任务
continue
except Exception as e:
# 捕获异常,避免线程崩溃
print(f"❌ [{agent.name}] Error: {e}")
3 DAG 构建
class DAGTaskOrchestrator:
"""🎼 DAG 任务编排器
核心:基于DAG管理任务依赖,保证任务按顺序执行
例如:必须先分析需求→再设计架构→最后编写代码
"""
def __init__(self, team: AgentTeam):
self.team = team # 绑定AI智能体团队
self.dag = DAGScheduler() # DAG调度器实例
self.task_results = {} # 存储任务执行结果
def add_task_with_deps(self, task: Dict, dependencies: List[str] = None):
"""添加带依赖关系的任务
:param task: 任务内容
:param dependencies: 依赖的任务ID列表
:return: 任务ID
"""
# 生成唯一任务ID
task_id = task.get("id", f"task_{datetime.now().timestamp()}")
# 将任务和依赖加入DAG调度器
self.dag.add_task(task_id, dependencies)
# 提交任务到团队队列,等待执行
self.team.submit_task({"id": task_id, **task})
return task_id
def get_ready_tasks(self) -> List[str]:
"""获取【所有依赖已完成】的可执行任务"""
return self.dag.get_ready_tasks()
def mark_task_done(self, task_id: str, result: str = None):
"""标记任务完成,并存储执行结果
:param task_id: 任务ID
:param result: 任务执行结果
"""
self.dag.mark_complete(task_id)
if result:
self.task_results[task_id] = result
def get_execution_order(self) -> List[str]:
"""获取DAG拓扑排序结果(任务执行顺序)"""
try:
return self.dag.topological_sort()
except ValueError as e:
print(f"❌ DAG Error: {e}")
return []
4 总结
上一部分的团队协作(TeamCollaboration) = 发消息 + 分配任务 + 记录日志;
它管的是: 谁给谁发了什么、任务分给谁、日志存哪里→ 只管任务的 “交付”,不管执行顺序;
DAG 任务编排(DAGTaskOrchestrator) = 控制执行顺序 + 管理依赖;
它管的是: 任务必须先做 A、再做 B、最后做 C→ 只管任务的 “先后规则”,不管谁来执行;
例子:
合并逻辑:
主管发消息(TeamCollaboration)
↓
创建带依赖的协作任务(DAG + TodoManager)
↓
DAG 判断:哪个任务能执行
↓
把可执行任务分配给智能体(AgentTeam)
↓
智能体完成后 → 标记DAG完成 → 解锁下一个任务
更多推荐


所有评论(0)