跳转到内容
搜索文档

快速入门

最后更新 查看 MarkdownAgent 设置

Cloudflare Queues 是一种灵活的消息队列,允许您将消息排队以进行异步处理。按照本指南,您将创建第一个队列、用于向该队列发布消息的生产者 Worker,以及用于从该队列消费消息的消费者 Worker。

前提条件

要使用 Queues,您需要:

  1. 注册 Cloudflare 账户
  2. 安装 Node.js

Node.js 版本管理器

使用 Voltanvm 等 Node 版本管理器,以避免权限问题并切换 Node.js 版本。本指南后续将介绍的 Wrangler 需要 Node 版本 16.17.0 或更高。

1. 创建 Worker 项目

您将从 Worker(生产者 Worker)访问队列。您必须创建至少一个生产者 Worker 才能向队列发布消息。若您使用 R2 Bucket 事件通知,则不需要生产者 Worker。

要创建生产者 Worker,请运行:

npm create cloudflare@latest -- producer-worker

进行设置时,请选择以下选项:

  • 对于 What would you like to start with?,选择 Hello World example
  • 对于 Which template would you like to use?,选择 Worker only
  • 对于 Which language do you want to use?,选择 TypeScript
  • 对于 Do you want to use git for version control?,选择 Yes
  • 对于 Do you want to deploy your application?,选择 No(部署前我们还会做一些修改)。

这将创建一个新目录,其中包含 src/index.ts Worker 脚本和 wrangler.jsonc 配置文件。创建 Worker 后,您将创建一个 Queue 以供访问。

进入新创建的目录:

cd producer-worker

2. 创建队列

要使用队列,您需要创建至少一个队列来发布和消费消息。

要创建队列,请运行:

npx wrangler queues create <MY-QUEUE-NAME>

选择一个具有描述性且与您打算用于此队列的消息类型相关的名称。描述性队列名称示例:debug-logsuser-clickstream-datapassword-reset-prod

队列名称必须为 1 到 63 个字符。队列名称不能包含破折号(-)以外的特殊字符,且必须以字母或数字开头和结尾。

设置队列名称后无法更改。创建队列后,您将设置生产者 Worker 以访问它。

3. 设置生产者 Worker

要将队列暴露给 Worker 内的代码,需要通过创建绑定(binding)将队列连接到 Worker。绑定(binding) 允许 Worker 访问 Cloudflare 开发者平台上的资源(如 Queues)。

要创建绑定,请打开新生成的 wrangler.jsonc 文件并添加以下内容:

{
	"queues": {
		"producers": [
			{
				"queue": "MY-QUEUE-NAME",
				"binding": "MY_QUEUE"
			}
		]
	}
}
[[queues.producers]]
queue = "MY-QUEUE-NAME"
binding = "MY_QUEUE"

MY-QUEUE-NAME 替换为步骤 2 中创建的队列名称。接下来,将 MY_QUEUE 替换为您想要的 binding 名称。绑定必须是有效的 JavaScript 变量名。这是您在 Worker 中引用此队列时使用的变量。

编写生产者 Worker

现在,您将配置生产者 Worker 以创建要发布到队列的消息。生产者 Worker 将:

  1. 接收来自浏览器的请求。
  2. 将请求转换为 JSON 格式。
  3. 直接将请求写入队列。

在 Worker 项目目录中,打开 src 文件夹并将以下内容添加到 index.ts 文件:

export default {
  async fetch(request, env, ctx): Promise<Response> {
    const log = {
      url: request.url,
      method: request.method,
      headers: Object.fromEntries(request.headers),
    };
    await env.<MY_QUEUE>.send(log);
    return new Response("Success!");
  },
} satisfies ExportedHandler<Env>;

MY_QUEUE 替换为 wrangler.jsonc 文件中设置的绑定名称。

还要在 index.tsEnv 接口中添加队列。

export interface Env {
   <MY_QUEUE>: Queue;
}

若写入失败,Worker 将返回错误(抛出异常)。若写入成功,它将向浏览器返回 HTTP 200 状态码及 Success

在生产应用程序中,您可能会使用 try...catch 语句捕获异常并直接处理(例如,返回自定义错误或重试)。

部署生产者 Worker

配置好 Wrangler 文件和 index.ts 文件后,即可部署生产者 Worker。要部署生产者 Worker,请运行:

npx wrangler deploy

您应看到类似下面的输出,默认包含 *.workers.dev URL。

Uploaded <YOUR-WORKER-NAME> (0.76 sec)
Published <YOUR-WORKER-NAME> (0.29 sec)
  https://<YOUR-WORKER-NAME>.<YOUR-ACCOUNT>.workers.dev

复制 *.workers.dev 子域并在新浏览器标签页中打开。刷新页面几次以开始向队列发布请求。每次将请求写入队列后,浏览器应返回 Success 响应。

您已构建队列和用于向队列发布消息的生产者 Worker。接下来,您将创建消费者 Worker 以消费发布到队列的消息。没有消费者 Worker,消息将保留在队列中直到过期,默认过期时间为四(4)天。

4. 创建消费者 Worker

消费者 Worker 从队列接收消息。当消费者 Worker 收到队列消息时,可以将其写入其他目标,例如日志控制台或存储对象。

在本指南中,您将创建消费者 Worker 并使用 wrangler tail 记录和检查消息。您将在创建生产者 Worker 的同一 Worker 项目中创建消费者 Worker。

要创建消费者 Worker,请打开 index.ts 文件并将以下 queue 处理程序添加到现有 fetch 处理程序:

export default {
  async fetch(request, env, ctx): Promise<Response> {
    const log = {
      url: request.url,
      method: request.method,
      headers: Object.fromEntries(request.headers),
    };
    await env.<MY_QUEUE>.send(log);
    return new Response("Success!");
  },
  async queue(batch, env, ctx): Promise<void> {
    for (const message of batch.messages) {
      console.log("consumed from our queue:", JSON.stringify(message.body));
    }
  },
} satisfies ExportedHandler<Env>;

MY_QUEUE 替换为 wrangler.jsonc 文件中设置的绑定名称。

每次向队列发布消息时,消费者 Worker 的 queue 处理程序(async queue)都会被调用,并传递一条或多条消息。

在本示例中,消费者 Worker 将队列的 JSON 格式消息转换为字符串并记录输出。在实际应用程序中,消费者 Worker 可配置为将消息写入对象存储(如 R2、写入数据库(如 D1、在调用外部 API(如邮件 API)或传统云提供商的数据仓库之前进一步处理消息。

在消费者处理程序中执行异步任务时,请使用 waitUntil() 以确保函数的响应得到处理。在此方法范围内不支持其他异步方法。

将消费者 Worker 连接到队列

配置好消费者 Worker 后,即可将其连接到队列。

每个队列只能连接一个消费者 Worker。若尝试将多个消费者连接到同一队列,在部署该 Worker 时将遇到错误。

要将队列连接到消费者 Worker,请打开 Wrangler 文件并在底部添加:

{
	"queues": {
		"consumers": [
			{
				"queue": "<MY-QUEUE-NAME>",
				// Required: this should match the name of the queue you created in step 3.
				// If you misspell the name, you will receive an error when attempting to publish your Worker.
				"max_batch_size": 10, // optional: defaults to 10
				"max_batch_timeout": 5 // optional: defaults to 5 seconds
			}
		]
	}
}
[[queues.consumers]]
queue = "<MY-QUEUE-NAME>"
max_batch_size = 10
max_batch_timeout = 5

MY-QUEUE-NAME 替换为步骤 2 中创建的队列。

在消费者 Worker 中,您使用 Queues 通过 max_batch_size 选项和 max_batch_timeout 选项自动批处理消息。消费者 Worker 将每 10 条消息或每 5 秒(以先发生者为准)接收一批消息。

max_batch_size(默认为 10)有助于减少消费者 Worker 需要被调用的次数。消费者不会为每条消息都被调用,而是在 10 条消息进入队列后才被调用。

max_batch_timeout(默认为 5 秒)有助于减少等待时间。若生产者 Worker 未向队列发送最多 10 条消息以触发消费者 Worker 调用,消费者 Worker 将每 5 秒被调用一次以接收队列中等待的消息。

部署消费者 Worker

配置好 Wrangler 文件和 index.ts 文件后,通过运行以下命令部署消费者 Worker:

npx wrangler deploy

5. 从队列读取消息

设置消费者 Worker 后,您可以从队列读取消息。

运行 wrangler tail 开始等待消费者记录收到的消息:

npx wrangler tail

wrangler tail 运行的情况下,打开步骤 3 中打开的 Worker URL。

您应在浏览器窗口中收到 Success 消息。

若收到 Success 消息,请刷新 URL 几次以生成消息并将其推送到队列。

wrangler tail 运行的情况下,消费者 Worker 将开始记录刷新产生的请求。

若刷新次数少于 10 次,消息可能需要几秒钟才会出现,因为批超时配置为 10 秒。10 秒后,消息应出现在终端中。

若刷新时出现错误,请检查步骤 2 中创建的队列名称与 Wrangler 文件中引用的队列是否相同。请确保生产者 Worker 返回 Success 而非错误。

完成本指南后,您已创建队列、向该队列发布消息的生产者 Worker,以及从该队列消费这些消息的消费者 Worker。

相关资源

  • 了解更多关于 Cloudflare Workers 以及您可以在 Cloudflare 上构建的应用程序。

这篇文档对您有帮助吗?