跳转到内容
搜索文档

将 Agent 与 Workflow 配合使用

最后更新 查看 MarkdownAgent 设置

什么是 Workflow?

Cloudflare Workflow 为需要 survive 失败、自动重试并等待外部事件的任务提供持久化多步 execution。与 Agent 集成时,Workflow 处理长时间运行的后台处理,Agent 管理实时通信。

Agent 与 Workflow

Agent 和 Workflow 具有互补优势:

能力 Agent Workflow
Execution 模型 在事件上唤醒的长生命周期身份 运行至完成
实时通信 WebSocket、HTTP 流式传输 不支持
状态持久化 内置 SQL 数据库 步骤级持久化
失败处理 应用定义 自动重试和恢复
外部事件 直接处理 暂停并等待事件
用户交互 直接(聊天、UI) 通过 Agent 回调

Agent 可循环、分支并直接与用户交互。Workflow 按顺序执行步骤,保证交付,并可暂停数天等待审批或外部数据。

何时使用各自

仅使用 Agent:

  • 聊天和消息应用
  • 快速 API 调用和响应
  • 实时协作功能
  • 30 秒以内的任务
  • 使用 submitMessages() 的单次持久化 Think 聊天轮次

Agent 与 Workflow 配合:

  • 数据处理管道
  • 报告生成
  • Human-in-the-loop 审批流程
  • 需要保证交付的任务
  • 需要重试的多步操作

仅使用 Workflow:

  • 有或无用户审批的后台作业
  • 定时数据同步
  • 事件驱动处理管道

Agent 与 Workflow 如何通信

AgentWorkflow 类(从 agents/workflows 导入)提供 Workflow 与其发起方 Agent 之间的双向通信。

Workflow 到 Agent

Workflow 可通过多种机制与 Agent 通信:

  • RPC 调用:通过 this.agent 直接调用 Agent 方法,具有完整类型安全
  • 进度报告:通过 this.reportProgress() 发送进度更新,触发 Agent 回调
  • 状态更新:通过 step.updateAgentState()step.mergeAgentState() 修改 Agent 状态,广播给已连接客户端
  • 客户端广播:通过 this.broadcastToClients() 向所有 WebSocket 客户端发送消息
// Inside a workflow's run() method
await this.agent.updateTaskStatus(taskId, "processing"); // RPC call
await this.reportProgress({ step: "process", percent: 0.5 }); // Progress (non-durable)
this.broadcastToClients({ type: "update", taskId }); // Broadcast (non-durable)
await step.mergeAgentState({ taskProgress: 0.5 }); // State update (durable)
// Inside a workflow's run() method
await this.agent.updateTaskStatus(taskId, "processing"); // RPC call
await this.reportProgress({ step: "process", percent: 0.5 }); // Progress (non-durable)
this.broadcastToClients({ type: "update", taskId }); // Broadcast (non-durable)
await step.mergeAgentState({ taskProgress: 0.5 }); // State update (durable)

Agent 到 Workflow

Agent 可通过以下方式与运行中的 Workflow 交互:

  • 启动 Workflow:使用 runWorkflow() 启动新 Workflow 实例
  • 发送事件:使用 sendWorkflowEvent() 分发事件
  • 审批/拒绝:使用 approveWorkflow() / rejectWorkflow() 响应审批请求
  • Workflow 控制:暂停、恢复、终止或重启 Workflow
  • 状态查询:使用 getWorkflow() / getWorkflows() 检查 Workflow 进度

持久化与非持久化操作

理解持久性对有效使用 Workflow 至关重要:

非持久化(重试时可能重复)

这些操作轻量,适合频繁更新,但 Workflow 重试时可能执行多次:

  • this.reportProgress() — 进度报告
  • this.broadcastToClients() — WebSocket 广播
  • this.agent 的直接 RPC 调用

持久化(幂等,不会重复)

这些操作使用 step 参数,保证恰好执行一次:

  • step.do() — 执行持久化步骤
  • step.reportComplete() / step.reportError() — 完成报告
  • step.sendEvent() — 自定义事件
  • step.updateAgentState() / step.mergeAgentState() — 状态同步

持久性保证

Workflow 通过基于步骤的 execution 提供持久性:

  1. 步骤完成是永久的 — 步骤完成后,即使 Workflow 重启也不会重新执行
  2. 自动重试 — 失败步骤以可配置退避重试
  3. 事件持久化 — Workflow 可等待事件最长一年
  4. 状态恢复 — Workflow 状态在基础设施故障中 survive

这意味着 Workflow 适合必须保留部分完成的任务,例如多阶段数据处理或跨多个系统的交易。

Workflow 跟踪

Agent 使用 runWorkflow() 启动 Workflow 时,Workflow 会自动跟踪在 Agent 内部数据库中。这支持:

  • 按 ID、名称或元数据查询 Workflow 状态,支持基于 cursor 的分页
  • 通过生命周期回调(onWorkflowProgressonWorkflowCompleteonWorkflowError)监控进度
  • Workflow 控制:暂停、恢复、终止、重启
  • 使用 deleteWorkflow() / deleteWorkflows() 清理已完成的 Workflow 记录
  • 通过元数据将 Workflow 与用户或会话关联

常见模式

带进度的后台处理

Agent 接收请求,为繁重处理启动 Workflow,并在 Workflow 执行各步骤时向已连接客户端广播进度更新。

// Workflow reports progress after each item
for (let i = 0; i < items.length; i++) {
	await step.do(`process-${i}`, async () => processItem(items[i]));
	await this.reportProgress({
		step: `process-${i}`,
		percent: (i + 1) / items.length,
		message: `Processed ${i + 1}/${items.length}`,
	});
}
// Workflow reports progress after each item
for (let i = 0; i < items.length; i++) {
	await step.do(`process-${i}`, async () => processItem(items[i]));
	await this.reportProgress({
		step: `process-${i}`,
		percent: (i + 1) / items.length,
		message: `Processed ${i + 1}/${items.length}`,
	});
}

Human-in-the-loop 审批

Workflow 准备请求,使用 waitForApproval() 暂停等待审批,Agent 提供 UI 供用户通过 approveWorkflow() / rejectWorkflow() 批准或拒绝。Workflow 根据决定恢复或抛出 WorkflowRejectedError

resilient 外部 API 调用

Workflow 将外部 API 调用包装在带重试逻辑的持久化步骤中。若 API 失败或 Workflow 重启,已完成的调用不会重复,失败的调用自动重试。

const result = await step.do(
	"call-api",
	{
		retries: { limit: 5, delay: "10 seconds", backoff: "exponential" },
		timeout: "5 minutes",
	},
	async () => {
		const response = await fetch("https://api.example.com/process");
		if (!response.ok) throw new Error(`API error: ${response.status}`);
		return response.json();
	},
);
const result = await step.do(
	"call-api",
	{
		retries: { limit: 5, delay: "10 seconds", backoff: "exponential" },
		timeout: "5 minutes",
	},
	async () => {
		const response = await fetch("https://api.example.com/process");
		if (!response.ok) throw new Error(`API error: ${response.status}`);
		return response.json();
	},
);

状态同步

Workflow 在关键里程碑使用 step.updateAgentState()step.mergeAgentState() 更新 Agent 状态。这些状态变更广播给所有已连接客户端,无需轮询即可保持 UI 同步。

相关资源

这篇文档对您有帮助吗?