用 DAG 解决 AI 工作流的阻塞问题
·
用 DAG 解决 AI 工作流的阻塞问题
嵌套地狱与同步延迟
做企业级 AI 应用时,经常要把多个大模型调用、API 请求和数据库查询串起来。如果只用简单的线性调用,代码很快就会变成“嵌套地狱”,维护起来很痛苦。
更麻烦的是同步阻塞带来的延迟。比如邮件处理流程里,“大模型分类”和“RAG 检索”本来是独立的,如果串行执行,总时间就是两者之和。用有向无环图(DAG)就能解决:在满足依赖关系的前提下,让能并行的节点一起跑,把总耗时压缩到最慢的那个节点的时间。
DAG 怎么工作
DAG 里,每个操作是一个节点,依赖关系是有向边。调度前先用 Kahn 算法做拓扑排序,检查有没有环路,再确定执行顺序。
graph LR
A[工作流入口] --> B[节点 A: 用户输入清洗]
B --> C[节点 B: 情感倾向分析 LLM]
B --> D[节点 C: 本地 FAQ 特征检索]
C --> E[节点 D: 智能邮件草稿生成]
D --> E
E --> F[工作流出口]
style C fill:#bbf,stroke:#333,stroke-width:2px
style D fill:#bbf,stroke:#333,stroke-width:2px
style E fill:#afa,stroke:#333,stroke-width:2px
调度时,节点 B 和 C 都依赖 A,且互不依赖,调度器会并发执行它们。总耗时取决于 B 和 C 中较慢的那个,而不是两者相加。
代码实现
下面是用 Node.js 写的一个简易 DAG 引擎,带环路检测和异步调度。
class WorkflowTask {
constructor(id, action) {
this.id = id;
this.action = action;
this.dependencies = [];
this.status = 'PENDING';
this.output = null;
}
dependsOn(depId) {
this.dependencies.push(depId);
}
}
class MicroWorkflowEngine {
constructor() {
this.tasks = new Map();
}
registerTask(task) {
this.tasks.set(task.id, task);
}
computeTopologicalOrder() {
const inDegree = new Map();
const adjacency = new Map();
const order = [];
for (const [id, _] of this.tasks) {
inDegree.set(id, 0);
adjacency.set(id, []);
}
for (const [id, task] of this.tasks) {
task.dependencies.forEach(depId => {
if (!this.tasks.has(depId)) {
throw new Error(`节点 [${id}] 依赖的节点 [${depId}] 尚未注册!`);
}
adjacency.get(depId).push(id);
inDegree.set(id, inDegree.get(id) + 1);
});
}
const queue = [];
for (const [id, deg] of inDegree.entries()) {
if (deg === 0) queue.push(id);
}
while (queue.length > 0) {
const curr = queue.shift();
order.push(curr);
const neighbors = adjacency.get(curr);
neighbors.forEach(nextId => {
inDegree.set(nextId, inDegree.get(nextId) - 1);
if (inDegree.get(nextId) === 0) {
queue.push(nextId);
}
});
}
if (order.length !== this.tasks.size) {
throw new Error("工作流拓扑校验异常: 依赖图中存在死锁环路,初始化失败!");
}
return order;
}
async run(ctx) {
const order = this.computeTopologicalOrder();
console.log("拓扑序列解析成功. 执行链条优先级:", order.join(' -> '));
const runningJobs = new Map();
const results = { ...ctx };
while (true) {
let activeTaskLaunched = false;
let unresolvedTasks = false;
for (const [id, task] of this.tasks) {
if (task.status === 'FINISHED' || task.status === 'ERROR') continue;
unresolvedTasks = true;
if (task.status === 'RUNNING') continue;
const ready = task.dependencies.every(depId => {
const t = this.tasks.get(depId);
return t && t.status === 'FINISHED';
});
if (ready) {
task.status = 'RUNNING';
activeTaskLaunched = true;
const promise = (async () => {
try {
const depData = {};
task.dependencies.forEach(depId => {
depData[depId] = this.tasks.get(depId).output;
});
task.output = await task.action(results, depData);
task.status = 'FINISHED';
} catch (err) {
task.status = 'ERROR';
throw err;
}
})();
runningJobs.set(id, promise);
}
}
if (!unresolvedTasks) break;
if (!activeTaskLaunched && runningJobs.size === 0) {
throw new Error("工作流执行挂起故障,陷入死锁");
}
await Promise.race(runningJobs.values());
for (const [id, p] of runningJobs) {
const t = this.tasks.get(id);
if (t.status === 'FINISHED' || t.status === 'ERROR') {
runningJobs.delete(id);
}
}
}
const finalOutput = {};
for (const [id, task] of this.tasks) {
finalOutput[id] = task.output;
}
return finalOutput;
}
}
// 测试
(async () => {
const engine = new MicroWorkflowEngine();
const task1 = new WorkflowTask("Sanitize", async (ctx) => ctx.text.trim());
const task2 = new WorkflowTask("AnalyzeSentiment", async (ctx, deps) => {
await new Promise(resolve => setTimeout(resolve, 400));
return deps.Sanitize.includes("赞") ? "POSITIVE" : "NEUTRAL";
});
task2.dependsOn("Sanitize");
const task3 = new WorkflowTask("Keywords", async (ctx, deps) => {
return deps.Sanitize.split(' ');
});
task3.dependsOn("Sanitize");
const task4 = new WorkflowTask("Report", async (ctx, deps) => {
return `倾向: ${deps.AnalyzeSentiment} | 词数: ${deps.Keywords.length}`;
});
task4.dependsOn("AnalyzeSentiment");
task4.dependsOn("Keywords");
engine.registerTask(task1);
engine.registerTask(task2);
engine.registerTask(task3);
engine.registerTask(task4);
const out = await engine.run({ text: "这个产品 赞" });
console.log("工作流执行完毕。输出数据报表:", out);
})();
分布式部署的坑
单机内存跑 DAG 很快,但真要上线到分布式环境,得考虑几个实际问题:
- 状态持久化:内存调度零网络开销,但一旦进程挂掉或云实例被抢占(Spot),工作流状态就丢了。用 Temporal 或 Redis 做状态锁能恢复,但每次状态转移都要网络 IO,延迟会明显上升。
- 重试与计费:下游节点(比如发短信)超时重试时,如果上游没做幂等,可能重复调用大模型,导致计费失控。大模型节点得用唯一主键防重复提交。
- 动态路由:静态 DAG 编译期就能检测环路。但大模型工作流经常要根据 LLM 输出动态决定下一步(Dynamic Routing)。支持动态路由意味着拓扑结构得动态扩展,依赖链和追踪会变得很乱。
小结
核心思路就是用图模型替代冗余的条件判断。通过 Kahn 算法做静态合法性检查,配合异步并发调度,能用较少的代码和服务器开销驱动多个大模型操作并发执行,降低延迟。
质量评分
| 维度 | 评估标准 | 得分 |
|---|---|---|
| 直接性 | 直接陈述事实还是绕圈宣告? | 9/10 |
| 节奏 | 句子长度是否变化? | 8/10 |
| 信任度 | 是否尊重读者智慧? | 9/10 |
| 真实性 | 听起来像真人说话吗? | 9/10 |
| 精炼度 | 还有可删减的内容吗? | 9/10 |
| 总分 | 44/50 |
所做更改:
- 删除了“嵌套地狱与等待瓶颈”、“多代理流式调度拓扑模型”等夸张标题,改为更直接的“嵌套地狱与同步延迟”、"DAG 怎么工作”。
- 去除了“毫秒级”、“闪电般”、“产品变现”、“底层支撑”等营销词汇。
- 简化了代码注释,去除了“快速运行测试”等冗余说明。
- 调整了段落结构,去除了“一、二、三、四、五”编号,改用自然标题。
- 增加了开发者口吻,如“维护起来很痛苦”、“得考虑几个实际问题”、“核心思路就是”。
- 保留了核心技术细节,如 Kahn 算法、幂等性、Spot Preemption 等。
- 去除了“此外”、“然而”、“不仅...而且...”等 AI 常用连接词。
- 避免了“作为...的证明”、“标志着”等夸大的象征意义。
- 简化了“总结”部分,去除了“为产品变现提供低延迟的底层支撑”等夸张表述。
更多推荐


所有评论(0)