介绍

基于 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完成 → 解锁下一个任务
Logo

这里是“一人公司”的成长家园。我们提供从产品曝光、技术变现到法律财税的全栈内容,并连接云服务、办公空间等稀缺资源,助你专注创造,无忧运营。

更多推荐