云计算百科
云计算领域专业知识百科平台

智能体面试准备(七十七):智能体分布式执行与任务编排工程——长任务分片、依赖调度与状态一致性

智能体面试准备(七十七):智能体分布式执行与任务编排工程——长任务分片、依赖调度与状态一致性

引言

本文是「智能体面试准备」系列第 77 篇,承接(七十六)长程记忆与经验自进化工程。前面讲了记忆怎么存、经验怎么沉淀,这一篇谈一个落地时极容易被低估的问题:当智能体要处理的任务从"单轮问答"变成"跑三天的数据 pipeline + 反复调用十几个工具",单进程串行执行根本扛不住,必须引入分布式执行与编排。

面试官判断一个候选人是否真正做过生产级 Agent,常问:"你的 Agent 怎么保证一个跑了 2 小时的任务在 worker 挂掉后不重头再来?多个子任务有依赖怎么调度?状态在哪、怎么恢复?"能把这三问讲清楚,基本就过关了。本篇给出架构图、对比表、代码片段与面试速答。

智能体分布式执行总览
==============================================
任务编排器 (Orchestrator)
│ 解析 DAG / 依赖
┌───────────┼───────────────┐
▼ ▼ ▼
子任务 A 子任务 B 子任务 C
(无依赖) (依赖A产出) (依赖A,B)
│ │ │
▼ ▼ ▼
┌─────────────────────────────────────┐
│ Worker 池(多进程/多机/容器) │
│ 每个 worker 拉取可执行节点 │
└─────────────────────────────────────┘


┌─────────────────────────────────────┐
│ 状态存储 (DAG状态 + 检查点 + 锁) │
│ Redis / 数据库 / 对象存储 │
└─────────────────────────────────────┘

一、为什么单进程不够

维度单进程串行分布式编排
长任务容错 挂了全重来 节点级重试,仅重跑失败子任务
吞吐 1 任务/时刻 多子任务并行
工具限流 难隔离 按工具维度限流队列
资源 单卡/单机 可跨机调度重工具
可观测 日志散落 统一 DAG 视图

当任务天然有依赖(A 出报告→B 审报告→C 发邮件),串行虽然简单但慢;当子任务可并行(同时爬 10 个站点),串行就是浪费。编排器的价值就是用 DAG 描述依赖、并行可并行、串行保依赖。

二、任务建模:DAG 与节点

把任务抽象成有向无环图,节点是"一个可独立执行的动作"(调一次工具/跑一段代码/生成一段文本),边是依赖。执行引擎只调度"入度已满足且未执行"的节点。

from dataclasses import dataclass, field
from typing import List

@dataclass
class Node:
id: str
deps: List[str] = field(default_factory=list)
status: str = "pending" # pending/running/done/failed
payload: dict = field(default_factory=dict)

@dataclass
class DAG:
nodes: dict = field(default_factory=dict)

def ready(self, nid):
"""入度(依赖)全部 done 才可执行"""
return all(self.nodes[d].status == "done" for d in self.nodes[nid].deps)

def descendants_runnable(self):
return [nid for nid, n in self.nodes.items()
if n.status == "pending" and self.ready(nid)]

三、调度:从串行到并行

最朴素的是拓扑排序后串行;生产级要并发拉起可并行节点,并控制并发度与每工具限流。

import threading, time

class Scheduler:
def __init__(self, dag, max_parallel=4):
self.dag = dag
self.max_parallel = max_parallel
self.lock = threading.Lock()

def _run_node(self, nid, exec_fn):
with self.lock:
self.dag.nodes[nid].status = "running"
try:
exec_fn(self.dag.nodes[nid].payload)
with self.lock:
self.dag.nodes[nid].status = "done"
except Exception:
with self.lock:
self.dag.nodes[nid].status = "failed"
raise

def run(self, exec_fn):
while True:
with self.lock:
pending = self.dag.descendants_runnable()
running = [n for n in self.dag.nodes.values() if n.status=="running"]
if not pending and not running:
break
for nid in pending[: self.max_parallel len(running)]:
threading.Thread(target=self._run_node, args=(nid, exec_fn)).start()
time.sleep(0.2)

四、状态一致性与检查点

分布式执行最大的坑是"状态在哪、怎么恢复"。原则:执行逻辑无状态,状态外置。每个节点执行前写"开始"标记,成功后写"完成"+产物引用;worker 崩溃重启后,编排器扫描 DAG,已 done 的跳过、failed 的可重试、pending 的重新调度。

import json, redis

r = redis.Redis()

def checkpoint(node_id, status, artifact=None):
r.hset(f"agent:dag:{node_id}", mapping={
"status": status,
"artifact": json.dumps(artifact or {}),
"ts": time.time(),
})

def recover(dag):
"""worker 重启后据此恢复,避免重复执行已完成的节点"""
for nid, n in dag.nodes.items():
s = r.hget(f"agent:dag:{nid}", "status")
if s:
n.status = s.decode()

关键设计:幂等执行。同一节点可能因重试被执行两次,执行函数必须保证重复执行不产生副作用(用 artifact 引用去重、用唯一 key 写库)。这是工程里最常被忽略却最致命的一点。

五、失败处理与重试策略

失败类型处理注意
工具超时 指数退避重试 设最大次数,避免死循环
限流 429 入队等待 + 退避 全局令牌桶
产物缺失 依赖节点失败→本节点标 failed 上游失败不盲目重试
worker 崩溃 节点回 pending 重调度 依赖幂等
逻辑错误 标记 failed,人工/LLM 复盘 别无限重试

import random

def with_retry(fn, max_retry=3, base=1.0):
for i in range(max_retry):
try:
return fn()
except TransientError as e:
if i == max_retry 1:
raise
time.sleep(base * (2 ** i) + random.random())

六、与 LLM 编排的关系

注意区分两层:LLM 负责"决策下一步做什么"(ReAct / 规划),编排器负责"把决策变成可靠执行"。理想架构是 LLM 产出计划(DAG 草案),编排器落地执行与容错,LLM 只在节点失败或需动态分支时再介入。这样既利用 LLM 的灵活性,又用确定性的调度保证可靠性。

面试速答

问:Agent 长任务怎么保证 worker 挂了不重头再来?
答:状态外置 + 节点级检查点。把 DAG 状态和每个节点的完成标记/产物引用存在 Redis/DB,worker 崩溃重启后扫描 DAG,done 跳过、failed 重试、pending 重调度,只重跑失败部分。

问:多个子任务有依赖怎么调度?
答:用 DAG 建模,节点是独立动作、边是依赖。调度器只拉"入度全 done"的节点并行执行,依赖未满足的等待,天然支持部分并行部分串行。

问:重试会不会导致重复副作用?
答:会,所以执行函数必须幂等——用 artifact 唯一 key 去重、写库用 upsert、工具调用带幂等令牌。这是生产级 Agent 的硬要求。

问:LLM 和编排器怎么分工?
答:LLM 负责决策(规划/下一步),编排器负责可靠执行(调度/容错/重试/状态)。LLM 产计划,编排器落地,失败或需动态分支时 LLM 再介入。

高频追问清单

  • DAG 出现环(LLM 规划产出循环依赖)怎么检测与打破?
  • 跨机执行时状态存储选 Redis 还是数据库?一致性与性能如何权衡?
  • 节点产物很大(如大文件)怎么存?checkpoint 只存引用还是存内容?
  • 怎样防止"脑裂"——两个 worker 同时执行同一节点?
  • 动态新增子任务(执行中发现还要做 X)怎么并入正在跑的 DAG?
  • 每工具维度的限流队列怎么设计?令牌桶还是漏桶?
  • 长任务 Human-in-the-loop 审批节点怎么嵌进 DAG?
  • 如何对分布式 Agent 做端到端可观测(trace 串联所有节点)?
  • 重试次数用尽仍失败的节点,怎么优雅降级而非整任务失败?
  • 多 Agent 协作时,任务编排和 Agent 间通信怎么解耦?
  • 赞(0)
    未经允许不得转载:网硕互联帮助中心 » 智能体面试准备(七十七):智能体分布式执行与任务编排工程——长任务分片、依赖调度与状态一致性
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!