跳转到内容
搜索文档

Queues 工作原理

最后更新 查看 MarkdownAgent 设置

Cloudflare Queues 是一种灵活的消息队列,允许您将消息排队以进行异步处理。消息队列非常适合解耦应用的组件,例如电商网站的结账和订单履行服务。解耦的服务更易于理解、部署和实现,使您能够交付令客户满意的功能,而无需担心同步复杂的部署。Queues 还允许您批处理和缓冲对下游服务和 API 的调用。

理解 Queues 有四个主要概念:

  1. 队列
  2. 生产者
  3. 消费者
  4. 消息

什么是队列

队列是一个缓冲区或列表,随着消息写入而自动扩展,并允许消费者 Worker 从同一队列拉取消息。

Queues 设计为可靠,一旦写入成功,写入队列的消息不应丢失。同样,消息在消费者成功消费之前不会从队列中删除。

Queues 不保证消息会按发布顺序传递给消费者。

开发者可以创建多个队列。创建多个队列可用于:

  • 分离不同的用例和处理要求:例如,日志队列 vs. 密码重置队列。
  • 通过使用多个队列进行水平扩展来提高整体吞吐量(每秒消息数)。
  • 为连接到队列的每个消费者配置不同的批处理策略。

对于大多数应用,每个队列一个生产者 Worker,加上一个从该队列消费消息的消费者 Worker,可让您在逻辑上分离每个队列的处理。

生产者

生产者是向队列发布或生产消息的客户端术语。生产者通过绑定(binding)队列到 Worker 并调用该绑定向队列写入消息来配置。

例如,如果我们将名为 my-first-queue 的队列绑定到 MY_FIRST_QUEUE,可以通过调用绑定上的 send() 向队列写入消息:

interface Env {
  readonly MY_FIRST_QUEUE: Queue;
}

export default {
  async fetch(req, env, ctx): Promise<Response> {
    const message = {
      url: req.url,
      method: req.method,
      headers: Object.fromEntries(req.headers),
    };

    await env.MY_FIRST_QUEUE.send(message); // This will throw an exception if the send fails for any reason
    return new Response("Sent!");
  },
} satisfies ExportedHandler<Env>;

一个队列可以有多个生产者 Worker。例如,您可能有多个生产者 Worker 根据用户的传入 HTTP 请求向共享队列写入事件或日志。写入单个队列的生产者 Worker 总数没有限制。

此外,多个队列可以绑定到单个 Worker。该 Worker 可以根据代码中定义的任何逻辑决定写入哪个队列(或写入多个)。

内容类型

发布到队列的消息可以不同格式发布,具体取决于与消费者所需的互操作性。默认内容类型为 json,这意味着任何可传递给 JSON.stringify() 的对象均可接受。

要显式设置内容类型或指定替代内容类型,请将 contentType 选项传递给队列的 send() 方法:

interface Env {
  readonly MY_FIRST_QUEUE: Queue;
}

export default {
  async fetch(req, env, ctx): Promise<Response> {
    const message = {
      url: req.url,
      method: req.method,
      headers: Object.fromEntries(req.headers),
    };
    try {
      await env.MY_FIRST_QUEUE.send(message, { contentType: "json" }); // "json" is the default
      return new Response("Sent!");
    } catch (e) {
      // Catch cases where send fails, including due to a mismatched content type
      const msg = e instanceof Error ? e.message : "Unknown error";
      return Response.json({ error: msg }, { status: 500 });
    }
  },
} satisfies ExportedHandler<Env>;

要在写入队列时仅接受简单字符串,请改为设置 { contentType: "text" }

interface Env {
  readonly MY_FIRST_QUEUE: Queue;
}

export default {
  async fetch(req, env, ctx): Promise<Response> {
    try {
      // This will throw an exception (error) if you pass a non-string to the queue,
      // such as a native JavaScript object or ArrayBuffer.
      await env.MY_FIRST_QUEUE.send("hello there", { contentType: "text" }); // explicitly set 'text'
      return new Response("Sent!");
    } catch (e) {
      const msg = e instanceof Error ? e.message : "Unknown error";
      return Response.json({ error: msg }, { status: 500 });
    }
  },
} satisfies ExportedHandler<Env>;

QueuesContentType API 文档描述了每种格式如何序列化到队列。

消费者

Queues 支持两种类型的消费者:

  1. 消费者 Worker,基于推送:当队列有消息要传递时调用 Worker。
  2. HTTP 拉取消费者,基于拉取:消费者通过 HTTP 调用队列端点以接收并确认消息。

一个队列只能配置一种类型的消费者。

创建消费者 Worker

消费者是从队列订阅或消费消息的客户端术语。最基本的形式是,通过在 Worker 中创建 queue 处理程序来定义消费者:

interface Env {
  // Add your bindings here, e.g. KV namespaces, R2 buckets, D1 databases
}

export default {
  async queue(batch, env, ctx): Promise<void> {
    // Do something with messages in the batch
    // i.e. write to R2 storage, D1 database, or POST to an external API
    for (const msg of batch.messages) {
      // Process each message
      console.log(msg.body);
    }
  },
} satisfies ExportedHandler<Env>;

然后使用 wrangler queues consumer <queue-name> <worker-script-name> 将消费者连接到队列,或在 Wrangler 配置文件 中手动定义 [[queues.consumers]] 配置:

{
	"queues": {
		"consumers": [
			{
				"queue": "<your-queue-name>",
				"max_batch_size": 100, // optional
				"max_batch_timeout": 30 // optional
			}
		]
	}
}
[[queues.consumers]]
queue = "<your-queue-name>"
max_batch_size = 100
max_batch_timeout = 30

重要的是,每个队列只能有一个活跃消费者。这使 Cloudflare Queues 能够实现至少一次交付,并最小化超出此范围的重复消息风险。

值得注意的是,您可以将同一消费者与多个队列一起使用。定义消费者 Worker 的 queue 处理程序将由其连接的队列调用。

  • 传递给 queue 处理程序的 MessageBatch 包含一个 queue 属性,其中包含读取批次的队列名称。
  • 这可以减少您需要编写的代码量,并允许您根据队列名称处理消息。

例如,配置为从多个队列消费消息的消费者如下所示:

interface Env {
  // Add your bindings here
}

export default {
  async queue(batch, env, ctx): Promise<void> {
    // MessageBatch has a `queue` property we can switch on
    switch (batch.queue) {
      case "log-queue":
        // Write the batch to R2
        break;
      case "debug-queue":
        // Write the message to the console or to another queue
        break;
      case "email-reset":
        // Trigger a password reset email via an external API
        break;
      default:
        // Handle messages we haven't mentioned explicitly (write a log, push to a DLQ)
        break;
    }
  },
} satisfies ExportedHandler<Env>;

移除消费者

要从项目中移除队列,请运行 wrangler queues consumer remove <queue-name> <script-name>,然后从 Wrangler 文件中 [[queues.consumers]] 下方移除所需队列。

拉取消费者

队列可以有基于 HTTP 的消费者从队列拉取,而不是将消息推送到 Worker。

该消费者可以是任何可通过 Internet 通信的 HTTP 服务。查看拉取消费者指南了解如何为队列配置基于拉取的消费。

消息

消息是您向队列生产以及从队列消费的对象。

任何 JSON 可序列化对象都可以发布到队列。对于大多数开发者,这意味着简单字符串或 JSON 对象。发送消息时可以显式设置内容类型

消息在传递给消费者时可以被批处理。默认情况下,批次内的消息在确定重试时被视为全有或全无。如果批次中最后一条消息处理失败,整个批次将被重试。您还可以选择显式确认已成功处理的消息,和/或标记个别消息进行重试。

这篇文档对您有帮助吗?