欢迎光临
我们一直在努力

独立产品 AI 工作流编排架构设计:多步骤 AI 任务的依赖管理与容错

独立产品 AI 工作流编排架构设计:多步骤 AI 任务的依赖管理与容错

一、独立产品 AI 化的策略性挑战:从单次调用到工作流编排的复杂度跃迁

独立产品引入 AI 能力的第一阶段通常是单次调用:用户输入 → 调用 LLM API → 返回结果。这个模式足以支撑简单的文本生成、翻译或摘要功能。但当一个产品需要实现"上传产品截图 → AI 分析设计问题 → 生成优化建议 → 输出设计规范 → 导出 Figma 插件配置"这样的多步骤链路时,单次调用的架构会迅速崩塌。

多步骤 AI 工作流的核心挑战包括:

  • 任务依赖管理:步骤 2 的输入依赖于步骤 1 的输出,如果步骤 1 返回格式不符合预期,整个链路中断。
  • 中间状态持久化:用户需要看到每一步的进度、中间结果,并且支持从任意步骤重试。
  • 容错与降级:LLM 调用存在一定的不确定性(响应超时、格式错误、内容审核拦截),需要有兜底策略。
  • 成本控制:多步骤调用可能产生数倍于单次调用的 Token 消耗,需要合理的中断和缓存机制。
  • 并发编排:部分步骤可以并行执行(如图像分析和文本分析),需要在 DAG(有向无环图)中识别并行机会。

传统方案通常直接在业务代码中用 await 串联调用,但这种写法在步骤超过 3 步时就会变得难以维护——错误处理混杂在业务逻辑中、重试策略分散在各处、步骤间状态传递依赖闭包变量。

二、AI 工作流编排的核心抽象

2.1 步骤(Step)的定义

每一个 AI 工作流步骤需要包含以下维度:

  • 输入/输出定义:使用 TypeScript 泛型约束步骤间的数据传递类型。
  • 执行策略:最大重试次数、退避间隔、超时时间、是否可跳过。
  • 降级函数:当主逻辑失败时执行的备选方案。例如 LLM 调用失败时返回缓存结果或预设默认值。
  • 依赖声明:明确声明该步骤依赖哪些前置步骤的输出。

2.2 DAG 的执行调度

工作流的依赖关系构成一个有向无环图(DAG)。调度器通过拓扑排序确定执行顺序——入度为 0 的节点可以立即执行;当一个节点执行完成后,将其后继节点的入度减 1,入度变为 0 的后继节点进入就绪队列。

并行执行的关键在于正确识别独立子图。两个节点如果没有直接或间接的依赖关系,它们分属不同的拓扑层级,可以并行调度。对于 I/O 密集型的 AI 调用(如同时向两个不同的 LLM 服务发送请求),并行执行可以显著缩短总耗时。

2.3 中间状态与检查点

工作流执行过程中需要持久化每个步骤的状态。状态数据应包含:

  • 步骤标识和执行状态(pending / running / completed / failed / skipped)
  • 步骤的输出数据(用于传递给下游步骤)
  • 执行时间戳和重试次数
  • 错误信息(用于调试和重试决策)

这些数据可以存储在 Redis 或文件系统中。当工作流中断后,用户可以从上次失败或最后一个成功步骤恢复执行,而不必从头开始——这对 Token 成本控制至关重要。

三、轻量级 AI 工作流引擎的实现

/**
* 独立产品 AI 工作流编排引擎
* 支持 DAG 依赖管理、并行调度、断点续传和降级策略
*/

type StepStatus = 'pending' | 'running' | 'completed' | 'failed' | 'skipped';

interface StepDefinition<TInput, TOutput> {
id: string;
name: string;
/** 依赖的前置步骤 ID 列表 */
dependsOn: string[];
/** 步骤执行的最大重试次数 */
maxRetries?: number;
/** 重试退避间隔(毫秒) */
retryBackoffMs?: number;
/** 步骤超时时间(毫秒) */
timeoutMs?: number;
/** 步骤失败时是否可以跳过 */
skippable?: boolean;
/** 输入数据转换:从上游输出中提取本步骤所需输入 */
extractInput: (context: WorkflowContext) => TInput;
/** 核心执行函数 */
execute: (input: TInput) => Promise<TOutput>;
/** 降级函数:execute 失败时的备选方案 */
fallback?: (input: TInput, error: Error) => Promise<TOutput>;
}

interface StepResult {
stepId: string;
status: StepStatus;
output?: unknown;
error?: string;
startedAt: number;
completedAt: number;
retryCount: number;
}

interface WorkflowContext {
/** 步骤输出数据映射表 */
results: Map<string, StepResult>;
}

interface WorkflowDefinition {
id: string;
steps: StepDefinition<any, any>[];
/** 全局超时(毫秒) */
globalTimeoutMs?: number;
}

interface WorkflowState {
workflowId: string;
status: 'running' | 'completed' | 'failed';
steps: Map<string, StepResult>;
startedAt: number;
}

class AIWorkflowEngine {
private runningWorkflows = new Map<string, AbortController>();

/**
* 执行工作流
*/
async execute(
definition: WorkflowDefinition,
options: {
/** 断点续传:从已有的 workflowState 恢复 */
resumeFrom?: WorkflowState;
} = {}
): Promise<WorkflowState> {
const abortController = new AbortController();
this.runningWorkflows.set(definition.id, abortController);

// 构建步骤索引表
const stepMap = new Map<string, StepDefinition<any, any>>();
for (const step of definition.steps) {
stepMap.set(step.id, step);
}

// 初始化上下文和状态
const context: WorkflowContext = { results: new Map() };
const stepResults = new Map<string, StepResult>();

// 断点续传:恢复已完成步骤的结果
if (options.resumeFrom) {
for (const [id, result] of options.resumeFrom.steps) {
if (result.status === 'completed' || result.status === 'skipped') {
stepResults.set(id, result);
context.results.set(id, result);
}
}
}

// 计算入度(依赖计数)
const inDegree = new Map<string, number>();
const successors = new Map<string, string[]>(); // 下游步骤

for (const step of definition.steps) {
// 断点续传:已完成的步骤不需要重新执行
if (stepResults.has(step.id)) {
inDegree.set(step.id, 0);
successors.set(step.id, []);
continue;
}
inDegree.set(step.id, step.dependsOn.length);
if (!successors.has(step.id)) {
successors.set(step.id, []);
}
}

// 构建后继关系
for (const step of definition.steps) {
for (const depId of step.dependsOn) {
const succ = successors.get(depId) || [];
succ.push(step.id);
successors.set(depId, succ);
}
}

// 找到所有入度为 0 的节点(就绪队列)
const readyQueue: string[] = [];
for (const [stepId, degree] of inDegree) {
if (degree === 0 && !stepResults.has(stepId)) {
readyQueue.push(stepId);
}
}

const startedAt = options.resumeFrom?.startedAt ?? Date.now();
let hasFailure = false;

// 核心调度循环:只要还有就绪或运行中的步骤,就继续
while (readyQueue.length > 0 || this.hasRunningStep(context)) {
// 检查全局超时
if (
definition.globalTimeoutMs &&
Date.now() – startedAt > definition.globalTimeoutMs
) {
throw new Error(`工作流执行超时(${definition.globalTimeoutMs}ms)`);
}

if (abortController.signal.aborted) {
throw new Error('工作流已被取消');
}

// 从就绪队列中取步骤并行执行
const batch: string[] = [];
while (readyQueue.length > 0) {
const id = readyQueue.shift()!;
batch.push(id);
}

if (batch.length > 0) {
const promises = batch.map((stepId) =>
this.executeStep(stepMap.get(stepId)!, context, abortController.signal)
);

const results = await Promise.allSettled(promises);

for (let i = 0; i < batch.length; i++) {
const stepId = batch[i];
const result = results[i];

if (result.status === 'fulfilled') {
const stepResult = result.value;
stepResults.set(stepId, stepResult);
context.results.set(stepId, stepResult);

if (stepResult.status === 'failed') {
hasFailure = true;
// 如果某步骤失败且不可跳过,整个工作流标记为失败
const step = stepMap.get(stepId)!;
if (!step.skippable) {
return this.buildFailedState(definition.id, stepResults, startedAt);
}
}

// 将后继步骤的入度减 1
for (const succId of successors.get(stepId) || []) {
const currentDegree = (inDegree.get(succId) || 0) – 1;
inDegree.set(succId, currentDegree);
if (currentDegree === 0) {
readyQueue.push(succId);
}
}
} else {
// Promise 被拒绝(系统级错误)
hasFailure = true;
const error = result.reason instanceof Error ? result.reason : new Error(String(result.reason));
const stepResult: StepResult = {
stepId,
status: 'failed',
error: error.message,
startedAt: 0,
completedAt: Date.now(),
retryCount: 0,
};
stepResults.set(stepId, stepResult);
context.results.set(stepId, stepResult);
}
}
}

// 如果就绪队列为空且没有失败,短暂等待检查是否有新节点就绪
// (处理所有节点完成的情况)
if (readyQueue.length === 0 && !hasFailure && this.allCompleted(stepResults, definition.steps)) {
break;
}
}

// 最终状态
return {
workflowId: definition.id,
status: hasFailure ? 'failed' : 'completed',
steps: stepResults,
startedAt,
};
}

/**
* 执行单个步骤(含重试和降级)
*/
private async executeStep(
step: StepDefinition<any, any>,
context: WorkflowContext,
signal: AbortSignal
): Promise<StepResult> {
const maxRetries = step.maxRetries ?? 2;
const backoff = step.retryBackoffMs ?? 1000;
const startedAt = Date.now();
let lastError: Error | undefined;

for (let attempt = 0; attempt <= maxRetries; attempt++) {
if (signal.aborted) {
return {
stepId: step.id,
status: 'failed',
error: '工作流已被取消',
startedAt,
completedAt: Date.now(),
retryCount: attempt,
};
}

try {
// 从上下文中提取本步骤的输入
const input = step.extractInput(context);

// 执行核心逻辑,添加超时控制
const output = await this.withTimeout(
step.execute(input),
step.timeoutMs ?? 60_000,
signal
);

return {
stepId: step.id,
status: 'completed',
output,
startedAt,
completedAt: Date.now(),
retryCount: attempt,
};
} catch (error) {
lastError = error instanceof Error ? error : new Error(String(error));
if (attempt < maxRetries) {
// 退避等待后重试
await this.sleep(backoff * Math.pow(2, attempt));
}
}
}

// 所有重试都失败,尝试降级
if (step.fallback) {
try {
const input = step.extractInput(context);
const fallbackOutput = await step.fallback(input, lastError!);
return {
stepId: step.id,
status: 'completed', // 降级成功视为完成
output: fallbackOutput,
startedAt,
completedAt: Date.now(),
retryCount: maxRetries + 1,
};
} catch (fallbackError) {
lastError =
fallbackError instanceof Error ? fallbackError : new Error(String(fallbackError));
}
}

return {
stepId: step.id,
status: 'failed',
error: lastError?.message ?? '未知错误',
startedAt,
completedAt: Date.now(),
retryCount: maxRetries + 1,
};
}

/**
* 带超时和取消信号的 Promise 包装
*/
private async withTimeout<T>(
promise: Promise<T>,
timeoutMs: number,
signal: AbortSignal
): Promise<T> {
return new Promise<T>((resolve, reject) => {
const timer = setTimeout(() => reject(new Error('步骤执行超时')), timeoutMs);
const onAbort = () => {
clearTimeout(timer);
reject(new Error('工作流已取消'));
};
signal.addEventListener('abort', onAbort, { once: true });

promise
.then((result) => {
clearTimeout(timer);
signal.removeEventListener('abort', onAbort);
resolve(result);
})
.catch((error) => {
clearTimeout(timer);
signal.removeEventListener('abort', onAbort);
reject(error);
});
});
}

/**
* 取消运行中的工作流
*/
cancel(workflowId: string): void {
const controller = this.runningWorkflows.get(workflowId);
if (controller) {
controller.abort();
this.runningWorkflows.delete(workflowId);
}
}

/**
* 检查是否还有步骤在执行中
*/
private hasRunningStep(context: WorkflowContext): boolean {
// 实际实现中需要维护运行中步骤的计数器
// 这里为简化不做实时追踪
return false;
}

/**
* 检查所有步骤是否已完成
*/
private allCompleted(
results: Map<string, StepResult>,
steps: StepDefinition<any, any>[]
): boolean {
return steps.every((s) => {
const r = results.get(s.id);
return r && (r.status === 'completed' || r.status === 'skipped' || r.status === 'failed');
});
}

private buildFailedState(
id: string,
results: Map<string, StepResult>,
startedAt: number
): WorkflowState {
return {
workflowId: id,
status: 'failed',
steps: results,
startedAt,
};
}

private sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}
}

// —- 使用示例:截图分析工作流 —-

async function createScreenshotAnalysisWorkflow() {
const engine = new AIWorkflowEngine();

const workflow: WorkflowDefinition = {
id: 'screenshot-analysis',
globalTimeoutMs: 120_000,
steps: [
{
id: 'analyze-ui',
name: 'AI 分析截图 UI 结构',
dependsOn: [],
maxRetries: 2,
timeoutMs: 30_000,
extractInput: (ctx) => ({
imageUrl: '2026-09-17zayt5jftqfj.png',
}),
execute: async (input) => {
// 调用视觉模型分析截图
return { components: ['Header', 'Card', 'Footer'], layout: 'vertical-stack' };
},
fallback: async (input, error) => {
// LLM 不可用时返回缓存的预设分析结果
return { components: [], layout: 'unknown', fallback: true };
},
},
{
id: 'generate-suggestions',
name: '生成优化建议',
dependsOn: ['analyze-ui'],
maxRetries: 1,
extractInput: (ctx) => {
const uiResult = ctx.results.get('analyze-ui')!;
return { components: uiResult.output };
},
execute: async (input) => {
return { suggestions: ['增加间距', '统一圆角'] };
},
},
],
};

return engine.execute(workflow);
}

export { AIWorkflowEngine };
export type { StepDefinition, WorkflowDefinition, WorkflowState, WorkflowContext, StepResult };

四、工作流编排的边界条件与失效模式

4.1 步骤间类型安全

当步骤 A 的输出类型变更时,步骤 B 的输入提取逻辑不会自动感知到。纯 TypeScript 的类型系统无法在运行时保证这种跨步骤的类型一致性。两个补救措施:

  • 使用 Zod 或 io-ts 在每个步骤的输入提取函数中做运行时校验,将不符合预期的上游输出转化为明确的错误消息。
  • 为工作流定义编写集成测试,验证完整链路的数据流转——从第一个步骤的输出断言到最后一个步骤的输入断言。

4.2 循环依赖检测

在 DAG 的有效性校验中,循环依赖的检测是启动前必须执行的安全检查。如果步骤 A 依赖步骤 B,步骤 B 又依赖步骤 A,拓扑排序将无法完成,调度器会陷入死锁。实现一个基于 DFS 的环检测算法,在工作流注册阶段就拒绝包含环的 DAG 定义。

4.3 部分成功场景的处理

工作流中如果前 3 个步骤成功、第 4 个步骤失败,已产生的结果和消耗的 Token 是否应该被丢弃?这取决于具体业务语义:对于生成式任务(如生成长文),中间结果仍有价值,可以标记为"部分完成"供用户查看;对于事务性任务(如订单处理),需要全有或全无的回滚语义。

五、总结

AI 工作流编排的核心在于将多步骤 LLM 调用的依赖关系建模为 DAG,通过拓扑排序实现自动调度,并在步骤层面提供重试、降级和超时控制。对于独立产品,不需要引入 LangChain 或 Temporal 等重型框架——一个自实现的 DAG 调度器配合 Redis 状态存储,足以支撑数十个步骤的 AI 任务编排。

落地路线从最关键的三个能力开始:步骤的依赖声明与拓扑排序(第一优先级)、中间状态的持久化与断点续传(降低重试成本)、以及步骤级的降级策略(保障用户体验)。并发调度和全局超时可以在后续版本中引入,作为性能优化而非功能阻塞项。

赞(0)
未经允许不得转载:171主机测评 » 独立产品 AI 工作流编排架构设计:多步骤 AI 任务的依赖管理与容错
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址