为下面三个工作流选择合适的状态管理模式(脚本变量、检查点、外部存储)并说明理由:
Level 1: 选择状态管理模式工作流 A: 批量压缩 20 张图片,每张耗时 5 秒,总共 100 秒
工作流 B: 训练一个机器学习模型,需要 50 个 epoch,每个 epoch 10 分钟,总共 500 分钟(8 小时)
工作流 C: 审查 100 个 PR,每个 PR 需要人工批准后才能合并,整个流程可能持续几天
学习目标:
- 区分工作流状态和 Agent 上下文
- 掌握三种状态管理模式
- 理解检查点和恢复机制
前置要求:第 3 课:如何拆解复杂任务 | 下一课 第 5 课 >>
你设计了一个完美的工作流,10 个步骤,清晰的依赖关系。第 8 步时,服务器重启了。工作流崩溃。
重新运行?那前 7 步的工作(可能花了 30 分钟)全部白费。
这就是没有状态管理的代价。
状态管理解决三个问题:1
没有状态管理,Agent 只能通过对话历史传递信息。对话历史会溢出,会丢失,会被 Agent 遗忘。
有了状态管理,工作流有一个明确的"内存",持久化、可查询、可恢复。2
这三个词容易混淆,先厘清概念:1
状态(State)
上下文(Context)
内存(Memory)
举个例子:
关键原则: 状态是全局的,上下文是局部的。3
适用场景: 短工作流(< 10 分钟),不需要跨进程或跨机器。
优点: 简单、快速、无需外部依赖。
缺点: 进程崩溃后状态丢失,无法恢复。
状态存在哪里? 在函数的局部变量里(processed、results、errors)。
如果进程崩溃? 状态全部丢失,必须从头开始。
好处: 状态结构清晰,容易传递给其他函数,容易序列化(如果需要持久化)。
适用场景: 中等时长工作流(10-60 分钟),耗时操作后需要保存进度。
优点: 崩溃后可以从最近的检查点恢复,避免重复工作。
缺点: 需要设计检查点位置和恢复逻辑。1
检查点策略:
适用场景: 长时间运行的工作流(> 1 小时)、需要跨机器协调、需要人工审批。
优点: 状态持久化,进程崩溃、机器重启都不影响,支持暂停/恢复。
缺点: 需要外部依赖(数据库、Redis)、增加复杂度。4
关键模式: 状态机(State Machine)2
工作流的阶段就是状态机的状态:
每个阶段转换都保存到外部存储,确保工作流可以从任何阶段恢复。
为什么? 上下文越大,Agent 越容易分心,推理质量下降,成本上升。3
为什么? 结构化上下文更容易被 Agent 理解,也更容易被你调试。
累积式上下文: 每步的结果加到上下文中,越来越大。
重置式上下文: 每步清空上下文,只保留必要信息。
选择: 大部分情况用重置式,避免上下文爆炸。累积式只在后续步骤真的需要前面所有结果时使用(如最后的汇总步骤)。5
好的工作流应该能回答这些问题:
下一课: 第 5 课:错误处理和重试策略 — 学习如何让工作流在失败时优雅恢复,而不是直接崩溃
MachineLearningMastery:AI 代理中的持久化内存和状态的 5 种架构模式 — https://machinelearningmastery.com/5-architectural-patterns-for-persistent-memory-and-state-in-ai-agents/ ↩ ↩2 ↩3
MindStudio:工作流状态 vs 会话状态 — https://www.mindstudio.ai/blog/workflow-state-vs-session-state-ai-agents ↩ ↩2
Chrono Innovation:可扩展的代理 AI 工作流架构 — https://www.chronoinnovation.com/resources/agentic-ai-workflows-architecture/ ↩ ↩2
Appamass:可靠 AI 代理工作流的状态管理模式 — https://appamass.com/en/blog/state-management-patterns-for-reliable-ai-agent-workflows-5yemlru6ui6cacast3l5 ↩
Ranjan Kumar:构建会记忆的代理 — https://ranjankumar.in/building-agents-that-remember-state-management-in-multi-agent-ai-systems ↩
工作流 A: 批量压缩 20 张图片,每张耗时 5 秒,总共 100 秒
工作流 B: 训练一个机器学习模型,需要 50 个 epoch,每个 epoch 10 分钟,总共 500 分钟(8 小时)
工作流 C: 审查 100 个 PR,每个 PR 需要人工批准后才能合并,整个流程可能持续几天
要求:
Jot down thoughts, sticking points, things you didn't get. Written to this course's appendix only — the lesson file is never touched.
// 状态:工作流知道的所有信息
const workflowState = {
phase: 'testing',
filesProcessed: 47,
totalFiles: 100,
issues: [/* 前面步骤发现的所有问题 */],
currentBatch: [/* 当前正在处理的文件 */]
};
// 上下文:给这个 Agent 的信息(从状态中提取)
const agentContext = {
file: workflowState.currentBatch[0],
previousIssues: workflowState.issues.filter(i => i.severity === 'high')
};
// Agent 调用
const result = await agent({
task: '测试文件',
context: agentContext // 只给相关信息,不是全部状态
});
// 更新状态
workflowState.filesProcessed++;
workflowState.issues.push(...result.newIssues);
async function simpleWorkflow(files) {
// 状态就是普通的 JavaScript 变量
let processed = 0;
let results = [];
let errors = [];
for (const file of files) {
try {
const result = await processFile(file);
results.push(result);
processed++;
console.log(`进度: ${processed}/${files.length}`);
} catch (error) {
errors.push({ file, error });
}
}
return { results, errors, total: files.length };
}
async function betterWorkflow(files) {
// 用对象组织状态,更清晰
const state = {
input: { files, total: files.length },
progress: { current: 0, phase: 'processing' },
output: { results: [], errors: [] },
metadata: { startTime: Date.now() }
};
for (const file of state.input.files) {
try {
const result = await processFile(file);
state.output.results.push(result);
state.progress.current++;
} catch (error) {
state.output.errors.push({ file, error });
}
}
state.progress.phase = 'completed';
state.metadata.endTime = Date.now();
state.metadata.duration = state.metadata.endTime - state.metadata.startTime;
return state;
}
async function workflowWithCheckpoints(tasks) {
const checkpointFile = '.workflow-state.json';
// 尝试恢复之前的状态
let state = await loadCheckpoint(checkpointFile) || {
completed: [],
pending: tasks,
phase: 'processing'
};
console.log(`恢复: 已完成 ${state.completed.length}/${tasks.length}`);
while (state.pending.length > 0) {
const task = state.pending.shift();
// 执行任务
const result = await executeTask(task);
state.completed.push({ task, result });
// 检查点:每 10 个任务后保存
if (state.completed.length % 10 === 0) {
await saveCheckpoint(checkpointFile, state);
console.log(`检查点: ${state.completed.length} 个任务已完成`);
}
}
state.phase = 'completed';
await saveCheckpoint(checkpointFile, state);
return state;
}
async function saveCheckpoint(file, state) {
await fs.writeFile(file, JSON.stringify(state, null, 2));
}
async function loadCheckpoint(file) {
try {
const data = await fs.readFile(file, 'utf-8');
return JSON.parse(data);
} catch {
return null; // 文件不存在,从头开始
}
}
// 状态存储接口
class WorkflowStateStore {
constructor(db) {
this.db = db;
}
async save(workflowId, state) {
await this.db.set(`workflow:${workflowId}`, JSON.stringify(state));
}
async load(workflowId) {
const data = await this.db.get(`workflow:${workflowId}`);
return data ? JSON.parse(data) : null;
}
async delete(workflowId) {
await this.db.del(`workflow:${workflowId}`);
}
}
// 使用外部存储的工作流
async function persistentWorkflow(workflowId, tasks) {
const store = new WorkflowStateStore(redis);
// 加载状态(如果存在)
let state = await store.load(workflowId) || {
id: workflowId,
phase: 'init',
completed: [],
pending: tasks,
createdAt: Date.now(),
updatedAt: Date.now()
};
console.log(`工作流 ${workflowId}: 阶段 ${state.phase},
进度 ${state.completed.length}/${tasks.length}`);
// 阶段 1: 处理任务
if (state.phase === 'init' || state.phase === 'processing') {
state.phase = 'processing';
while (state.pending.length > 0) {
const task = state.pending.shift();
const result = await executeTask(task);
state.completed.push({ task, result });
state.updatedAt = Date.now();
// 每个任务后保存状态
await store.save(workflowId, state);
}
state.phase = 'awaiting_approval';
await store.save(workflowId, state);
}
// 阶段 2: 等待人工审批(可能在另一个进程/机器上恢复)
if (state.phase === 'awaiting_approval') {
console.log('等待审批...');
// 这里可以返回,让另一个进程(或几小时后)继续
return { workflowId, status: 'awaiting_approval' };
}
// 阶段 3: 执行最终操作(审批通过后)
if (state.phase === 'approved') {
state.phase = 'finalizing';
await store.save(workflowId, state);
await executeFinalAction(state.completed);
state.phase = 'completed';
state.completedAt = Date.now();
await store.save(workflowId, state);
}
return state;
}
// 审批工作流
async function approveWorkflow(workflowId) {
const store = new WorkflowStateStore(redis);
const state = await store.load(workflowId);
if (!state) throw new Error('工作流不存在');
if (state.phase !== 'awaiting_approval') {
throw new Error(`无法审批: 当前阶段是 ${state.phase}`);
}
state.phase = 'approved';
state.approvedAt = Date.now();
await store.save(workflowId, state);
// 继续执行工作流
return await persistentWorkflow(workflowId, []);
}
init → processing → awaiting_approval → approved → finalizing → completed ↓ rejected → cancelled// ❌ 不好: 把所有状态都给 Agent
const result = await agent({
task: '分析这个文件',
context: workflowState // 包含 100 个文件的分析结果、配置、日志...
});
// ✓ 好: 只给相关信息
const result = await agent({
task: '分析这个文件',
context: {
file: currentFile,
guidelines: workflowState.config.analysisGuidelines,
similarIssues: workflowState.results
.filter(r => r.file.type === currentFile.type)
.slice(0, 3) // 最多 3 个相似案例
}
});
// ❌ 不好: 非结构化的文本
const context = `
之前分析了 47 个文件,发现了 23 个问题。
当前文件是 src/utils.js,大小 350 行。
配置要求检查 SQL 注入和 XSS。
`;
// ✓ 好: 结构化对象
const context = {
progress: { filesAnalyzed: 47, issuesFound: 23 },
currentFile: { path: 'src/utils.js', lines: 350 },
checkTypes: ['sql_injection', 'xss']
};
let context = { task: 'refactor codebase' };
for (const file of files) {
const result = await agent({ task: 'analyze', context });
context.results = context.results || [];
context.results.push(result); // 累积
}
// 最后 context 包含所有文件的结果,可能非常大
const allResults = [];
for (const file of files) {
const context = {
file,
guidelines: config.guidelines,
exampleIssues: allResults.slice(-3) // 只看最近 3 个
};
const result = await agent({ task: 'analyze', context });
allResults.push(result); // 存在工作流状态里,不在上下文里
}
class ObservableWorkflow {
constructor(name, totalSteps) {
this.state = {
name,
totalSteps,
currentStep: 0,
phase: 'init',
startTime: Date.now(),
errors: [],
results: []
};
}
async executeStep(stepName, fn) {
this.state.currentStep++;
this.state.phase = stepName;
console.log(`[${this.state.name}]
步骤 ${this.state.currentStep}/${this.state.totalSteps}:
${stepName}`);
const stepStart = Date.now();
try {
const result = await fn();
this.state.results.push({ stepName, result, duration: Date.now() - stepStart });
return result;
} catch (error) {
this.state.errors.push({ stepName, error: error.message });
throw error;
}
}
getStatus() {
const progress = (this.state.currentStep / this.state.totalSteps) * 100;
const elapsed = Date.now() - this.state.startTime;
const avgStepTime = elapsed / this.state.currentStep;
const remainingSteps = this.state.totalSteps - this.state.currentStep;
const estimatedRemaining = avgStepTime * remainingSteps;
return {
progress: `${progress.toFixed(1)}%`,
currentPhase: this.state.phase,
elapsed: `${(elapsed / 1000).toFixed(1)}s`,
estimatedRemaining: `${(estimatedRemaining / 1000).toFixed(1)}s`,
errors: this.state.errors.length
};
}
}
// 使用
async function myWorkflow() {
const wf = new ObservableWorkflow('数据迁移', 4);
const data = await wf.executeStep('读取源数据', async () => {
return await readSourceData();
});
const transformed = await wf.executeStep('转换格式', async () => {
return await transformData(data);
});
await wf.executeStep('写入目标数据库', async () => {
return await writeToTarget(transformed);
});
await wf.executeStep('验证', async () => {
return await validateMigration();
});
console.log('最终状态:', wf.getStatus());
}