跳转到内容
搜索文档

事件与参数

最后更新 查看 MarkdownAgent 设置

触发 Workflow 时,它可以接收可选事件。该事件可以包含 Workflow 可操作的数据,包括请求详情、从数据库(如 D1 或 KV)或 Webhook 获取的用户数据,或来自 Queue consumer 的消息。

事件是 Workflow 的重要组成部分,因为你通常希望 Workflow 对数据采取行动。由于给定 Workflow 实例持久执行,事件是向 Workflow 提供应不可变(不变化)和/或代表 Workflow 在该时间点需要操作的数据的有用方式。

向 Workflow 传递数据

你可以通过三种方式向 Workflow 传递参数:

  • 从 Worker 触发 Workflow 时,作为 Workflow 绑定(binding)create 方法的可选参数。
  • 使用 wrangler CLI 触发 Workflow 时,通过 --params 标志。
  • 通过 step.waitForEvent API,允许 Workflow 实例在运行中等待事件(及可选数据)。Workflow 实例可通过 HTTP 或 Workflows 的 Workers API 从外部服务接收事件。

你可以传递任何 JSON 可序列化对象作为参数。

export default {
	async fetch(req, env) {
		let someEvent = { url: req.url, createdTimestamp: Date.now() };
		// 触发 Workflow
		// 将事件作为 `create` 方法的第二个参数
		// 传递给我们的 Workflow 绑定。
		let instance = await env.MY_WORKFLOW.create({
			id: crypto.randomUUID(),
			params: someEvent,
		});

		return Response.json({
			id: instance.id,
			details: await instance.status(),
		});
	},
};
export default {
	async fetch(req: Request, env: Env) {
		let someEvent = { url: req.url, createdTimestamp: Date.now() };
		// 触发 Workflow
		// 将事件作为 `create` 方法的第二个参数
		// 传递给我们的 Workflow 绑定。
		let instance = await env.MY_WORKFLOW.create({
			id: crypto.randomUUID(),
			params: someEvent,
		});

		return Response.json({
			id: instance.id,
			details: await instance.status(),
		});
	},
};

要通过 wrangler 命令行界面传递参数,将 JSON 字符串作为第二个参数传递给 workflows trigger 子命令:

npx wrangler@latest workflows trigger workflows-starter '{"some":"data"}'
🚀 Workflow instance "57c7913b-8e1d-4a78-a0dd-dce5a0b7aa30" has been queued successfully

等待事件

运行中的 Workflow 可以通过在 Workflow 内调用 step.waitForEvent 等待事件,允许你通过以下两种方式之一向 Workflow 发送事件:

  1. 通过 Workers API 绑定(binding):调用 instance.sendEvent 向特定 Workflow 实例发送事件。
  2. 使用 REST API (HTTP API) 的事件端点

由于 waitForEventWorkflowStep API 的一部分,你可以在 Workflow 内多次调用它,并使用控制流有条件地等待事件。

调用 waitForEvent 需要指定 type(最多 100 个字符 1),用于在向 Workflow 实例发送事件时匹配相应的 type

例如,等待 billing webhook:

export class MyWorkflow extends WorkflowEntrypoint {
	async run(event, step) {
		// Workflow 中的其他步骤
		let stripeEvent = await step.waitForEvent(
			"receive invoice paid webhook from Stripe",
			{ type: "stripe-webhook", timeout: "1 hour" },
		);
		// Workflow 的其余部分
	}
}
export class MyWorkflow extends WorkflowEntrypoint<Env, Params> {
	async run(event: WorkflowEvent<Params>, step: WorkflowStep) {
		// Workflow 中的其他步骤
		let stripeEvent = await step.waitForEvent<IncomingStripeWebhook>(
			"receive invoice paid webhook from Stripe",
			{ type: "stripe-webhook", timeout: "1 hour" },
		);
		// Workflow 的其余部分
	}
}

上述示例:

  • 使用 typestripe-webhook 调用 waitForEvent - 相应的 sendEvent 调用为 await instance.sendEvent({type: "stripe-webhook", payload: webhookPayload})
  • 使用 TypeScript 类型参数step.waitForEvent 的返回值类型化为 IncomingStripeWebhook
  • 继续 Workflow 的其余部分。

waitForEvent 调用的默认超时为 24 小时,可以通过向 waitForEvent 调用传递 { timeout: WorkflowTimeoutDuration } 作为第二个参数来更改。

let event = await step.waitForEvent("wait for human approval", {
	type: "approval-flow",
	timeout: "15 minutes",
});
let event = await step.waitForEvent(
		"wait for human approval",
		{ type: "approval-flow", timeout: "15 minutes" },
	);

你可以指定 1 秒到最多 365 天之间的超时。

向运行中的 Workflows 发送事件

使用 waitForEvent API 等待事件的 Workflow 实例可以使用 instance.sendEvent API 接收事件:

export default {
	async fetch(req, env) {
		const instanceId = new URL(req.url).searchParams.get("instanceId");
		const webhookPayload = await req.json();

		let instance = await env.MY_WORKFLOW.get(instanceId);
		// 发送事件,`type` 与
		// step.waitForEvent 调用中定义的事件类型匹配
		await instance.sendEvent({
			type: "stripe-webhook",
			payload: webhookPayload,
		});

		return Response.json({
			status: await instance.status(),
		});
	},
};
export default {
	async fetch(req: Request, env: Env) {
		const instanceId = new URL(req.url).searchParams.get("instanceId");
		const webhookPayload = await req.json<Payload>();

		let instance = await env.MY_WORKFLOW.get(instanceId);
		// 发送事件,`type` 与
		// step.waitForEvent 调用中定义的事件类型匹配
		await instance.sendEvent({
			type: "stripe-webhook",
			payload: webhookPayload,
		});

		return Response.json({
			status: await instance.status(),
		});
	},
};
  • 与本指南中的 waitForEvent 示例类似,waitForEventsendEvent 字段中的 type 属性必须匹配。
  • 要向具有多个 waitForEvent 调用的 Workflow 发送多个事件,调用 sendEvent 并设置相应的 type 属性(最多 100 个字符 1)。
  • 也可以使用 REST API (HTTP API) 的事件端点 发送事件。

TypeScript 与类型参数

默认情况下,传递给 Workflow 定义 run 方法的 WorkflowEvent 类型符合以下内容:

export type WorkflowCronSchedule = {
	/** Cron expression that triggered this event. */
	cron: string;
	/** Timestamp of the scheduled trigger, in milliseconds since the Unix epoch. */
	scheduledTime: number;
};

export type WorkflowEvent<T> = {
	/** The data passed as the parameter when the Workflow instance was triggered. */
	payload: Readonly<T>;
	/** The timestamp that the Workflow was triggered. */
	timestamp: Date;
	/** ID of the current Workflow instance. */
	instanceId: string;
	/** Name of the current Workflow. */
	workflowName: string;
	/** Metadata for Workflow instances created by a cron schedule. */
	schedule?: WorkflowCronSchedule;
};

当 Workflow 实例由 Workflow 绑定上配置的 cron 调度创建时,event.schedule 包括创建实例的 cron 表达式和调度触发时间:

export class MyWorkflow extends WorkflowEntrypoint<Env> {
	async run(event: WorkflowEvent<unknown>, step: WorkflowStep) {
		if (event.schedule) {
			console.log(event.schedule.cron);
			console.log(new Date(event.schedule.scheduledTime));
		}
	}
}

你可以通过定义自己的类型并将其作为类型参数传递给 WorkflowEvent 来可选地类型化这些事件:

// 定义符合 Workflow 实例实例化时
// 使用的事件的类型
interface YourEventType {
	userEmail: string;
	createdTimestamp: number;
	metadata?: Record<string, string>;
}

YourEventType 作为类型参数传递给 WorkflowEvent 时,event.payload 属性在整个 Workflow 定义中将具有 YourEventType 类型:

src/index.tsts
// 导入 Workflow 定义
import { WorkflowEntrypoint, WorkflowStep, WorkflowEvent} from 'cloudflare:workers';

export class MyWorkflow extends WorkflowEntrypoint {
	// 将类型作为类型参数传递给 WorkflowEvent
	// 'payload' 属性将具有你的参数类型。
	async run(event: WorkflowEvent<YourEventType>, step: WorkflowStep) {
		let state = await step.do("my first step", async () => {
			// 通过 event.payload 访问属性
          let userEmail = event.payload.userEmail
          let createdTimestamp = event.payload.createdTimestamp
        })

        await step.do("my second step", async () => { /* your code here */ })
	}
}

你也可以在使用 Workers APIcreate 方法创建(触发)Workflow 实例时,为 Workflows 类型提供类型参数。请注意,这不会将类型信息传播到 Workflow 内部,因为 TypeScript 类型是构建时构造。

Footnotes

  1. 匹配模式:^[a-zA-Z0-9_][a-zA-Z0-9-_]*$ 2

这篇文档对您有帮助吗?